# benchmarks/scale: multi-backend scale and load suite This suite runs one workload against all three ways to deploy the mentat core, up to hundreds of millions of datoms. It is meant to be re-run to validate future releases. | backend | what runs | how the suite drives it | |---|---|---| | `embedded` | the `mentat` Rust crate (`Store`) | `runner/`: the `mentat-scale` binary, one `Store` per thread | | `sqlite-ext` | `crates/sqlite/ext`, `libmentat_sqlite.so` in a host SQLite | `bench.py`: Python `sqlite3` + `load_extension`, one process per client | | `pg` | `crates/pg/pg_mentat` in PostgreSQL 16 | `pgbench -f pgbench/*.sql` (`\set` random params), `psql` for load | | `duckdb` | `crates/duckdb`, `mentat.duckdb_extension` in DuckDB v1.5.5 | `bench.py`: Python `duckdb==1.5.5` (needs Python >= 3.10), one process per client | | `duckdb-quack` | the same extension, loaded in ONE long-lived DuckDB Quack server (`crates/duckdb/server/serve.sh`, started and stopped by `run.sh` per scale) | `bench.py`: client processes send the duckdb backend's SQL through `quack_query` (Python `duckdb` + `LOAD quack`, no mentat in the client) | The workload is phase2's issue tracker, the same `schema.edn`, SEED and value distributions as `../phase2/gen_dataset.py`. `gen.py` imports phase2's constants and shards the output, so it can stream hundreds of millions of datoms to disk in parallel. ## Layout ``` gen.py dataset generator (scales xs/s/m/l/xl); writes pg/*.sql, store/*.edn, meta.json (truth) queries/*.edn the Datalog queries every backend runs (q1..q4, since, inputs, pull) pgbench/*.sql pgbench scripts for the same queries (+ write_state.sql writer) runner/ Rust crate `mentat_scale_runner` -> `mentat-scale` (embedded backend) bench.py correctness checks; sqlite-ext/duckdb load+bench; pgbench log parsing; medians run.sh the driver report.py timings.csv + loads.csv -> summary.md tables (+ findings.md if present) compare.py diff two result dirs, flag regressions ec2/bootstrap.sh fresh AL2023 host: NVMe RAID-0, toolchains, PG 16 (no cassert), duckdb CLI + venv ec2/tune.sh kernel + postgresql.conf tuning (shared_buffers = 85% RAM, hugepages...) ``` ## Scales | name | users | issues | labels | ~datoms | |---|---:|---:|---:|---:| | xs | 200 | 13,000 | 50 | 0.1M (local smoke) | | s | 1,600 | 130,000 | 200 | 1M | | m | 16,000 | 1,300,000 | 500 | 10M | | l | 160,000 | 13,000,000 | 1,000 | 100M | | xl | 500,000 | 40,000,000 | 2,000 | 300M | Each issue has 6 cardinality-one attributes and 0 to 3 labels. About 5% of issues get one later `:issue/state` change in a *history phase* after the initial load. The loader records two tx ids: `t_mid`, the last tx of the initial load, and `t_since`, the last tx before the final history file. Those make `as_of` and `since` meaningful and checkable. `meta.json` holds the exact expected answers. ## Scenarios | scenario | what it measures | check (fails the run if wrong) | |---|---|---| | `bulk_load` | load from empty. Wall time and datoms/s. PG adds ANALYZE and VACUUM time and `pg_database_size` vs `shared_buffers`. | the datom count follows from the checks below | | `point_lookup` | q1: user by unique `:user/email`, random user | exactly 1 row, the right name | | `ref_traversal` | q2: issues assigned to a random user (email -> ref -> title, state) | the set of (issue, state) pairs matches truth | | `aggregate` | q3: count by state (full scan of one attribute) | counts sum to n_issues and equal the per-state truth | | `predicate_scan` | q4: open issues with priority >= 4 | row count equals truth | | `pull` | `edn_pull('[*]', e)` on a random issue (embedded: `(pull ?e [*])`) | title, state and priority of fixed issues | | `as_of` | q2 with `{"asOf": t_mid}` | returns the **pre-update** states | | `since` | `[:find ?i :where [?i :issue/state _]]` with `{"since": t_since}` | distinct ?i equals the number of updates in the last history file | | `input_bindings` | `:in $ [?email ...]` with 100 random emails | the right names for a fixed set | | `write_mixed` | 1 writer (state update or label add, 1 datom per tx) plus MIXED_READERS readers doing q1/q2 | a post-write check: point lookup, aggregate sum, as_of, input bindings | | `concurrency_sweep` | read mix (q1, q2, pull in equal parts) at CLIENTS = 1 8 32 64 128 | (the checks above) | | `cold_vs_warm` | restart (PG: `pg_ctl restart`; embedded: new process + `Store::open`), drop the OS page cache, then the first call vs the next 30 | | | `sustained` | SUSTAINED_S seconds of read mix at SUSTAINED_CLIENTS plus 1 writer on the largest scale. RSS, iostat, pg_stat_io/bgwriter/database sampled every 10 s. 10-second latency windows for drift. | | Concurrency model: - **embedded**: threads, one `Store` per thread on the same file (`Store` is `!Sync`). Exactly one writer, because each `Store` keeps its own partition map and two writers would allocate the same tx id. - **sqlite-ext/duckdb**: processes, each with its own host connection. Since the store cache (1.10.0) each process reuses its open mentat store across calls; before it, every call was a `Store::open`. DuckDB is in-process with a single writer, so the sweep runs parallel reader processes. - **duckdb-quack**: client processes, each with its own Quack client connection, all executing inside one server process (its worker threads, each with its own cached store). - **pg**: `pgbench -c N -j min(N, nproc) -M simple`. Prepared/extended mode would rewrite the `:keywords` inside the Datalog literals into `$N` parameters. ## Output Results go to `benchmarks/results/scale-/`: - `raw.csv`: every (scenario, backend, scale, clients, op, rep) row. - `timings.csv`: the median over reps, with the schema `scenario,backend,scale,n_datoms,clients,op,count,p50_ms,p95_ms,p99_ms,max_ms,throughput_ops_s,errors,reps`. `op` is `read`, `write`, `ceiling` (the first call exceeded PROBE_S, so the scenario was skipped; the latency shown is that first call), `open`, or `cold_*`/`warm_*`. - `loads.csv`: bulk-load rows. `ceiling=1` means LOAD_MAX_S was hit, and the row carries the datoms/s reached. - `checks.txt`: every correctness check. - `sizes.txt`: the PG DB size vs `shared_buffers`. - `plans/`: `mentat_explain` of q1 to q4 per scale. - `logs/`: pgbench output, pg_stat_statements, sampler/iostat, load logs. - `env.txt`: hardware, kernel, tuning, the full non-default `pg_settings`, versions, git SHA, every knob. - `summary.md`: generated tables plus the hand-written `findings.md`. Methodology: 3 warm-up calls are discarded. Each point runs at least MIN_S (10) seconds and MIN_N (30) samples, capped at MAX_S (60) s. There are REPS (3) repetitions, and timings.csv reports the median. The seeds are fixed. The run exits non-zero if any check fails, so a fast wrong answer never looks like a win. ## Run it ### Locally, at a small scale (a quick check that nothing is broken) ```bash cargo build --release -p mentat_scale_runner -p mentat_sqlite_ext # PG: a pg16 with pg_mentat installed (cargo pgrx install --release) and # pg_stat_statements preloaded. PGHOST/PGPORT pointing at it. SCALES=xs BACKENDS="embedded sqlite-ext pg" REPS=1 MIN_S=2 MAX_S=10 MIXED_S=5 \ CLIENTS="1 4" benchmarks/scale/run.sh # duckdb: add it to BACKENDS and set PY=/path/to/python3.11-venv/bin/python # (with duckdb==1.5.5) and DUCKDB_EXT=crates/duckdb/build/release/mentat.duckdb_extension ``` ### On EC2 at full scale (what produced `results/scale-*`) ```bash # r6id.metal (128 vCPU, 1 TiB, 4x1.9 TB NVMe), AL2023, us-east-2. rsync -a --exclude .git --exclude target ./ HOST:/nvme/mentat/ # + crates/duckdb/extension-ci-tools/ ssh HOST 'bash /nvme/mentat/benchmarks/scale/ec2/bootstrap.sh' # RAID-0 /nvme, PG 16, toolchains ssh HOST 'cd /nvme/mentat && cargo build --release -p mentat_sqlite_ext -p mentat_scale_runner && (cd crates/pg/pg_mentat && cargo pgrx install --release --pg-config /nvme/pg16/bin/pg_config) && (cd crates/duckdb && make configure PYTHON_BIN=python3.11 && make release)' ssh HOST 'bash /nvme/mentat/benchmarks/scale/ec2/tune.sh' # hugepages, s_b=85%, start PG ssh HOST 'cd /nvme/mentat && PGHOST=/tmp PGBIN=/nvme/pg16/bin PGDATA=/nvme/pgdata PY=/nvme/venv/bin/python \ DATA_ROOT=/nvme/data SCALES="s m l" SUSTAINED_S=1800 MENTAT_GIT= \ nohup setsid benchmarks/scale/run.sh > /nvme/run.log 2>&1 &' ``` PG's database must stay **smaller than shared_buffers**, and `sizes.txt` records the ratio at each scale. On r6id.metal, s_b is about 870 GiB. ### Knobs (env vars, all recorded in env.txt) - `PHASE`: `all` (default), `load` (load + check only), `bench` (reuse already-loaded stores/DBs), or `sustained` (only the sustained run). Several invocations can share one `OUT` dir: raw.csv, loads.csv and checks.txt are appended, and env.txt records every invocation's knobs. - `SCALES`, `BACKENDS`: select what runs. - `SCENARIOS`: the single-client query scenarios. `SCENARIOS=` (empty) means none. - `EXTRA`: `write_mixed concurrency_sweep cold_vs_warm`. - `CLIENTS`, `REPS`, `MIN_S`, `MIN_N`, `MAX_S`: sampling. - `PROBE_S`: the per-call ceiling (default 20 s). - `LOAD_MAX_S`: the per-backend bulk-load cap (default 5400 s). A backend that cannot load a scale in time gets a `ceiling` row with its achieved rate, and its scenarios for that scale are skipped. - `MIXED_S`, `MIXED_READERS`, `SUSTAINED_S`, `SUSTAINED_CLIENTS`. - `EXT_SCENARIO_FILTER`: scenarios to skip on sqlite-ext and duckdb. - `PG_LOAD_JOBS`: parallel psql loaders. - `OPEN_PAR`: embedded stores opened at once (default 8). 128 concurrent `Store::open` calls OOM-killed the runner at 10M datoms. - `CHECK_TIMEOUT_S`: per-call cap for the as_of check (default 120). A timeout is a SKIP. - `PG_MAX_RESULT_ROWS`, `PG_TEMP_FILE_LIMIT`, `PG_SLOW_QUERY_MS`: pg_mentat GUCs set per bench DB. - `DATA_ROOT`, `WORK`, `OUT`, `PY`, `RUNNER`, `SQLITE_EXT`, `DUCKDB_EXT`, `PGBIN`, `PGDATA`, `MENTAT_GIT`. - `DUCKDB_CLI` (DuckDB v1.5.5 CLI that runs the duckdb-quack server), `QUACK_PORT` (default 9494), `MENTAT_QUACK_TOKEN` (default: random per run). - `SUSTAINED_S` > 0 also runs `sustained` for sqlite-ext, duckdb and duckdb-quack when they are in BACKENDS (`bench.py mixed … sustained`: 1 writer + SUSTAINED_CLIENTS reader processes). The first full run is `benchmarks/results/scale-2026-09-27T010840Z/`, on r6id.metal at s/m/l/xl. Read its `summary.md` (findings plus tables) before comparing against it. Several of its phases overlapped, and its `findings.md` lists them. ## Compare two runs ```bash benchmarks/scale/compare.py benchmarks/results/scale-OLD benchmarks/results/scale-NEW --threshold 10 ``` Rows are matched on (scenario, backend, scale, clients, op). The script flags p50/p99 increases and throughput drops beyond the threshold. It also flags new errors, a scenario that became a ceiling, and rows missing from NEW. It exits 1 on any regression, so it can gate a release. Compare runs from the same instance type only. ## Known limits of the suite (read before citing numbers) - **The embedded bulk load must use transactions under 5461 datoms.** `crates/sqlite/db/src/db.rs` `insert_non_fts_searches` asserts `6 * n < 32766`, and 5461 datoms trip it. `gen.py` batches 500 issues (about 3.8K datoms) per tx. - **Before 1.10.0 every sqlite-ext and duckdb call was a full `Store::open`**, which read the partition map through the `parts` view (a GROUP BY over all of `timelined_transactions`), so the per-call cost grew with history (about 480 ms at 1M datoms). 1.10.0 persists the partition marks (O(1) open) and caches the open store per path; runs before and after are not comparable for the ext backends. - **Embedded as_of is slow.** It uses a correlated `NOT EXISTS` over `timelined_transactions`, which has no (e, a, v) index. It usually shows up as a `ceiling`. - **pg_mentat errors when a result has more than `mentat.max_result_rows` rows** (default 100000). q4 and `since` pass that cap from scale m up, so `run.sh` sets `ALTER DATABASE … SET mentat.max_result_rows = 0` (override with `PG_MAX_RESULT_ROWS`). The same goes for `mentat.temp_file_limit` (default 1GB, applied with SET LOCAL per query). q3's `COUNT(DISTINCT)` sort goes past it at xl, so run.sh sets it to 100GB (`PG_TEMP_FILE_LIMIT`). - **An interrupted embedded query panics.** The PROBE_S watchdog uses `sqlite3_interrupt`, and mentat's projector unwraps the row iterator (`query-projector/src/projectors/simple.rs`), which poisons the Store's mutex. The runner catches the unwind and reopens its stores. - The pg loader uses explicit entids (bands of 1e10), because pg_mentat accepts caller-chosen ids. The embedded, sqlite-ext and duckdb loaders use tempids and lookup-refs, because embedded mentat only accepts entids it allocated. The datoms are the same, but the transaction shape differs.