Skip to content

Commit 68de018

Browse files
committed
mason: ck-mc write amplification design r2: exact migration SQL, reader map, fail-closed path, plugin context.db measurements
Answers every item of the adversarial review of r1. Adds the migcheck probe (the proposed migration 63 as plain SQL, run with SQLite 3.46.0 and verified row by row on a clone: 1,400 of 1,400 exact) and a SQL write trace in drive.ts that measures the plugin's per-pass lkg_slots and session_meta writes.
1 parent 69c9729 commit 68de018

8 files changed

Lines changed: 1326 additions & 160 deletions

File tree

‎docs/reports/ckmc-write-amplification-design.md‎

Lines changed: 457 additions & 159 deletions
Large diffs are not rendered by default.

‎scripts/ckmc-write-probe/README.md‎

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,29 @@ Per pass it records the module's disk-write counter (`rusage.py`, macOS `proc_pi
2626

2727
Confirm isolation with `lsof -p <pid>` on the probe's `ck-mc` and `ck-subc`: every regular file must be under the run directory.
2828

29+
`PROBE_SQL_TRACE=1` also logs every write statement the plugin process runs against `context.db` to `$PROBE_RUN/sqltrace.jsonl`: the pass, the table, the columns it sets, its bound bytes, the WAL frames it appended (exact for autocommit statements under `PROBE_PIN=1`), and for each bound string of 64 KiB or more how much of it matches the same statement's previous value. `python3 sqlsum.py $PROBE_RUN` prints the per-pass summary. `pluginrows.py <context.db clone> <session prefix>...` lists the bytes per column of a session's `session_meta` and `lkg_slots` rows.
30+
31+
## The proposed migration 63
32+
33+
`migcheck/` holds the exact migration text (`migration63.sql`) and a Rust probe that runs it with the SQLite ck-mc links (3.46.0, through `rusqlite =0.32.1` with `bundled`). The crate is outside the repository's Cargo workspace; build it from a copy so no `Cargo.lock` or `target/` lands here:
34+
35+
```sh
36+
W=$TMPDIR/magic-context/ckmc-writes-r2
37+
cp -R scripts/ckmc-write-probe/migcheck $W/migcheck-src
38+
(cd $W/migcheck-src && CARGO_TARGET_DIR=$W/target-exact cargo build --release -j 2 --features exact)
39+
(cd $W/migcheck-src && CARGO_TARGET_DIR=$W/target-fast cargo build --release -j 2)
40+
cp -c $W/golden/mc/store.db $W/mig/store.db
41+
$W/target-exact/release/migcheck migrate $W/mig/store.db
42+
$W/target-exact/release/migcheck verify $W/golden/mc/store.db $W/mig/store.db
43+
$W/target-exact/release/migcheck fixtures
44+
$W/target-fast/release/migcheck loadcost $W/golden/mc/store.db $W/mig/store.db <session> 60
45+
```
46+
47+
`verify` compares every moved value with the serde-parsed original; `--features exact` compares numbers by their text and keeps key order, so it also reports byte identity. `MIGCHECK_SQL=<file>` runs a different migration text, which is how a deliberately broken aggregate shows that `verify` catches corruption. `loadcost` parses into `serde_json::Value`, a proxy for the typed structs, so compare the two layouts rather than the absolute times.
48+
2949
## SQL-level replay of one commit
3050

31-
`sqlexp.py <golden>/mc/store.db <work dir> <session> [scenario ...]` (the golden clone has its WAL folded in) runs the statements `McStore::commit_transform` issues, with a ck-mc connection's pragmas, on per-scenario clones, for today's layout and for the proposed split layout (including the cost of migrating every session).
51+
`sqlexp.py <golden>/mc/store.db <work dir> <session> [scenario ...]` (the golden clone has its WAL folded in) runs the statements `McStore::commit_transform` issues, with a ck-mc connection's pragmas, on per-scenario clones, for today's layout and for the proposed split layout. Its `migrate()` is a Python split used only to build the split layout for the cost model; it is not migration 63, and its timing is not the migration's (use `migcheck` for that).
3252

