- 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>
523 lines
13 KiB
Text
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();
|
|
}
|