writeonce/docs/examples/db-bench/main.wo
shoney.arickathil b7e6c24075 feat(db-bench): mix + msgrate — concurrent modes (T3)
- Mixer actors: 90/10 read/write, per-actor histograms merged through
  the store itself (Hist rows) — exact aggregate percentiles
- msgrate: one-way flood at a worker-placed sink; measured 15.3M
  msgs/s same-heap vs 2.06M cross-shard — the mutex-inbox number
- finding: point lookups are O(table) (probe walks all slabs), so
  read-heavy mix is quadratic in store size — all-mode calibrated to
  N/10 mix ops; the number 22 exists to publish
- TSan clean both shard counts (setarch -R, fibers-gate pattern)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-21 16:42:01 +02:00

523 lines
13 KiB
Text

use time
-- db-bench — iteration 22's load generator. Every measured mode prints
-- one machine-parsable line per operation class:
--
-- <op> <count> <ops/sec> <p50us> <p99us>
--
-- Timing is per-operation via time.ticks (CLOCK_MONOTONIC µs);
-- percentiles come from a 1µs-bucket histogram clamped at HIST_CLAMP —
-- exact to the microsecond below the clamp, and the clamp bucket keeps
-- the tail honest (a p99 AT the clamp means "clamp or worse").
-- The wal mode prints a running `acked <n>` line after every insert
-- RETURNS (the return IS the ack): the crash battery kills this mode
-- mid-run and verify-acked proves every acknowledged row survived.
-- ---- deterministic helpers ----
fn lcg(seed: Int) -> Int {
let x = seed * 1103515245 + 12345;
if x < 0 {
x = 0 - x;
}
return x;
}
fn item_v(i: Int) -> Int {
return (i * 37) % 1000;
}
-- ---- the histogram (percentiles without a sort) ----
fn hist_add(mut h: map<Int, Int>, us: Int) {
let b = us;
if b < 0 {
b = 0;
}
if b > 20000 {
b = 20000;
}
if has(h, b) {
set(h, b, get(h, b) + 1);
} else {
set(h, b, 1);
}
}
fn hist_pct(h: map<Int, Int>, total: Int, pct: Int) -> Int {
let target = total * pct / 100;
if target < 1 {
target = 1;
}
let seen = 0;
let b = 0;
while b <= 20000 {
if has(h, b) {
seen = seen + get(h, b);
if seen >= target {
return b;
}
}
b = b + 1;
}
return 20000;
}
fn report(op: Text, n: Int, total_us: Int, h: map<Int, Int>) {
let us = total_us;
if us < 1 {
us = 1;
}
let rate = n * 1000000 / us;
print("${op} ${n} ${rate} ${hist_pct(h, n, 50)} ${hist_pct(h, n, 99)}");
}
-- ---- modes ----
-- seed N: N children, one parent per 100, k = i % (N/10) (10 rows per
-- key), v deterministic. Meta rows record the expectations verify reads.
fn seed(n: Int) -> Int {
let h: map<Int, Int> = {};
let kmod = n / 10;
if kmod < 1 {
kmod = 1;
}
-- bucket-major: one parent, then its 100 children, using a single ref
-- local. (A hand-built `multi Bucket` of insert results SEGVs on drop —
-- the compiler classifies the elements OWNED while table refs are
-- scalar ids; recorded as a standing finding, not this iteration's fix.
-- Query-built multis are runtime-typed and safe.)
let vsum = 0;
let t0 = time.ticks();
let i = 1;
let b = 0;
while i <= n {
let bref = insert Bucket { tag: "b${b}" };
b = b + 1;
let j = 0;
while j < 100 and i <= n {
let o0 = time.ticks();
insert Item { k: i % kmod, v: item_v(i), bucket: bref };
hist_add(h, time.ticks() - o0);
vsum = vsum + item_v(i);
i = i + 1;
j = j + 1;
}
}
let t1 = time.ticks();
insert Meta { tag: "count", val: n };
insert Meta { tag: "vsum", val: vsum };
insert Meta { tag: "kmod", val: kmod };
report("seed", n, t1 - t0, h);
return 0;
}
fn meta_val(tag: Text) -> Int {
let ms = from m in Meta where m.tag == tag take 1 select m;
if len(ms) == 0 {
return -1;
}
return ms[0].val;
}
-- read N: indexed take-1 point lookups (the point-read this surface
-- offers), keys spread by LCG over the seeded key range.
fn read_mode(n: Int) -> Int {
let kmod = meta_val("kmod");
if kmod < 1 {
print_err("read: seed first");
return 1;
}
let h: map<Int, Int> = {};
let sink = 0;
let s = 42;
let t0 = time.ticks();
let i = 0;
while i < n {
s = lcg(s);
let key = s % kmod;
let o0 = time.ticks();
let xs = from x in Item where x.k == key take 1 select x;
if len(xs) > 0 {
sink = sink + xs[0].v;
}
hist_add(h, time.ticks() - o0);
i = i + 1;
}
let t1 = time.ticks();
report("read", n, t1 - t0, h);
if sink < 0 {
print("impossible ${sink}");
}
return 0;
}
-- query N: full equality probes on the k index (≈10 rows per key),
-- each materialized and counted.
fn query_mode(n: Int) -> Int {
let kmod = meta_val("kmod");
if kmod < 1 {
print_err("query: seed first");
return 1;
}
let h: map<Int, Int> = {};
let rows = 0;
let s = 7;
let t0 = time.ticks();
let i = 0;
while i < n {
s = lcg(s);
let key = s % kmod;
let o0 = time.ticks();
for x in from x in Item where x.k == key select x {
rows = rows + 1;
}
hist_add(h, time.ticks() - o0);
i = i + 1;
}
let t1 = time.ticks();
report("query", n, t1 - t0, h);
print("query rows ${rows}");
return 0;
}
-- write N: alternating inserts (disjoint k range 2e6+) and updates
-- through a query result. Corrupts vsum by design — the durability legs
-- run on their own fresh store.
fn write_mode(n: Int) -> Int {
let kmod = meta_val("kmod");
if kmod < 1 {
print_err("write: seed first");
return 1;
}
let bs = from b in Bucket where b.tag == "b0" take 1 select b;
if len(bs) == 0 {
print_err("write: no buckets");
return 1;
}
let h: map<Int, Int> = {};
let s = 99;
let t0 = time.ticks();
let i = 0;
while i < n {
let o0 = time.ticks();
if i % 2 == 0 {
insert Item { k: 2000000 + i, v: item_v(i), bucket: bs[0] };
} else {
s = lcg(s);
let key = s % kmod;
let xs = from x in Item where x.k == key take 1 select x;
if len(xs) > 0 {
xs[0].v = xs[0].v + 1;
}
}
hist_add(h, time.ticks() - o0);
i = i + 1;
}
let t1 = time.ticks();
report("write", n, t1 - t0, h);
return 0;
}
-- wal N: the crash battery's vehicle — insert-only, disjoint k range
-- (1e6+), `acked <i>` printed AFTER each insert returns (the return is
-- the ack: RAM applied, record staged, ONE commit done).
fn wal_mode(n: Int) -> Int {
let bs = from b in Bucket where b.tag == "b0" take 1 select b;
if len(bs) == 0 {
push(bs, insert Bucket { tag: "b0" });
}
let i = 1;
while i <= n {
insert Item { k: 1000000 + i, v: item_v(i), bucket: bs[0] };
print("acked ${i}");
i = i + 1;
}
return 0;
}
-- verify: the store against its own Meta expectations — count, checksum,
-- one unique-index probe. Exit 3 on any mismatch.
fn verify() -> Int {
let want_n = meta_val("count");
let want_sum = meta_val("vsum");
if want_n < 0 or want_sum < 0 {
print_err("verify: no meta (seed first)");
return 3;
}
let got_n = 0;
let got_sum = 0;
for x in from x in Item select x {
if x.k < 1000000 {
got_n = got_n + 1;
got_sum = got_sum + x.v;
}
}
if got_n != want_n or got_sum != want_sum {
print_err("verify: count ${got_n}/${want_n} sum ${got_sum}/${want_sum}");
return 3;
}
let bs = from b in Bucket where b.tag == "b0" take 1 select b;
if len(bs) == 0 {
print_err("verify: unique probe b0 missing");
return 3;
}
print("verify ok ${got_n} rows sum ${got_sum}");
return 0;
}
-- verify-acked M: after a kill -9 mid-wal — rows 1..M (k = 1e6+i) must
-- exist with the right v; rows beyond M are allowed (acked after the
-- last print landed). Exit 3 on any missing/wrong row.
fn verify_acked(m: Int) -> Int {
let i = 1;
while i <= m {
let key = 1000000 + i;
let xs = from x in Item where x.k == key take 1 select x;
if len(xs) == 0 {
print_err("verify-acked: row ${i} missing");
return 3;
}
if xs[0].v != item_v(i) {
print_err("verify-acked: row ${i} v ${xs[0].v} != ${item_v(i)}");
return 3;
}
i = i + 1;
}
print("verify-acked ok ${m} rows");
return 0;
}
-- ---- the concurrent modes (mix, msgrate) ----
fn hist_dump(h: map<Int, Int>, kind: Int) {
let b = 0;
while b <= 20000 {
if has(h, b) {
insert Hist { kind: kind, b: b, c: get(h, b) };
}
b = b + 1;
}
}
-- One mixer = one actor: 90/10 read/write over the seeded store. On a
-- worker shard every statement below rides the stage-3 DB RPC — the
-- code must not know or care (transparency is the point). Done signal:
-- a Meta row main polls for (the coordination idiom this side of 31).
class Mixer {
id: Int
fn receive(msg: MixJob) {
let hr: map<Int, Int> = {};
let hw: map<Int, Int> = {};
let s = msg.seed;
let sink = 0;
let i = 0;
while i < msg.ops {
s = lcg(s);
let key = s % msg.kmod;
let o0 = time.ticks();
if i % 10 == 9 {
let xs = from x in Item where x.k == key take 1 select x;
if len(xs) > 0 {
xs[0].v = xs[0].v + 1;
}
hist_add(hw, time.ticks() - o0);
} else {
let xs = from x in Item where x.k == key take 1 select x;
if len(xs) > 0 {
sink = sink + xs[0].v;
}
hist_add(hr, time.ticks() - o0);
}
i = i + 1;
}
hist_dump(hr, 0);
hist_dump(hw, 1);
insert Meta { tag: "mixdone${self.id}", val: sink };
}
}
fn mix_mode(total: Int, c: Int) -> Int {
let kmod = meta_val("kmod");
if kmod < 1 {
print_err("mix: seed first");
return 1;
}
let per = total / c;
if per < 1 {
per = 1;
}
let wall0 = time.ticks();
let i = 0;
while i < c {
let a: actor MixJob = spawn Mixer { id: i };
send(a, MixJob { ops: per, seed: 1000 + i * 7919, kmod: kmod });
i = i + 1;
}
-- poll until every mixer's done row exists
let done = 0;
while done < c {
time.sleep(20);
done = 0;
i = 0;
while i < c {
if meta_val("mixdone${i}") >= 0 {
done = done + 1;
}
i = i + 1;
}
}
let wall = time.ticks() - wall0;
-- merge the dumped histograms; wall time is shared by both classes
let hr: map<Int, Int> = {};
let hw: map<Int, Int> = {};
let nr = 0;
let nw = 0;
for x in from x in Hist select x {
if x.kind == 0 {
if has(hr, x.b) {
set(hr, x.b, get(hr, x.b) + x.c);
} else {
set(hr, x.b, x.c);
}
nr = nr + x.c;
} else {
if has(hw, x.b) {
set(hw, x.b, get(hw, x.b) + x.c);
} else {
set(hw, x.b, x.c);
}
nw = nw + x.c;
}
}
report("mixread", nr, wall, hr);
report("mixwrite", nw, wall, hw);
return 0;
}
-- msgrate: one-way flood — main sends N messages at a sink actor; the
-- sink counts and writes the done row at N. Spawn TWO sinks and flood
-- the second: round-robin placement puts it off the primary whenever
-- more than one shard exists, so the multi-shard number prices the
-- mutex inbox (stage-2 deviation 4's number); single-shard prices the
-- same-heap path.
class Sink {
got: Int
fn receive(msg: Flood) {
self.got = self.got + 1;
if self.got == msg.n {
insert Meta { tag: "flooddone", val: self.got };
}
}
}
fn msgrate_mode(n: Int) -> Int {
let first: actor Flood = spawn Sink { got: 0 };
let a: actor Flood = spawn Sink { got: 0 };
if first == a {
print_err("msgrate: impossible");
}
let t0 = time.ticks();
let i = 0;
while i < n {
send(a, Flood { n: n });
i = i + 1;
}
while meta_val("flooddone") < 0 {
time.sleep(5);
}
let us = time.ticks() - t0;
if us < 1 {
us = 1;
}
print("msgrate ${n} ${n * 1000000 / us}");
return 0;
}
-- all N: the throughput campaign in ONE process — without WO_DATA the
-- store is RAM and dies with the process, so seed and the measured
-- modes must share a run; under WO_DATA the same mode prices the
-- durable flavor. Restart/crash legs use the separate modes.
fn all_mode(n: Int) -> Int {
let rc = seed(n);
if rc != 0 {
return rc;
}
rc = read_mode(n / 2);
if rc != 0 {
return rc;
}
rc = query_mode(n / 10);
if rc != 0 {
return rc;
}
rc = write_mode(n / 2);
if rc != 0 {
return rc;
}
-- mix at N/10: every point lookup is O(table) today (the probe walks
-- all slabs — a headline finding, not a bug to hide), so a read-heavy
-- mix over a seeded store is quadratic in N. The campaign driver
-- chooses absolute sizes; this keeps `all` finishing in minutes.
return mix_mode(n / 10, 4);
}
fn usage() -> Int {
print_err("usage: db-bench <mode>");
print_err(" all N | seed N | read N | query N | write N | wal N");
print_err(" mix N C | msgrate N | verify | verify-acked M");
return 2;
}
fn main(args: multi Text) -> Int {
if len(args) < 1 {
return usage();
}
if args[0] == "verify" {
return verify();
}
if len(args) < 2 {
return usage();
}
let n = parse_int(args[1]);
if n == nil or n < 1 {
print_err("db-bench: <n> must be a positive number");
return 2;
}
if args[0] == "all" {
return all_mode(n);
}
if args[0] == "seed" {
return seed(n);
}
if args[0] == "read" {
return read_mode(n);
}
if args[0] == "query" {
return query_mode(n);
}
if args[0] == "write" {
return write_mode(n);
}
if args[0] == "wal" {
return wal_mode(n);
}
if args[0] == "verify-acked" {
return verify_acked(n);
}
if args[0] == "msgrate" {
return msgrate_mode(n);
}
if args[0] == "mix" {
if len(args) < 3 {
return usage();
}
let c = parse_int(args[2]);
if c == nil or c < 1 {
print_err("db-bench: <c> must be a positive number");
return 2;
}
return mix_mode(n, c);
}
return usage();
}