3353
## Live evidence without opening the live store
3454

‎scripts/ckmc-write-probe/drive.ts‎

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,11 @@
2020
* this makes frame attribution exact but defers checkpoint writes
2121
* PROBE_BROCA=1 start the hermetic historian producer so historian runs can publish
2222
* PROBE_OUT JSONL output path (default $PROBE_RUN/passes.jsonl)
23+
* PROBE_SQL_TRACE=1 record every write statement this process (the plugin) runs against
24+
* context.db to $PROBE_RUN/sqltrace.jsonl, with the pass it ran in and the
25+
* WAL frames it appended (exact for autocommit statements under PROBE_PIN=1)
2326
*/
27+
import { Database as BunDatabase } from "bun:sqlite";
2428
import { spawn, spawnSync, type ChildProcess } from "node:child_process";
2529
import { appendFileSync, existsSync, mkdirSync, readdirSync, statSync, writeFileSync } from "node:fs";
2630
import { join, resolve } from "node:path";
@@ -58,6 +62,99 @@ process.env.OPENCODE_DB = join(RUN, "oc", "opencode.db");
5862
process.env.TMPDIR = join(RUN, "tmp");
5963
mkdirSync(process.env.TMPDIR, { recursive: true });
6064

65+
/** Index of the plan entry being replayed, so each SQL write-trace entry names its pass. */
66+
let currentPass: number | string = "boot";
67+
const SQL_TRACE = join(RUN, "sqltrace.jsonl");
68+
const WAL_FRAME_BYTES = 4096 + 24;
69+
70+
/**
71+
* Wrap bun:sqlite so every write statement against context.db is logged. The plugin reaches
72+
* SQLite only through `Database` from shared/sqlite, which is bun:sqlite under Bun, so
73+
* patching the prototype sees every statement it runs. Frames are the growth of the
74+
* context.db WAL across the statement: exact for an autocommit statement while a pinned
75+
* reader stops the WAL from resetting, and zero for a statement inside a transaction, whose
76+
* frames land at its COMMIT.
77+
*/
78+
function installSqlTrace(): void {
79+
const walSize = (): number => statSync(`${CONTEXT_DB}-wal`, { throwIfNoEntry: false })?.size ?? 0;
80+
const byteLength = (value: unknown): number =>
81+
typeof value === "string" ? Buffer.byteLength(value) : value instanceof Uint8Array ? value.length : 8;
82+
const columnsOf = (sql: string): string[] => {
83+
const set = sql.match(/\bSET\b([\s\S]*?)(\bWHERE\b|$)/i)?.[1];
84+
if (set) return [...set.matchAll(/([A-Za-z_][A-Za-z0-9_]*)\s*=/g)].map((m) => m[1]);
85+
const insert = sql.match(/\bINTO\s+[A-Za-z_][A-Za-z0-9_]*\s*\(([^)]*)\)/i)?.[1];
86+
return insert ? insert.split(",").map((c) => c.trim()) : [];
87+
};
88+
// For every bound string of 64 KiB or more, how much of it the same statement's previous
89+
// run already wrote: the common prefix, and how many 64 KiB positional chunks differ. This
90+
// is what a chunked layout for the plugin's large columns would have to rewrite.
91+
const CHUNK_CHARS = 64 * 1024;
92+
const previousLarge = new Map<string, string>();
93+
const largeDiffs = (sql: string, values: unknown[]): unknown[] =>
94+
values.flatMap((value, index) => {
95+
if (typeof value !== "string" || value.length < CHUNK_CHARS) return [];
96+
const key = `${sql}\u0000${index}`;
97+
const previous = previousLarge.get(key);
98+
previousLarge.set(key, value);
99+
if (previous === undefined) return [{ index, chars: value.length, previous_chars: null }];
100+
let common = 0;
101+
const limit = Math.min(previous.length, value.length);
102+
while (common < limit && previous.charCodeAt(common) === value.charCodeAt(common)) common++;
103+
const chunks = Math.ceil(value.length / CHUNK_CHARS);
104+
let changedChunks = 0;
105+
for (let chunk = 0; chunk < chunks; chunk++) {
106+
const start = chunk * CHUNK_CHARS;
107+
if (value.slice(start, start + CHUNK_CHARS) !== previous.slice(start, start + CHUNK_CHARS)) changedChunks++;
108+
}
109+
return [{ index, chars: value.length, previous_chars: previous.length, common_prefix_chars: common, chunks, changed_chunks: changedChunks }];
110+
});
111+
const log = (db: BunDatabase, sql: string, args: unknown[], run: () => unknown): unknown => {
112+
const before = walSize();
113+
const inTransaction = db.inTransaction;
114+
const result = run() as { changes?: number } | undefined;
115+
const flat = args.length === 1 && Array.isArray(args[0]) ? (args[0] as unknown[]) : args;
116+
const large = largeDiffs(sql, flat);
117+
appendFileSync(
118+
SQL_TRACE,
119+
`${JSON.stringify({
120+
pass: currentPass,
121+
table: sql.match(/\b(?:INTO|UPDATE|FROM)\s+([A-Za-z_][A-Za-z0-9_]*)/i)?.[1] ?? "?",
122+
in_transaction: inTransaction,
123+
frames: (walSize() - before) / WAL_FRAME_BYTES,
124+
param_bytes: flat.reduce((sum: number, value) => sum + byteLength(value), 0),
125+
columns: columnsOf(sql),
126+
changes: result?.changes ?? null,
127+
large,
128+
sql: sql.replace(/\s+/g, " ").slice(0, 160),
129+
})}\n`,
130+
);
131+
return result;
132+
};
133+
const isWrite = (db: BunDatabase, sql: string): boolean =>
134+
db.filename.endsWith("context.db") && /^\s*(INSERT|UPDATE|REPLACE|DELETE)\b/i.test(sql);
135+
// biome-ignore lint/suspicious/noExplicitAny: patching bun:sqlite's runtime prototype.
136+
const proto = BunDatabase.prototype as any;
137+
// biome-ignore lint/suspicious/noExplicitAny: bun:sqlite statements are patched in place.
138+
const wrap = (db: BunDatabase, sql: string, statement: any): any => {
139+
if (!isWrite(db, sql) || statement.__probeTraced) return statement;
140+
const run = statement.run.bind(statement);
141+
statement.run = (...args: unknown[]) => log(db, sql, args, () => run(...args));
142+
statement.__probeTraced = true;
143+
return statement;
144+
};
145+
for (const method of ["prepare", "query"] as const) {
146+
const original = proto[method];
147+
proto[method] = function (this: BunDatabase, sql: string, ...rest: unknown[]) {
148+
return wrap(this, sql, original.call(this, sql, ...rest));
149+
};
150+
}
151+
const originalRun = proto.run;
152+
proto.run = function (this: BunDatabase, sql: string, ...args: unknown[]) {
153+
if (!isWrite(this, sql)) return originalRun.call(this, sql, ...args);
154+
return log(this, sql, args, () => originalRun.call(this, sql, ...args));
155+
};
156+
}
157+
61158
const children: ChildProcess[] = [];
62159
function cleanup(): void {
63160
for (const child of children) child.kill("SIGTERM");
@@ -114,6 +211,7 @@ function startProcess(label: string, cmd: string, args: string[], env: Record<st
114211
}
115212

116213
async function main(): Promise<void> {
214+
if (process.env.PROBE_SQL_TRACE === "1") installSqlTrace();
117215
mkdirSync(RUNTIME_DIR, { recursive: true });
118216
mkdirSync(HOME, { recursive: true });
119217
const daemonConfig = join(DATA, "cortexkit", "_daemon-config");
@@ -352,6 +450,7 @@ async function main(): Promise<void> {
352450
}
353451

354452
for (const [index, kind] of PLAN.entries()) {
453+
currentPass = index;
355454
const before = frames();
356455
const beforeSide = sideFileBytes();
357456
const now = Date.now();
@@ -400,6 +499,7 @@ async function main(): Promise<void> {
400499
...diff(before, after, module.pid!, daemon.pid!),
401500
});
402501
}
502+
currentPass = "end";
403503
if (process.env.PROBE_PIN === "1") {
404504
record({
405505
kind: "run_total_context_db",
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
# A throwaway probe, deliberately outside the repository's Cargo workspace. Build it from a copy
2+
# outside the repository so no Cargo.lock or target directory lands here; README.md has the
3+
# exact commands. rusqlite 0.32.1 with `bundled` pulls in libsqlite3-sys 0.30.1, the same
4+
# SQLite (3.46.0) that ck-mc links.
5+
[package]
6+
name = "migcheck"
7+
version = "0.0.0"
8+
edition = "2021"
9+
publish = false
10+
11+
[workspace]
12+
13+
[dependencies]
14+
rusqlite = { version = "=0.32.1", features = ["bundled"] }
15+
serde_json = "1"
16+
17+
[features]
18+
# Compare numbers by their text and keep object key order, so `verify` can also report
19+
# whether a migrated body is byte-identical to a serde re-serialization.
20+
exact = ["serde_json/arbitrary_precision", "serde_json/preserve_order"]
21+
22+
[profile.release]
23+
debug = false
Lines changed: 138 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,138 @@
1+
-- The proposed store.db migration 63, exactly as the design note specifies it.
2+
-- migcheck.rs runs this file unchanged against a clone of the whole store and checks every
3+
-- row it moves. It must stay plain SQL: cortexkit-store migrations are static SQL batches
4+
-- run inside one transaction together with their version record.
5+
--
6+
-- SQLite 3.46.0 (libsqlite3-sys 0.30.1, the version ck-mc links) or newer is required for
7+
-- ORDER BY inside an aggregate.
8+
9+
-- 1. Shape guards. Each INSERT counts the rows the codec could not round-trip, and the CHECK
10+
-- fails the whole migration, rolling it back, if any count is not zero. The store then
11+
-- stays at version 62 and an older binary can still open it.
12+
CREATE TEMP TABLE mc_migration_63_guard (
13+
problem TEXT NOT NULL,
14+
bad_rows INTEGER NOT NULL CHECK (bad_rows = 0)
15+
);
16+
INSERT INTO mc_migration_63_guard
17+
SELECT 'core_state.frozen_units is present but not an array', COUNT(*)
18+
FROM mc_cache_state
19+
WHERE json_type(core_state, '$.frozen_units') NOT IN ('array');
20+
INSERT INTO mc_migration_63_guard
21+
SELECT 'a frozen unit is not a JSON object', COUNT(*)
22+
FROM mc_cache_state AS s, json_each(s.core_state, '$.frozen_units') AS e
23+
WHERE json_type(s.core_state, '$.frozen_units') = 'array' AND e.type <> 'object';
24+
INSERT INTO mc_migration_63_guard
25+
SELECT 'meta.resolved_compartment_boundaries is present but not an array', COUNT(*)
26+
FROM mc_cache_state
27+
WHERE json_type(meta, '$.resolved_compartment_boundaries') NOT IN ('array', 'null');
28+
INSERT INTO mc_migration_63_guard
29+
SELECT 'a resolved compartment boundary is not a JSON object', COUNT(*)
30+
FROM mc_cache_state AS s, json_each(s.meta, '$.resolved_compartment_boundaries') AS e
31+
WHERE json_type(s.meta, '$.resolved_compartment_boundaries') = 'array' AND e.type <> 'object';
32+
INSERT INTO mc_migration_63_guard
33+
SELECT 'meta.tail_hygiene_baseline is present but not an object', COUNT(*)
34+
FROM mc_cache_state
35+
WHERE json_type(meta, '$.tail_hygiene_baseline') NOT IN ('object', 'null');
36+
DROP TABLE mc_migration_63_guard;
37+
38+
-- 2. Frozen units: fixed-size positional chunks of 64 units.
39+
-- json(e.value) is required. Under GROUP BY, SQLite's sorter drops the JSON subtype of
40+
-- json_each.value, and json_group_array(e.value) then stores every unit as a JSON string.
41+
CREATE TABLE mc_cache_frozen_chunks (
42+
session_id TEXT NOT NULL,
43+
chunk INTEGER NOT NULL,
44+
body TEXT NOT NULL,
45+
PRIMARY KEY (session_id, chunk)
46+
);
47+
INSERT INTO mc_cache_frozen_chunks (session_id, chunk, body)
48+
SELECT s.session_id, e.key / 64, json_group_array(json(e.value) ORDER BY e.key)
49+
FROM mc_cache_state AS s, json_each(s.core_state, '$.frozen_units') AS e
50+
GROUP BY s.session_id, e.key / 64;
51+
52+
-- 3. Values replaced wholesale: one row per non-empty section. No row means the section is
53+
-- absent (an empty boundary list, or no tail baseline).
54+
CREATE TABLE mc_cache_sections (
55+
session_id TEXT NOT NULL,
56+
section TEXT NOT NULL,
57+
body TEXT NOT NULL,
58+
PRIMARY KEY (session_id, section)
59+
);
60+
INSERT INTO mc_cache_sections (session_id, section, body)
61+
SELECT session_id, 'resolved_compartment_boundaries',
62+
json_extract(meta, '$.resolved_compartment_boundaries')
63+
FROM mc_cache_state
64+
WHERE json_type(meta, '$.resolved_compartment_boundaries') = 'array'
65+
AND json_array_length(meta, '$.resolved_compartment_boundaries') > 0;
66+
INSERT INTO mc_cache_sections (session_id, section, body)
67+
SELECT session_id, 'tail_hygiene_baseline', json_extract(meta, '$.tail_hygiene_baseline')
68+
FROM mc_cache_state
69+
WHERE json_type(meta, '$.tail_hygiene_baseline') = 'object';
70+
71+
-- 4. The small row. section_index records what the chunk and section rows must hold.
72+
-- "sv" is the sections version: 0 marks rows written by this migration, whose digests
73+
-- are filled in by the open-time backfill. Only codec writers set it, always to >= 1.
74+
-- json_patch drops the keys whose value is NULL, so an absent section has no key.
75+
-- Every SET expression reads the row as it was before this UPDATE.
76+
ALTER TABLE mc_cache_state ADD COLUMN section_index TEXT NOT NULL DEFAULT '{}';
77+
UPDATE mc_cache_state SET
78+
section_index = json_patch('{}', json_object(
79+
'sv', 0,
80+
'f', json_object(
81+
'n', COALESCE(json_array_length(core_state, '$.frozen_units'), 0),
82+
'c', (COALESCE(json_array_length(core_state, '$.frozen_units'), 0) + 63) / 64),
83+
'b', CASE WHEN json_array_length(meta, '$.resolved_compartment_boundaries') > 0
84+
THEN json_object('n', json_array_length(meta, '$.resolved_compartment_boundaries'))
85+
END,
86+
't', CASE WHEN json_type(meta, '$.tail_hygiene_baseline') = 'object'
87+
THEN json('{}')
88+
END)),
89+
core_state = json_remove(core_state, '$.frozen_units'),
90+
meta = json_remove(meta, '$.resolved_compartment_boundaries', '$.tail_hygiene_baseline');
91+
92+
-- 5. Pass-trace histories as ring rows. The sequence orders entries; the slot only bounds
93+
-- storage, so readers order by seq and never by slot. The request history lives today in a
94+
-- carrier entry inside scheduler_interesting_history; only the newest carrier counts,
95+
-- as in pass_trace_meta_parts.
96+
CREATE TABLE mc_pass_trace_history (
97+
session_id TEXT NOT NULL,
98+
kind TEXT NOT NULL, -- 'scheduler' | 'interesting' | 'request'
99+
slot INTEGER NOT NULL, -- seq % 256, or seq % 32 for 'request'
100+
seq INTEGER NOT NULL,
101+
entry TEXT NOT NULL,
102+
PRIMARY KEY (session_id, kind, slot)
103+
);
104+
INSERT INTO mc_pass_trace_history (session_id, kind, slot, seq, entry)
105+
SELECT t.session_id, 'scheduler', e.key % 256, e.key, json(e.value)
106+
FROM mc_pass_trace AS t, json_each(t.scheduler_history) AS e;
107+
INSERT INTO mc_pass_trace_history (session_id, kind, slot, seq, entry)
108+
SELECT t.session_id, 'interesting', e.key % 256, e.key, json(e.value)
109+
FROM mc_pass_trace AS t, json_each(t.scheduler_interesting_history) AS e
110+
WHERE json_extract(e.value, '$.scheduler_decision') IS NOT '__request_trace_history__';
111+
INSERT INTO mc_pass_trace_history (session_id, kind, slot, seq, entry)
112+
SELECT t.session_id, 'request', r.key % 32, r.key, json(r.value)
113+
FROM mc_pass_trace AS t, json_each(t.scheduler_interesting_history) AS c,
114+
json_each(c.value, '$.request_history') AS r
115+
WHERE c.key = (SELECT MAX(c2.key) FROM json_each(t.scheduler_interesting_history) AS c2
116+
WHERE json_extract(c2.value, '$.scheduler_decision') = '__request_trace_history__');
117+
ALTER TABLE mc_pass_trace ADD COLUMN scheduler_next_seq INTEGER NOT NULL DEFAULT 0;
118+
ALTER TABLE mc_pass_trace ADD COLUMN interesting_next_seq INTEGER NOT NULL DEFAULT 0;
119+
ALTER TABLE mc_pass_trace ADD COLUMN request_next_seq INTEGER NOT NULL DEFAULT 0;
120+
UPDATE mc_pass_trace SET
121+
scheduler_next_seq = COALESCE(json_array_length(scheduler_history), 0),
122+
interesting_next_seq = COALESCE(json_array_length(scheduler_interesting_history), 0),
123+
request_next_seq = COALESCE((SELECT MAX(seq) + 1 FROM mc_pass_trace_history AS h
124+
WHERE h.session_id = mc_pass_trace.session_id
125+
AND h.kind = 'request'), 0);
126+
-- Dropping the array columns makes a reader that names them fail loudly instead of reading
127+
-- a frozen copy. Readers that select * must move to the view below.
128+
ALTER TABLE mc_pass_trace DROP COLUMN scheduler_history;
129+
ALTER TABLE mc_pass_trace DROP COLUMN scheduler_interesting_history;
130+
CREATE VIEW mc_pass_trace_history_arrays AS
131+
SELECT t.session_id,
132+
(SELECT json_group_array(json(h.entry) ORDER BY h.seq) FROM mc_pass_trace_history AS h
133+
WHERE h.session_id = t.session_id AND h.kind = 'scheduler') AS scheduler_history,
134+
(SELECT json_group_array(json(h.entry) ORDER BY h.seq) FROM mc_pass_trace_history AS h
135+
WHERE h.session_id = t.session_id AND h.kind = 'interesting') AS scheduler_interesting_history,
136+
(SELECT json_group_array(json(h.entry) ORDER BY h.seq) FROM mc_pass_trace_history AS h
137+
WHERE h.session_id = t.session_id AND h.kind = 'request') AS request_history
138+
FROM mc_pass_trace AS t;

0 commit comments

Comments
 (0)