From 85e2190731edac5a57f0d460a47326e9c65db95c Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 23 Sep 2026 00:10:49 -0700 Subject: [PATCH 01/10] perf(sqlite): reuse canonical encoding within verified history reads Signed-off-by: Lihua <1017343802@qq.com> --- .../coordination/authority_state_log.ts | 7 +- .../coordination/authority_store_codec.ts | 65 +++++++++++++++++++ .../coordination/sqlite_authority_store.ts | 27 +++++--- 3 files changed, 89 insertions(+), 10 deletions(-) diff --git a/loopx/control_plane/coordination/authority_state_log.ts b/loopx/control_plane/coordination/authority_state_log.ts index 12e656e899..5376221a37 100644 --- a/loopx/control_plane/coordination/authority_state_log.ts +++ b/loopx/control_plane/coordination/authority_state_log.ts @@ -13,6 +13,7 @@ * never decides Todo, lease, quota or promotion semantics. */ import type {JsonObject} from "../effect_program.ts"; +import type {CanonicalAuthorityDigest} from "./authority_store_codec.ts"; import { AuthorityStoreProtocolError, authorityUnicodeCompare, @@ -61,8 +62,10 @@ function canonicalBytesEqual(left: unknown, right: unknown): boolean { } /** Digest of one committed projection; the anchor every reconstruction checks. */ -export function authorityStateDigest(projection: JsonObject): string { - return canonicalAuthoritySha256(projection); +export function authorityStateDigest( + projection: JsonObject, digest: CanonicalAuthorityDigest = canonicalAuthoritySha256, +): string { + return digest(projection); } /** The one checkpoint cursor that covers a positive commit cursor. */ diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index a386352b7a..33c31112eb 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -108,6 +108,71 @@ export function canonicalAuthoritySha256(value: unknown): string { return createHash("sha256").update(canonicalAuthorityBytes(value)).digest("hex"); } +export type CanonicalAuthorityDigest = (value: unknown) => string; + +/** + * Recompute the exact canonical v0 digest while reusing encoded long strings + * within one verification window. Retained projections often share a large + * unchanged value across many deltas; serializing it for every state and + * commit proof dominates historical reads. The cache is bounded and never + * survives the caller's window, so a later read still verifies fresh bytes. + */ +export function createCanonicalAuthorityDigestWindow(): CanonicalAuthorityDigest { + const maxEntryBytes = 2 * 1024 * 1024; + const maxCachedBytes = 4 * 1024 * 1024; + const encoded = new Map(); + let cachedBytes = 0; + const stringBytes = (value: string): string | Buffer => { + if (value.length < 1024) return JSON.stringify(value); + const previous = encoded.get(value); + if (previous) return previous; + const bytes = Buffer.from(JSON.stringify(value), "utf8"); + if (bytes.byteLength > maxEntryBytes) return bytes; + while (cachedBytes + bytes.byteLength > maxCachedBytes) { + const oldest = encoded.keys().next().value; + if (oldest === undefined) break; + cachedBytes -= encoded.get(oldest)!.byteLength; + encoded.delete(oldest); + } + encoded.set(value, bytes); + cachedBytes += bytes.byteLength; + return bytes; + }; + return (value: unknown): string => { + const canonical = canonicalAuthorityJson(value); + const hash = createHash("sha256"); + const visit = (part: unknown): void => { + if (part === null) { hash.update("null"); return; } + if (typeof part === "string") { hash.update(stringBytes(part)); return; } + if (typeof part === "number" || typeof part === "boolean") { + hash.update(JSON.stringify(part)); return; + } + if (Array.isArray(part)) { + hash.update("["); + for (let index = 0; index < part.length; index++) { + if (index) hash.update(","); + if (Object.hasOwn(part, index)) visit(part[index]); + else hash.update("null"); // JSON.stringify represents sparse slots as null. + } + hash.update("]"); + return; + } + hash.update("{"); + let first = true; + for (const key of Object.keys(part as JsonObject)) { + if (!first) hash.update(","); + first = false; + hash.update(stringBytes(key)); + hash.update(":"); + visit((part as JsonObject)[key]); + } + hash.update("}"); + }; + visit(canonical); + return hash.digest("hex"); + }; +} + export function parseAuthorityCursor(value: string | null): bigint { if (value === null) return 0n; if (typeof value !== "string" || !/^[1-9]\d*$/.test(value)) { diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index ca54a225a3..7c3cc640ba 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -10,8 +10,9 @@ import type { AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult, AuthorityStoreIdentityResult, AuthorityStoreLoadResult, AuthorityStoreReadFailure, AuthorityStoreReceiptResult, AuthorityStoreScanResult } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityBytes, canonicalAuthorityObject, - canonicalAuthorityObjectList, canonicalAuthoritySha256, normalizeAuthorityStoreCommit, - requireAuthorityStoreId } from "./authority_store_codec.ts"; + canonicalAuthorityObjectList, canonicalAuthoritySha256, createCanonicalAuthorityDigestWindow, + normalizeAuthorityStoreCommit, requireAuthorityStoreId, + type CanonicalAuthorityDigest } from "./authority_store_codec.ts"; import { AUTHORITY_STATE_CHECKPOINT_INTERVAL, applyAuthorityStateDelta, @@ -258,6 +259,7 @@ export class SqliteAuthorityStore implements AuthorityStore { row: SqliteCommitRow, identity: string, base: {kind: "predecessor" | "sealed"; state: SqliteStateCursor}, + digest: CanonicalAuthorityDigest, ): SqliteVerifiedCommit { let projection: JsonObject; if (base.kind === "sealed") { @@ -273,11 +275,17 @@ export class SqliteAuthorityStore implements AuthorityStore { : row.parent_state_digest === base.state.digest; if (!parentMatches) protocol("SQLite authority state log parent lineage is invalid"); projection = applyAuthorityStateDelta(base.state.projection, row.delta); - if (authorityStateDigest(projection) !== row.state_digest) { + // A later empty delta preserves the already-verified predecessor state. + // Keep the exact commit proof below, including its events and receipts, + // but avoid rehashing an unchanged large projection just to rediscover + // the same state digest. The root has no verified predecessor digest. + const stateDigest = row.cursor > 1n && row.delta.operations.length === 0 + ? base.state.digest : authorityStateDigest(projection, digest); + if (stateDigest !== row.state_digest) { protocol("SQLite authority state log digest mismatch"); } } - const expected = commitDigest(identity, row.cursor, row.operation_id, projection, row.events, row.receipts); + const expected = commitDigest(identity, row.cursor, row.operation_id, projection, row.events, row.receipts, digest); if (row.commit_digest !== expected) protocol("SQLite committed row digest mismatch"); return { state: {cursor: row.cursor, projection, digest: row.state_digest}, @@ -321,12 +329,13 @@ export class SqliteAuthorityStore implements AuthorityStore { const identity = this.identity(db); const checkpoint = this.loadCheckpoint(db, from); const rows = this.windowRows(db, checkpoint.cursor, to); + const digest = createCanonicalAuthorityDigestWindow(); let state: SqliteStateCursor = checkpoint; const transactions: AuthorityStoreCommittedTransaction[] = []; for (const raw of rows) { const row = this.decodeCommitRow(raw); const verified = this.verifyCommitRow(row, identity, - row.cursor === checkpoint.cursor ? {kind: "sealed", state} : {kind: "predecessor", state}); + row.cursor === checkpoint.cursor ? {kind: "sealed", state} : {kind: "predecessor", state}, digest); state = verified.state; if (row.cursor >= from) transactions.push(verified.transaction); } @@ -568,6 +577,7 @@ export class SqliteAuthorityStore implements AuthorityStore { return {schema_version: "loopx_sqlite_authority_history_audit_v0", status: "verified", commits: 0, checkpoints: 0}; } const counted = db.prepare("SELECT COUNT(*) AS count FROM checkpoints").get(); + const digest = createCanonicalAuthorityDigestWindow(); let state: SqliteStateCursor | null = null; let cursor = 0n; let commits = 0; @@ -581,10 +591,10 @@ export class SqliteAuthorityStore implements AuthorityStore { // delta can never diverge from the history it claims to extend. const predecessor: SqliteStateCursor = state ?? {cursor: 0n, projection: {}, digest: ""}; const replayed: SqliteStateCursor = this.verifyCommitRow(row, identity, - {kind: "predecessor", state: predecessor}).state; + {kind: "predecessor", state: predecessor}, digest).state; if (isAuthorityStateCheckpoint(row.cursor)) { const checkpoint = this.loadCheckpoint(db, row.cursor); - const sealed = this.verifyCommitRow(row, identity, {kind: "sealed", state: checkpoint}).state; + const sealed = this.verifyCommitRow(row, identity, {kind: "sealed", state: checkpoint}, digest).state; if (sealed.digest !== replayed.digest || !canonicalAuthorityBytes(sealed.projection).equals(canonicalAuthorityBytes(replayed.projection))) { protocol("SQLite authority checkpoint does not match retained history"); @@ -653,8 +663,9 @@ export function commitDigest( projection: JsonObject, events: readonly JsonObject[], receipts: readonly JsonObject[], + digest: CanonicalAuthorityDigest = canonicalAuthoritySha256, ): string { - return canonicalAuthoritySha256({ + return digest({ expected_provider_revision: cursor === 1n ? null : `${identity}:${cursor - 1n}`, operation_id: operationId, next_projection: projection, events, receipts, }); From 918d3deaf4435358188c860f5cfad2d559027187 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Wed, 23 Sep 2026 00:11:26 -0700 Subject: [PATCH 02/10] test(sqlite): prove cached digest parity and corruption rejection Signed-off-by: Lihua <1017343802@qq.com> --- .../authority_store_digest_window.test.ts | 43 +++++++++++++++++ .../sqlite_authority_store.test.ts | 46 +++++++++++++++++++ 2 files changed, 89 insertions(+) create mode 100644 tests/control_plane_ts/authority_store_digest_window.test.ts diff --git a/tests/control_plane_ts/authority_store_digest_window.test.ts b/tests/control_plane_ts/authority_store_digest_window.test.ts new file mode 100644 index 0000000000..6dc275f9e8 --- /dev/null +++ b/tests/control_plane_ts/authority_store_digest_window.test.ts @@ -0,0 +1,43 @@ +import assert from "node:assert/strict"; +import {performance} from "node:perf_hooks"; +import test from "node:test"; + +import {canonicalAuthoritySha256, createCanonicalAuthorityDigestWindow} from + "../../loopx/control_plane/coordination/authority_store_codec.ts"; + +test("window digest keeps the exact canonical v0 hash across JSON shapes and cache eviction", () => { + const digest = createCanonicalAuthorityDigestWindow(); + const sparse: unknown[] = []; + sparse.length = 3; + sparse[1] = "middle"; + const fixtures: unknown[] = [null, "quote\"slash\\\n", -0, 1.25e30, true, sparse, + JSON.parse('{"": {"marker": "empty"}, "__proto__": {"marker": "proto"}}'), + JSON.parse('{"10":1,"2":2,"Ω":"🙂","nested":{"b":1,"a":2}}')]; + for (const value of fixtures) assert.equal(digest(value), canonicalAuthoritySha256(value)); + for (let index = 0; index < 8; index++) { + const value = {payload: String(index).repeat(1024 * 1024), ordinal: index}; + assert.equal(digest(value), canonicalAuthoritySha256(value)); + } + assert.throws(() => digest({invalid: Number.NaN})); +}); + +test("one window reuses stable large JSON strings across many distinct proofs", () => { + const padding = "p".repeat(1024 * 1024); + const inputs = Array.from({length: 100}, (_, index) => ({padding, ordinal: index})); + const digest = createCanonicalAuthorityDigestWindow(); + const baseline = inputs.map(value => canonicalAuthoritySha256(value)); + assert.deepEqual(inputs.map(value => digest(value)), baseline); + const duration = (hash: (value: unknown) => string) => { + const samples: number[] = []; + for (let trial = 0; trial < 3; trial++) { + const start = performance.now(); + for (const value of inputs) hash(value); + samples.push(performance.now() - start); + } + return samples.sort((left, right) => left - right)[1]!; + }; + const legacyMs = duration(canonicalAuthoritySha256); + const cachedMs = duration(digest); + assert.ok(cachedMs * 2 < legacyMs, + `window digest ${cachedMs.toFixed(1)}ms must beat repeated canonical encoding ${legacyMs.toFixed(1)}ms`); +}); diff --git a/tests/control_plane_ts/sqlite_authority_store.test.ts b/tests/control_plane_ts/sqlite_authority_store.test.ts index 9631010443..8261c0b6a3 100644 --- a/tests/control_plane_ts/sqlite_authority_store.test.ts +++ b/tests/control_plane_ts/sqlite_authority_store.test.ts @@ -62,6 +62,52 @@ test("SQLite commits and reads back every JSON object key", {timeout: 30000}, as assert.equal(({} as Record).marker, undefined); }); +test("SQLite large unchanged projections retain independent receipts and reject a forged digest chain", async t => { + const {store} = await fixture(t); + const projection = {capacity_padding: "p".repeat(1024 * 1024), marker: "constant"}; + let revision: string | null = null; + for (let index = 1; index <= 3; index++) { + const result = await store.commitAuthority({expected_provider_revision: revision, + operation_id: `large-${index}`, next_projection: projection, + events: [{index}], receipts: [{operation_id: `large-${index}`, index}]}); + assert.equal(result.status, "applied"); + if (result.status !== "applied") return; + revision = result.provider_revision; + } + const found = await store.readReceipt("large-3"); + assert.equal(found.status, "found"); + if (found.status === "found") assert.equal(found.receipts[0]?.index, 3); + const page = await store.scanCommitted(null, 3); + assert.equal(page.status, "page"); + if (page.status === "page") { + assert.deepEqual(page.transactions.map(row => row.operation_id), ["large-1", "large-2", "large-3"]); + (page.transactions[0]!.projection as {marker: string}).marker = "edited only in returned data"; + assert.equal((page.transactions[1]!.projection as {marker: string}).marker, "constant"); + } + const {DatabaseSync} = createRequire(import.meta.url)("node:sqlite"); + const db = new DatabaseSync(store.path); + try { + const forged = "0".repeat(64); + db.prepare("UPDATE commits SET state_digest=? WHERE cursor=2").run(forged); + db.prepare("UPDATE commits SET parent_state_digest=? WHERE cursor=3").run(forged); + } finally { db.close(); } + // The current head can still load; a historical read must prove the empty + // delta's claimed state digest rather than trusting the forged chain. + assert.equal((await store.loadAuthority()).status, "loaded"); + for (const result of [await store.readReceipt("large-3"), await store.scanCommitted(null, 3)]) { + assert.equal(result.status, "failed"); + if (result.status === "failed") assert.equal(result.reason_code, "provider_protocol_violation"); + } +}); + +test("SQLite first empty projection still derives its root digest", async t => { + const {store} = await fixture(t); + const committed = await store.commitAuthority({expected_provider_revision: null, + operation_id: "empty-root", next_projection: {}, events: [], receipts: [{operation_id: "empty-root"}]}); + assert.equal(committed.status, "applied"); + assert.equal((await store.readReceipt("empty-root")).status, "found"); +}); + test("SQLite head continuity is independent of retained history", {timeout: 30000}, async t => { const {store} = await fixture(t); assert.equal((await store.storeIdentity()).status, "available"); From ec80b39747b2df6b63501bccca41724a99e41ed7 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:02:43 +0800 Subject: [PATCH 03/10] perf(sqlite): preserve native encoding for ordinary record proofs Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/authority_store_codec.ts | 45 ++++++++++++++----- .../authority_store_digest_window.test.ts | 36 +++++++++------ 2 files changed, 55 insertions(+), 26 deletions(-) diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index 70c2ddad8e..6a589d1038 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -142,37 +142,58 @@ export function createCanonicalAuthorityDigestWindow(): CanonicalAuthorityDigest cachedBytes += bytes.byteLength; return bytes; }; + const hasLongString = (part: unknown): boolean => { + if (typeof part === "string") return part.length >= 1024; + if (part === null || typeof part !== "object") return false; + if (Array.isArray(part)) return part.some(hasLongString); + return Object.keys(part).some(key => key.length >= 1024 || hasLongString((part as JsonObject)[key])); + }; return (value: unknown): string => { const canonical = canonicalAuthorityJson(value); const hash = createHash("sha256"); + // Ordinary record-rich states have nothing reusable in this cache. Keep + // their native JSON encoder instead of paying per-token JS traversal. + if (!hasLongString(canonical)) return hash.update(JSON.stringify(canonical)).digest("hex"); + // Batch small tokens: calling the native hash binding per punctuation/key + // regresses ordinary many-field Todo projections without large strings. + let pending = ""; + const flush = (): void => { + if (pending) { hash.update(pending); pending = ""; } + }; + const append = (chunk: string | Buffer): void => { + if (typeof chunk !== "string") { flush(); hash.update(chunk); return; } + pending += chunk; + if (pending.length >= 8192) flush(); + }; const visit = (part: unknown): void => { - if (part === null) { hash.update("null"); return; } - if (typeof part === "string") { hash.update(stringBytes(part)); return; } + if (part === null) { append("null"); return; } + if (typeof part === "string") { append(stringBytes(part)); return; } if (typeof part === "number" || typeof part === "boolean") { - hash.update(JSON.stringify(part)); return; + append(JSON.stringify(part)); return; } if (Array.isArray(part)) { - hash.update("["); + append("["); for (let index = 0; index < part.length; index++) { - if (index) hash.update(","); + if (index) append(","); if (Object.hasOwn(part, index)) visit(part[index]); - else hash.update("null"); // JSON.stringify represents sparse slots as null. + else append("null"); // JSON.stringify represents sparse slots as null. } - hash.update("]"); + append("]"); return; } - hash.update("{"); + append("{"); let first = true; for (const key of Object.keys(part as JsonObject)) { - if (!first) hash.update(","); + if (!first) append(","); first = false; - hash.update(stringBytes(key)); - hash.update(":"); + append(stringBytes(key)); + append(":"); visit((part as JsonObject)[key]); } - hash.update("}"); + append("}"); }; visit(canonical); + flush(); return hash.digest("hex"); }; } diff --git a/tests/control_plane_ts/authority_store_digest_window.test.ts b/tests/control_plane_ts/authority_store_digest_window.test.ts index 6dc275f9e8..fcf45639a1 100644 --- a/tests/control_plane_ts/authority_store_digest_window.test.ts +++ b/tests/control_plane_ts/authority_store_digest_window.test.ts @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import {performance} from "node:perf_hooks"; +import {createHash} from "node:crypto"; import test from "node:test"; import {canonicalAuthoritySha256, createCanonicalAuthorityDigestWindow} from @@ -27,17 +27,25 @@ test("one window reuses stable large JSON strings across many distinct proofs", const digest = createCanonicalAuthorityDigestWindow(); const baseline = inputs.map(value => canonicalAuthoritySha256(value)); assert.deepEqual(inputs.map(value => digest(value)), baseline); - const duration = (hash: (value: unknown) => string) => { - const samples: number[] = []; - for (let trial = 0; trial < 3; trial++) { - const start = performance.now(); - for (const value of inputs) hash(value); - samples.push(performance.now() - start); - } - return samples.sort((left, right) => left - right)[1]!; - }; - const legacyMs = duration(canonicalAuthoritySha256); - const cachedMs = duration(digest); - assert.ok(cachedMs * 2 < legacyMs, - `window digest ${cachedMs.toFixed(1)}ms must beat repeated canonical encoding ${legacyMs.toFixed(1)}ms`); + // Relative speed belongs in the explicit experiment, not a timing-sensitive + // correctness test running beside unrelated CI workers. +}); + +test("window digest preserves strict input validation and exact UTF-8 bytes", () => { + const digest = createCanonicalAuthorityDigestWindow(); + const cycle: unknown[] = []; cycle.push(cycle); + for (const value of [undefined, NaN, Infinity, 1n, new Date(), {bad: undefined}, cycle]) { + assert.throws(() => digest(value)); + } + // Literal canonical bytes are independent of the implementation under test. + const plain = {b: [true, null, -0], a: "\u4e2d"}; + assert.equal(digest(plain), createHash("sha256").update('{"a":"中","b":[true,null,0]}').digest("hex")); + const strings = ["\ud800".repeat(1200), '"\\\n'.repeat(1200), "🙂".repeat(600), "x".repeat(3 * 1024 ** 2)]; + for (const value of strings) { + assert.equal(digest(value), createHash("sha256").update(JSON.stringify(value)).digest("hex")); + } + const changing = {padding: "p".repeat(4096), metadata: {attempt: 1}}; + const before = digest(changing); changing.metadata.attempt = 2; + assert.notEqual(digest(changing), before); + assert.equal(digest(changing), canonicalAuthoritySha256(changing)); }); From 0c6f6170aeba23927248c3ea11191300cfd22c2a Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:02:43 +0800 Subject: [PATCH 04/10] test(authority): add matched disposable local provider comparison Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/local-provider-comparison.ts | 148 ++++++++++++++++++ 1 file changed, 148 insertions(+) create mode 100644 examples/coordination/local-provider-comparison.ts diff --git a/examples/coordination/local-provider-comparison.ts b/examples/coordination/local-provider-comparison.ts new file mode 100644 index 0000000000..e7fe33e431 --- /dev/null +++ b/examples/coordination/local-provider-comparison.ts @@ -0,0 +1,148 @@ +/** Matched, disposable local-provider experiment; never selects a live provider. */ +import assert from "node:assert/strict"; +import {createHash} from "node:crypto"; +import {spawnSync} from "node:child_process"; +import {mkdtempSync, readFileSync, readdirSync, rmSync, statfsSync, statSync, writeFileSync} from "node:fs"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import {performance} from "node:perf_hooks"; +import {parseArgs} from "node:util"; +import {fileURLToPath} from "node:url"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {sqliteRuntimeIdentity} from "../../loopx/control_plane/coordination/sqlite_runtime.ts"; +import {productionScaleCoordinationFixture, productionScaleHistoryProjection, productionScaleObservationStep} from + "../../tests/control_plane_ts/production_scale_coordination_fixture.ts"; +import {authorityProjectionFixture} from "../../tests/control_plane_ts/authority_projection_fixture.ts"; +import {latency} from "./sqlite-capacity-report.ts"; + +const {values} = parseArgs({options: {provider: {type: "string"}, workload: {type: "string", default: "mixed"}, + commits: {type: "string", default: "128"}, samples: {type: "string", default: "20"}, + output: {type: "string"}, "cold-root": {type: "string"}}}); +assert(values.provider === "file" || values.provider === "sqlite", "provider must be file or sqlite"); +assert(["mixed", "full", "fixed-64k", "changing-1m"].includes(values.workload!), "unknown workload"); +const count = Number(values.commits), samples = Number(values.samples); +assert(Number.isSafeInteger(count) && count >= 100 && count <= 2048, "commits must be 100..2048; use D2 runner for capacity"); +assert(Number.isSafeInteger(samples) && samples >= 3 && samples <= 100, "samples must be 3..100"); +const goal = "provider-comparison"; +const openStore = (root: string): AuthorityStore => values.provider === "file" + ? new FileAuthorityStore(root, goal) : new SqliteAuthorityStore(root, goal); +if (values["cold-root"]) { + const result = await openStore(values["cold-root"]).loadAuthority(); + assert.equal(result.status, "loaded"); + process.stdout.write(JSON.stringify(result)); +} else { + const space = statfsSync(tmpdir()); + assert(space.bavail * space.bsize > 512 * 1024 ** 2, "need 512 MiB free before bounded comparison"); + const root = mkdtempSync(join(tmpdir(), "loopx-provider-comparison-")); + try { + const report = await measure(root); + const output = JSON.stringify(report, null, 2) + "\n"; + if (values.output) writeFileSync(values.output, output); + process.stdout.write(output); + } finally { rmSync(root, {recursive: true, force: true}); } +} + +function projectionAt(index: number): JsonObject { + const base = (values.workload === "full" ? productionScaleCoordinationFixture(goal, "native") + : productionScaleHistoryProjection(goal, "native")).projection as JsonObject; + if (values.workload === "fixed-64k") { + const padding = 65536 - Buffer.byteLength(JSON.stringify({...base, padding: ""})); + assert(padding >= 0, "mixed fixture no longer fits fixed 64 KiB profile"); + return {...base, padding: "p".repeat(padding)}; + } + const step = productionScaleObservationStep(base, index); + const todos = (base.todos as JsonObject[]).map(todo => todo.todo_id === step.todo_id ? step.mutation.todo : todo); + const extras = {...base}; + for (const key of ["todos", "leases", "todo_read_model", "goal_id"]) delete extras[key]; + const projection = authorityProjectionFixture(goal, todos, base.leases as JsonObject[], "native", extras); + if (values.workload === "changing-1m") { + const padding = 1024 ** 2 - Buffer.byteLength(JSON.stringify({...projection, padding: ""})); + assert(padding >= 0); projection.padding = "p".repeat(padding); + } + return projection; +} + +function sourceIdentity() { + const git = (args: string[]) => { + const result = spawnSync("git", args, {encoding: "utf8"}); + assert.equal(result.status, 0, result.stderr); return result.stdout; + }; + return {revision: git(["rev-parse", "HEAD"]).trim(), + tracked_diff_sha256: createHash("sha256").update(git(["diff", "HEAD", "--", "loopx", "tests/control_plane_ts"])).digest("hex"), + runner_sha256: createHash("sha256").update(readFileSync(fileURLToPath(import.meta.url))).digest("hex")}; +} + +async function measure(root: string) { + const source = sourceIdentity(); + const store = openStore(root), commits: number[] = [], loads: number[] = [], reads: number[] = [], scans: number[] = []; + // Fixture construction is outside timings. Each provider receives identical full records. + const projections = Array.from({length: count}, (_, index) => projectionAt(index)); + let revision: string | null = null; + let first: AuthorityStoreCommit | undefined; + let filePublicationBytes = 0; + const fileBytes = (): number => readdirSync(root).reduce((sum, name) => sum + statSync(join(root, name)).size, 0); + const timed = async (action: () => Promise, into: number[]): Promise => { + const start = performance.now(); const result = await action(); into.push(performance.now() - start); return result; + }; + const start = performance.now(); + for (let index = 0; index < count; index++) { + const input: AuthorityStoreCommit = {expected_provider_revision: revision, operation_id: `op-${index}`, + next_projection: projections[index]!, events: [{kind: "observation", index}], + receipts: [{operation_id: `op-${index}`, index, metadata: {checked: true, labels: ["synthetic", "保留"]}}]}; + const result = await timed(() => store.commitAuthority(input), commits); + assert.equal(result.status, "applied"); if (result.status !== "applied") throw new Error("commit rejected"); + revision = result.provider_revision; + if (index === 0) first = input; + if (values.provider === "file") filePublicationBytes += statSync((store as FileAuthorityStore).path).size; + if ((index + 1) % 128 === 0) process.stderr.write(`${values.provider} ${values.workload}: ${index + 1}/${count}\n`); + } + const fillMs = performance.now() - start; + for (let index = 0; index < samples; index++) { + const head = await timed(() => store.loadAuthority(), loads); + assert.equal(head.status, "loaded"); if (head.status === "loaded") assert.deepEqual(head.head, projections.at(-1)); + const receiptIndex = Math.floor(index * (count - 1) / (samples - 1)); + const receipt = await timed(() => store.readReceipt(`op-${receiptIndex}`), reads); + assert.equal(receipt.status, "found"); + const page = await timed(() => store.scanCommitted(String(count - 100), 100), scans); + assert.equal(page.status, "page"); + if (page.status === "page") { + assert.equal(page.transactions.length, 100); + for (const [offset, row] of page.transactions.entries()) { + const ordinal = count - 100 + offset; + assert.equal(row.operation_id, `op-${ordinal}`); + assert.deepEqual(row.projection, projections[ordinal]); + assert.deepEqual(row.receipts, [{operation_id: `op-${ordinal}`, index: ordinal, + metadata: {checked: true, labels: ["synthetic", "保留"]}}]); + } + } + } + // A separately opened process verifies the full head, not merely a successful exit code. + const cold: number[] = []; + for (let index = 0; index < Math.min(samples, 5); index++) { + const started = performance.now(); + const child = spawnSync(process.execPath, ["--no-warnings", "--experimental-strip-types", "--experimental-sqlite", + fileURLToPath(import.meta.url), "--provider", values.provider!, "--cold-root", root], + {encoding: "utf8", timeout: 120000, maxBuffer: 4 * 1024 ** 2}); + cold.push(performance.now() - started); assert.equal(child.status, 0, child.stderr); + assert.deepEqual(JSON.parse(child.stdout).head, projections.at(-1)); + } + const reopened = openStore(root), replay = await reopened.commitAuthority(first!); + assert.equal(replay.status, "conflict"); // Current store contract reconciles via readReceipt. + const original = await reopened.readReceipt(first!.operation_id); + assert.equal(original.status, "found"); + if (original.status === "found") assert.deepEqual(original.receipts, first!.receipts); + const after = await reopened.loadAuthority(); + assert.equal(after.status, "loaded"); if (after.status === "loaded") assert.equal(after.provider_revision, revision); + assert.deepEqual(sourceIdentity(), source, "measurement source changed while running"); + return {schema_version: "loopx_local_provider_comparison_v0", provider: values.provider, workload: values.workload, + source, node: process.version, sqlite: sqliteRuntimeIdentity(), platform: process.platform, arch: process.arch, + commits: count, projection_json_bytes: Buffer.byteLength(JSON.stringify(projections.at(-1))), fill_ms: fillMs, + commit_first_100: latency(commits.slice(0, 100)), commit_last_100: latency(commits.slice(-100)), + warm_head: latency(loads), historical_receipt: latency(reads), scan_100: latency(scans), cold_process_head: latency(cold), + final_store_bytes: fileBytes(), file_document_publication_bytes: values.provider === "file" ? filePublicationBytes : null, + complete_record_and_receipt_checks: "passed", original_receipt_recovery_after_reopen: "passed", + limits: "bounded sequential store experiment; cold process includes module loading, not cold OS cache; no CLI, concurrent writers, crash, soak or formal D2 qualification; File publication bytes are application bytes, not physical writes; SQLite WAL traffic is measured by the separate capacity runner"}; +} From e8193ce838d961183072e40b3f8b2221f91bbc9e Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:05:54 +0800 Subject: [PATCH 05/10] test(authority): bound comparison fixture memory independently of history Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/local-provider-comparison.ts | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/examples/coordination/local-provider-comparison.ts b/examples/coordination/local-provider-comparison.ts index e7fe33e431..17db88d862 100644 --- a/examples/coordination/local-provider-comparison.ts +++ b/examples/coordination/local-provider-comparison.ts @@ -27,6 +27,7 @@ const count = Number(values.commits), samples = Number(values.samples); assert(Number.isSafeInteger(count) && count >= 100 && count <= 2048, "commits must be 100..2048; use D2 runner for capacity"); assert(Number.isSafeInteger(samples) && samples >= 3 && samples <= 100, "samples must be 3..100"); const goal = "provider-comparison"; +let fixtureBase: JsonObject | undefined; const openStore = (root: string): AuthorityStore => values.provider === "file" ? new FileAuthorityStore(root, goal) : new SqliteAuthorityStore(root, goal); if (values["cold-root"]) { @@ -46,7 +47,7 @@ if (values["cold-root"]) { } function projectionAt(index: number): JsonObject { - const base = (values.workload === "full" ? productionScaleCoordinationFixture(goal, "native") + const base = fixtureBase ??= (values.workload === "full" ? productionScaleCoordinationFixture(goal, "native") : productionScaleHistoryProjection(goal, "native")).projection as JsonObject; if (values.workload === "fixed-64k") { const padding = 65536 - Buffer.byteLength(JSON.stringify({...base, padding: ""})); @@ -78,8 +79,9 @@ function sourceIdentity() { async function measure(root: string) { const source = sourceIdentity(); const store = openStore(root), commits: number[] = [], loads: number[] = [], reads: number[] = [], scans: number[] = []; - // Fixture construction is outside timings. Each provider receives identical full records. - const projections = Array.from({length: count}, (_, index) => projectionAt(index)); + // Generate one input at a time outside timings; retaining N full expected + // snapshots here would add artificial memory pressure to the provider test. + const finalProjection = projectionAt(count - 1); let revision: string | null = null; let first: AuthorityStoreCommit | undefined; let filePublicationBytes = 0; @@ -90,7 +92,7 @@ async function measure(root: string) { const start = performance.now(); for (let index = 0; index < count; index++) { const input: AuthorityStoreCommit = {expected_provider_revision: revision, operation_id: `op-${index}`, - next_projection: projections[index]!, events: [{kind: "observation", index}], + next_projection: projectionAt(index), events: [{kind: "observation", index}], receipts: [{operation_id: `op-${index}`, index, metadata: {checked: true, labels: ["synthetic", "保留"]}}]}; const result = await timed(() => store.commitAuthority(input), commits); assert.equal(result.status, "applied"); if (result.status !== "applied") throw new Error("commit rejected"); @@ -100,9 +102,10 @@ async function measure(root: string) { if ((index + 1) % 128 === 0) process.stderr.write(`${values.provider} ${values.workload}: ${index + 1}/${count}\n`); } const fillMs = performance.now() - start; + const postFillRss = process.memoryUsage().rss; for (let index = 0; index < samples; index++) { const head = await timed(() => store.loadAuthority(), loads); - assert.equal(head.status, "loaded"); if (head.status === "loaded") assert.deepEqual(head.head, projections.at(-1)); + assert.equal(head.status, "loaded"); if (head.status === "loaded") assert.deepEqual(head.head, finalProjection); const receiptIndex = Math.floor(index * (count - 1) / (samples - 1)); const receipt = await timed(() => store.readReceipt(`op-${receiptIndex}`), reads); assert.equal(receipt.status, "found"); @@ -113,7 +116,7 @@ async function measure(root: string) { for (const [offset, row] of page.transactions.entries()) { const ordinal = count - 100 + offset; assert.equal(row.operation_id, `op-${ordinal}`); - assert.deepEqual(row.projection, projections[ordinal]); + assert.deepEqual(row.projection, projectionAt(ordinal)); assert.deepEqual(row.receipts, [{operation_id: `op-${ordinal}`, index: ordinal, metadata: {checked: true, labels: ["synthetic", "保留"]}}]); } @@ -127,7 +130,7 @@ async function measure(root: string) { fileURLToPath(import.meta.url), "--provider", values.provider!, "--cold-root", root], {encoding: "utf8", timeout: 120000, maxBuffer: 4 * 1024 ** 2}); cold.push(performance.now() - started); assert.equal(child.status, 0, child.stderr); - assert.deepEqual(JSON.parse(child.stdout).head, projections.at(-1)); + assert.deepEqual(JSON.parse(child.stdout).head, finalProjection); } const reopened = openStore(root), replay = await reopened.commitAuthority(first!); assert.equal(replay.status, "conflict"); // Current store contract reconciles via readReceipt. @@ -139,10 +142,11 @@ async function measure(root: string) { assert.deepEqual(sourceIdentity(), source, "measurement source changed while running"); return {schema_version: "loopx_local_provider_comparison_v0", provider: values.provider, workload: values.workload, source, node: process.version, sqlite: sqliteRuntimeIdentity(), platform: process.platform, arch: process.arch, - commits: count, projection_json_bytes: Buffer.byteLength(JSON.stringify(projections.at(-1))), fill_ms: fillMs, + commits: count, projection_json_bytes: Buffer.byteLength(JSON.stringify(finalProjection)), fill_ms: fillMs, commit_first_100: latency(commits.slice(0, 100)), commit_last_100: latency(commits.slice(-100)), warm_head: latency(loads), historical_receipt: latency(reads), scan_100: latency(scans), cold_process_head: latency(cold), + post_fill_rss_bytes: postFillRss, final_store_bytes: fileBytes(), file_document_publication_bytes: values.provider === "file" ? filePublicationBytes : null, complete_record_and_receipt_checks: "passed", original_receipt_recovery_after_reopen: "passed", - limits: "bounded sequential store experiment; cold process includes module loading, not cold OS cache; no CLI, concurrent writers, crash, soak or formal D2 qualification; File publication bytes are application bytes, not physical writes; SQLite WAL traffic is measured by the separate capacity runner"}; + limits: "bounded sequential store experiment; cold process includes module loading, not cold OS cache; RSS includes fixture and verification allocations, not a steady-state qualification; no CLI, concurrent writers, crash, soak or formal D2 qualification; File publication bytes are application bytes, not physical writes; SQLite WAL traffic is measured by the separate capacity runner"}; } From 96ce9efb35880edc998c7c9a14ba555d22bca32b Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:17:09 +0800 Subject: [PATCH 06/10] docs(authority): ground local default choice in matched provider measurements Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- ...shared-goal-authority-state-provider-v0.md | 32 ++++--- ...-goal-authority-state-provider-v0.zh-CN.md | 24 +++-- docs/reference/sqlite-authority-store.md | 89 ++++++++++++++++++- 3 files changed, 126 insertions(+), 19 deletions(-) diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index e47af23098..544878ea0a 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -1205,15 +1205,22 @@ Keep live-state size fixed when isolating history growth, then grow live state separately. No goal-wide unbounded list of completed Todos or receipts may be hidden inside the supposedly fixed live projection. -For the current `FileAuthorityStore`, a fixed projection of P bytes retained in -each of N transactions costs approximately P*N final history bytes and -P*N*(N+1)/2 cumulative document-publication bytes, before head, event, receipt, -and envelope overhead. Normal reads also decode and validate the full chain. -With P=15 KiB, the renewal-only case gives about **534 GiB** of cumulative -publication at day 10 and **4.69 TiB** at day 30. The former 380 MiB estimate was -only N*P at day 30, not the cumulative rewrite of retained projections. These -are analytical payload estimates, not physical SSD writes or measured latency; -growing receipt indexes inside every projection can make the model worse. +The original full-projection File journal retained approximately P*N payload +bytes and republished approximately P*N*(N+1)/2 bytes across N commits. That +historical model must not be applied to the current checkpoint/delta format: +#5102 retired that layout from ordinary reads and writes. + +The current File provider retains a checkpoint every 64 commits plus deltas, +events and original receipts in one envelope. Its approximate retained bytes +are `H(N) = ceil(N/64)*P + sum(delta/event/receipt/metadata bytes)`, before the +live head and envelope overhead. Each commit still durably replaces the whole +envelope, so cumulative application publication is `sum(H(n))`. Warm reads +read/hash the envelope and may reuse its verified view; cold reads reconstruct +and verify the history. SQLite instead updates transactional indexed rows and +bounded checkpoint windows. These mechanisms motivate a matched experiment; +neither a formula nor a cache hit establishes a short-term default choice. +Report application publication separately from physical disk writes, and +compare current code on equal state, history, durability and cold/warm workload. #### Preferred local direction and compatibility boundary @@ -3201,7 +3208,12 @@ Qualify **one** long-lived local default profile. SQLite is the current D2 candidate; File remains the real reference/explicit profile and migration rehearsal backend. Do not publish two ambiguous defaults, declare the current File history layout long-horizon-qualified, or silently fall back from a -selected SQLite store. The final profile decision must cite its D2 evidence. +selected SQLite store. Release activation must cite its D2 evidence. The September 27 matched +short-history experiment also selects SQLite as the **short-term default +implementation target**: writes/head/restart beat current checkpoint/delta +File, while File retains faster warm history reads. Large-state receipt/scan +budgets remain unmet, so this is not permission to enable the default now. +[Measurements, reproduction and D2/D3/L9 dependencies](../../reference/sqlite-authority-store.md#short-term-default-decision-and-matched-experiment). PostgreSQL shares the TS semantic contracts but has independent service, tenant, restore and capacity qualification; its deployment must not delay the local profile's work. diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index c44ebd5d8a..736aff15fd 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -944,12 +944,17 @@ provider 已通过;经评审切换前,已交付规则仍是 `retain_all_v0` 一个模型 Turn 可能产生多次 commit。验证历史增长时固定 live state,之后单独增加 live state。不能把所有已完成 Todo 或 receipt 的无界列表藏在所谓固定的 live projection。 -当前 `FileAuthorityStore` 每笔保留 P 字节 projection,N 笔约产生 P*N 最终历史字节, -累计文档发布量约 P*N*(N+1)/2,尚未计 head、event、receipt 和 envelope。普通读取还会 -解码、验证完整链。P=15 KiB 时,仅 renew 的例子在第 10 天累计发布约 **534 GiB**, -第 30 天约 **4.69 TiB**。旧文中的 380 MiB 只算了第 30 天的 N*P,并非保留历次 -projection 后的累计重写。这是 payload 解析估算,不是 SSD 物理写入或实测延迟; -若每个 projection 自身还包含不断增长的 receipt index,成本可能更高。 +原先全量 projection 的 File journal 约保留 P*N 字节,并在 N 笔提交中累计发布 +P*N*(N+1)/2 字节。这个历史模型不能用于当前 checkpoint/delta 格式:#5102 +已从普通读写路径退役旧布局。 + +当前 File 每 64 笔保存 checkpoint,其余保存 delta、event 和原始 receipt,仍放在一个 +完整 envelope 中。保留量近似为 `H(N) = ceil(N/64)*P + sum(delta/event/receipt/metadata 字节)`, +另加 live head 和 envelope 开销;每次提交仍完整替换该 envelope,因此累计应用发布量为 +`sum(H(n))`。热读会读取/散列 envelope 并可能复用已验证视图,冷读需要重建和验证历史。 +SQLite 则更新事务化索引行并验证有界 checkpoint 窗口。这些机制是对照实验的依据, +不是短期默认选择的结论。必须在当前代码、相同状态/历史/持久性及冷热负载下实测, +并区分应用发布字节与物理磁盘写入。 #### 本地优先方向与兼容边界 @@ -2475,8 +2480,11 @@ provider 确认;权威空集合不回退到陈旧 Markdown。Legacy 与预览 长程默认应选定**一个**合格本地 profile。SQLite 是当前 D2 候选;File 保留为真实 对照、显式可选 profile 和迁移演练后端。不能发布两个含混的默认项,不能把现有 File -历史布局直接称为长程合格,也不能从选定 SQLite 静默回退。最终选择必须引用 D2 -证据。PostgreSQL 复用 TS 语义合同,但 service、tenant、restore 和 capacity 单独 +历史布局直接称为长程合格,也不能从选定 SQLite 静默回退。发布启用必须引用 D2 +证据。9 月 27 日同负载短历史实验也选择 SQLite 作为**短期默认实现目标**:写入、 +当前 head 与重启快于现有 checkpoint/delta File,而 File 的热历史读取仍更快。 +大状态 receipt/scan 仍未达预算,因此不是立即启用默认的许可。 +[实测、复现命令及 D2/D3/L9 依赖](../../reference/sqlite-authority-store.md#short-term-default-decision-and-matched-experiment)。PostgreSQL 复用 TS 语义合同,但 service、tenant、restore 和 capacity 单独 资格化;其部署不阻塞本地路线。 核对基线:#4286(命令回执/归档)、#4289(typed 工作/归属 intent)、#4292 diff --git a/docs/reference/sqlite-authority-store.md b/docs/reference/sqlite-authority-store.md index 1452f72795..c7e0661440 100644 --- a/docs/reference/sqlite-authority-store.md +++ b/docs/reference/sqlite-authority-store.md @@ -1,11 +1,98 @@ # SQLite authority provider SQLite is an **opt-in local conformance candidate**, behind the existing -TypeScript `AuthorityStore` interface. File remains the default. This slice +TypeScript `AuthorityStore` interface. File remains the canonical-provider +fallback when no selector is present; this does not make every new Goal +canonical or migrate an existing legacy Goal. This slice does not promote a goal, run a live cutover, enable cross-host writes, or qualify ten elapsed days of operation. It does provide the explicit version-1 to version-2 database migration described below. +## Short-term default decision and matched experiment + +The September 27 decision is to target **SQLite for the next qualified local +new-Goal default**, rather than first defaulting to File and moving again. +This is an implementation direction, **not default activation or completed D2 +qualification**. File remains an explicit provider, a conformance reference and +an export/recovery destination. Existing Goals retain their selected authority; +a rejected SQLite runtime must never silently open File instead. + +The decision uses current File checkpoint/delta storage, not its retired +full-projection-per-commit layout. On macOS arm64, Node 22.22.3 / SQLite 3.51.3, +measurement source `e8193ce83` produced the following p95 milliseconds. Arms ran +sequentially on one host; each uses 20 warm read samples, the final 100 writes, +and five fresh-process head reads. Cold-process timing includes module loading +and does not clear the OS page cache. These bounded observations are not a +population estimate or a formal capacity/soak result. + +| Projection / commits | Provider | Write | Warm head | Historical receipt | Scan 100 | Cold process head | +| --- | --- | ---: | ---: | ---: | ---: | ---: | +| Mixed 20 KiB / 128 | SQLite | 5.12 | 1.65 | 28.63 | 69.39 | 90.82 | +| Mixed 20 KiB / 128 | File | 25.23 | 6.22 | 4.62 | 32.82 | 140.14 | +| Mixed 20 KiB / 512 | SQLite | 4.76 | 1.62 | 28.36 | 66.78 | 77.69 | +| Mixed 20 KiB / 512 | File | 36.05 | 6.88 | 4.84 | 37.32 | 295.88 | +| Changing 1 MiB / 128 | SQLite | 48.50 | 9.87 | 153.46 | 377.60 | 91.62 | +| Changing 1 MiB / 128 | File | 262.20 | 34.10 | 53.80 | 137.00 | 763.60 | +| 464 Todos, 64 leases, 220 KiB / 128 | SQLite | 24.28 | 6.09 | 274.32 | 728.69 | 85.24 | +| 464 Todos, 64 leases, 220 KiB / 128 | File | 42.10 | 42.49 | 25.27 | 427.99 | 745.60 | + +The same runner on main `d23f1c87d` measured SQLite changing-1-MiB receipt/ +scan p95 at 458.25/1056.86 ms, versus 153.46/377.60 ms here. Ordinary mixed +receipt/scan was 26.00/68.29 ms versus 28.63/69.39 ms; the many-field fixture +was 270.74/690.99 ms versus 274.32/728.69 ms. The optimization benefits repeated +large strings; it does not establish a speedup for ordinary record-rich states, +and their small overhead remains visible. Do not generalize that speedup to +all Goal shapes. + +SQLite wins writes, current-state reads and process restart in these workloads. +File's verified warm history cache wins historical receipt and page reads. +The many-field fixture matters: optimizing repeated large strings alone does +not make a real Todo-rich projection cheap. Its SQLite receipt/scan timings +still exceed the RFC's 50/250 ms targets; no threshold has been increased. +These tests change a deterministic observation field in a retained full state; +they exercise storage, not the complete CLI/Turn workflow. In `changing-1m`, +padding is resized to retain exactly 1 MiB, so the large string itself can +change. It must not be reported as a stable-payload cache benchmark. + +Reproduce each arm from the same checkout and qualified Node runtime: + +```sh +node --experimental-sqlite --experimental-strip-types \ + examples/coordination/local-provider-comparison.ts \ + --provider sqlite --workload mixed --commits 128 --samples 20 +# Repeat with --provider file. Workloads: mixed, full, fixed-64k, changing-1m. +# Use --commits 512 to cross more checkpoint windows; --output writes JSON. +``` + +The runner creates and removes its own temporary store, checks complete +projections (including Todo metadata), original receipts and reopened state, +and records the source revision, runtime, runner hash and tracked source diff +hash. It does not open a selected live Goal. RSS includes fixture/checking +allocations; File publication bytes are application bytes, not physical disk +writes. Use the existing SQLite capacity runner for WAL traffic and D2 history +sizes. Keep performance experiments separate from concurrent test suites. + +Before changing release defaults, the existing owners must close these gaps: + +1. **L6 / D2:** rerun the unchanged reference capacity profiles after read-path + optimization; qualify many-field/changing-state history, crash/restore, + consumer lag, supported runtimes/platforms and the >=10-day elapsed soak. + PR #4931 contributes read-proof optimization, not a D2 pass. +2. **L8 / D3:** qualify the integrated Goal command/projection and migration + path. Reuse merged reviewed migration/retained audit (#5173), managed-host + protection (#5144) and obsolete Todo-source retirement (#5054), rather than + count them as new work. A storage benchmark does not certify long-running + execution or authorize migration of existing Goals. +3. **L9:** make new-Goal creation/onboarding choose that qualified profile, + including installed runtime admission, settings/readback and packaged entry + points. Keep explicit provider selection and reviewed backup/rollback. + +Both current local providers require Node >=22.22.3. SQLite uses built-in +`node:sqlite`: it adds no database service or external SQLite package. Its +actual embedded SQLite/finalization probe remains required, and a pre-existing +managed runtime must be restarted on the qualified executable as documented +below. PostgreSQL deployment is independent of this local default decision. + ## Placement and persistence The provider belongs to the existing shared-coordination authority boundary From dcc6836d21ccaf55ae7f8ae3bfdd87934af108e1 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:40:37 +0800 Subject: [PATCH 07/10] perf(authority): reuse owned canonical subtrees during history replay Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/authority_state_log.ts | 169 ++++++++++++------ .../coordination/authority_store_codec.ts | 86 --------- .../coordination/sqlite_authority_store.ts | 133 ++++++-------- .../authority_state_log.test.ts | 85 +++++++++ .../authority_store_digest_window.test.ts | 51 ------ 5 files changed, 252 insertions(+), 272 deletions(-) delete mode 100644 tests/control_plane_ts/authority_store_digest_window.test.ts diff --git a/loopx/control_plane/coordination/authority_state_log.ts b/loopx/control_plane/coordination/authority_state_log.ts index 5376221a37..bba780cddd 100644 --- a/loopx/control_plane/coordination/authority_state_log.ts +++ b/loopx/control_plane/coordination/authority_state_log.ts @@ -13,7 +13,8 @@ * never decides Todo, lease, quota or promotion semantics. */ import type {JsonObject} from "../effect_program.ts"; -import type {CanonicalAuthorityDigest} from "./authority_store_codec.ts"; +import {createHash} from "node:crypto"; +import type {AuthorityStoreCommit} from "./authority_store.ts"; import { AuthorityStoreProtocolError, authorityUnicodeCompare, @@ -62,10 +63,8 @@ function canonicalBytesEqual(left: unknown, right: unknown): boolean { } /** Digest of one committed projection; the anchor every reconstruction checks. */ -export function authorityStateDigest( - projection: JsonObject, digest: CanonicalAuthorityDigest = canonicalAuthoritySha256, -): string { - return digest(projection); +export function authorityStateDigest(projection: JsonObject): string { + return canonicalAuthoritySha256(projection); } /** The one checkpoint cursor that covers a positive commit cursor. */ @@ -160,10 +159,9 @@ export function applyAuthorityStateDelta( previous: JsonObject, delta: AuthorityStateDelta, ): JsonObject { - const decoded = decodeAuthorityStateDelta(delta); - const result = structuredClone(previous); - for (const operation of decoded.operations) applyAuthorityStateOperation(result, operation); - return canonicalAuthorityObject(result, "reconstructed authority state"); + const replay = new AuthorityStateReplay(previous); + replay.apply(delta); + return replay.snapshot(); } /** @@ -189,53 +187,118 @@ export function authorityStateDeltaReconstructs( } } -function applyAuthorityStateOperation(root: JsonObject, operation: AuthorityStateOperation): void { - const segments = operation.path; - if (segments.length === 0) protocol("authority state delta cannot target the root state"); - let container: unknown = root; - for (const segment of segments.slice(0, -1)) container = descend(container, segment); - const last = segments[segments.length - 1]!; - if (!isAuthorityJsonObject(container)) protocol("authority state delta path leaves the previous state"); - const present = Object.hasOwn(container, last); - if (operation.op === "splice") { - if (!present) protocol("authority state delta path leaves the previous state"); - const target = container[last]; - if (!Array.isArray(target)) protocol("authority state delta spliced a value that is not an array"); - if (operation.index < 0 || operation.remove < 0 || - operation.index + operation.remove > target.length) { - protocol("authority state delta splice is out of range"); +/** + * One read/audit's privately owned canonical state. Changed paths are copied; + * unchanged subtrees keep their identity and exact JSON encoding. Inputs and + * returned snapshots never share objects with this owner. Weak keys retain no + * historical roots; the long-string cache is separately capped at 4 MiB. + * No cached proof crosses a provider read or bypasses a row's digest check. + */ +export class AuthorityStateReplay { + #state: JsonObject; + #proofBytes: Buffer | undefined; + #encoded = new WeakMap(); + #strings = new Map(); + #stringBytes = 0; + + constructor(projection: unknown) { + this.#state = canonicalAuthorityObject(projection, "authority replay state"); + } + + apply(delta: AuthorityStateDelta): void { + let next = this.#state; + // Decode before applying and publish only after the whole batch succeeds. + // A rejected suffix cannot leave a half-applied replay frontier. + for (const operation of decodeAuthorityStateDelta(delta).operations) { + next = applyAuthorityStateOperation(next, operation, 0); } - target.splice(operation.index, operation.remove, - ...operation.insert.map(item => structuredClone(item))); - return; + if (this.#state !== next) this.#proofBytes = undefined; + this.#state = next; } - if (operation.op === "remove") { - if (!present) protocol("authority state delta removed a value that was never stored"); - delete container[last]; + + snapshot(): JsonObject { return structuredClone(this.#state); } + + /** Immutable canonical text; storage readers may parse it into independent rows. */ + canonicalJson(): string { return this.#encode(this.#state); } + + stateDigest(): string { + return createHash("sha256").update(this.#bytes()).digest("hex"); } - else setOwnJsonKey(container, last, structuredClone(operation.value)); -} -/** - * Write one decoded key as an own data property. - * - * `container[key] = value` would run the inherited `__proto__` accessor and - * replace the reconstructed object's prototype with the stored value, so a - * retained projection carrying that key would silently lose it. Decoded - * deltas are data, so every key is created the same way the canonicalizer - * creates keys, which keeps reconstruction exact for every JSON object key. - */ -function setOwnJsonKey(container: JsonObject, key: string, value: unknown): void { - Object.defineProperty(container, key, { - value, writable: true, enumerable: true, configurable: true, - }); -} + commitDigest(fields: Omit): string { + // Canonicalize the ordinary envelope, inserting our owned projection's + // exact bytes at its key. This is the existing v0 digest, not a Merkle hash + // or a new persistent proof. Integer-like keys obey native JSON ordering. + const envelope = canonicalAuthorityObject({...fields, next_projection: null}, "commit proof"); + const keys = Object.keys(envelope), split = keys.indexOf("next_projection"); + const field = (key: string) => JSON.stringify(key) + ":" + JSON.stringify(envelope[key]); + const before = keys.slice(0, split).map(field), after = keys.slice(split + 1).map(field); + return createHash("sha256") + .update("{" + (before.length ? before.join(",") + "," : "") + '"next_projection":') + .update(this.#bytes()) + .update((after.length ? "," + after.join(",") : "") + "}").digest("hex"); + } -function descend(container: unknown, segment: string): unknown { - if (!isAuthorityJsonObject(container) || !Object.hasOwn(container, segment)) { - protocol("authority state delta path leaves the previous state"); + #bytes(): Buffer { + // Both proofs consume identical UTF-8 projection bytes. Convert once for + // this frontier, avoiding another full string concatenation/UTF-8 pass. + return this.#proofBytes ??= Buffer.from(this.canonicalJson(), "utf8"); + } + + #encode(value: unknown): string { + if (value === null || typeof value !== "object") { + if (typeof value !== "string" || value.length < 1024) return JSON.stringify(value); + const known = this.#strings.get(value); + if (known !== undefined) return known; + const encoded = JSON.stringify(value), bytes = Buffer.byteLength(encoded); + if (bytes > 2 * 1024 ** 2) return encoded; + while (this.#stringBytes + bytes > 4 * 1024 ** 2) { + const oldest = this.#strings.keys().next().value!; + this.#stringBytes -= Buffer.byteLength(this.#strings.get(oldest)!); + this.#strings.delete(oldest); + } + this.#strings.set(value, encoded); this.#stringBytes += bytes; + return encoded; + } + const known = this.#encoded.get(value); + if (known !== undefined) return known; + const encoded = Array.isArray(value) + ? "[" + Array.from({length: value.length}, (_, index) => + Object.hasOwn(value, index) ? this.#encode(value[index]) : "null").join(",") + "]" + : "{" + Object.keys(value).map(key => JSON.stringify(key) + ":" + + this.#encode((value as JsonObject)[key])).join(",") + "}"; + this.#encoded.set(value, encoded); + return encoded; } - return container[segment]; +} + +/** Copy only the modified object path. Canonical children remain private and immutable. */ +function applyAuthorityStateOperation( + container: JsonObject, operation: AuthorityStateOperation, depth: number, +): JsonObject { + if (!isAuthorityJsonObject(container)) protocol("authority state delta path leaves the previous state"); + const key = operation.path[depth]!; + const next = {...container}; + const set = (value: unknown): void => { + // Treat __proto__ as data, including keys newly introduced by a delta. + Object.defineProperty(next, key, {value, writable: true, enumerable: true, configurable: true}); + }; + const present = Object.hasOwn(container, key); + if (depth + 1 < operation.path.length) { + if (!present) protocol("authority state delta path leaves the previous state"); + set(applyAuthorityStateOperation(container[key] as JsonObject, operation, depth + 1)); + } else if (operation.op === "splice") { + if (!present) protocol("authority state delta path leaves the previous state"); + const target = container[key]; + if (!Array.isArray(target)) protocol("authority state delta spliced a value that is not an array"); + if (operation.index + operation.remove > target.length) protocol("authority state delta splice is out of range"); + set(target.slice(0, operation.index).concat(operation.insert, target.slice(operation.index + operation.remove))); + } else if (operation.op === "remove") { + if (!present) protocol("authority state delta removed a value that was never stored"); + delete next[key]; + } else set(operation.value); + // Only this changed container needs reordering, not every nested Todo. + return Object.fromEntries(Object.keys(next).sort(authorityUnicodeCompare).map(name => [name, next[name]])); } /** Boundary decoder: stored or transported deltas enter as `unknown`. */ @@ -247,7 +310,7 @@ export function decodeAuthorityStateDelta(value: unknown): AuthorityStateDelta { } requireExactKeys(value, ["schema_version", "operations"], "authority state delta"); return {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, - operations: value.operations.map((operation, index) => + operations: Array.from(value.operations, (operation, index) => decodeAuthorityStateOperation(operation, index))}; } @@ -270,14 +333,14 @@ function decodeAuthorityStateOperation(value: unknown, index: number): Authority protocol(`${label} splice bounds are invalid`); } return {op: "splice", path, index: value.index as number, remove: value.remove as number, - insert: value.insert.map(item => canonicalAuthorityJson(item))}; + insert: Array.from(value.insert, item => canonicalAuthorityJson(item))}; } return protocol(`${label} op is unsupported`); } function decodeAuthorityStatePath(value: unknown, label: string): AuthorityStatePath { if (!Array.isArray(value) || value.length === 0) protocol(`${label} path is invalid`); - return value.map(segment => { + return Array.from(value, segment => { // Any string is a legal JSON key, including the empty string. if (typeof segment === "string") return segment; return protocol(`${label} path segment is invalid`); diff --git a/loopx/control_plane/coordination/authority_store_codec.ts b/loopx/control_plane/coordination/authority_store_codec.ts index 6a589d1038..06ec9e553f 100644 --- a/loopx/control_plane/coordination/authority_store_codec.ts +++ b/loopx/control_plane/coordination/authority_store_codec.ts @@ -112,92 +112,6 @@ export function canonicalAuthoritySha256(value: unknown): string { return createHash("sha256").update(canonicalAuthorityBytes(value)).digest("hex"); } -export type CanonicalAuthorityDigest = (value: unknown) => string; - -/** - * Recompute the exact canonical v0 digest while reusing encoded long strings - * within one verification window. Retained projections often share a large - * unchanged value across many deltas; serializing it for every state and - * commit proof dominates historical reads. The cache is bounded and never - * survives the caller's window, so a later read still verifies fresh bytes. - */ -export function createCanonicalAuthorityDigestWindow(): CanonicalAuthorityDigest { - const maxEntryBytes = 2 * 1024 * 1024; - const maxCachedBytes = 4 * 1024 * 1024; - const encoded = new Map(); - let cachedBytes = 0; - const stringBytes = (value: string): string | Buffer => { - if (value.length < 1024) return JSON.stringify(value); - const previous = encoded.get(value); - if (previous) return previous; - const bytes = Buffer.from(JSON.stringify(value), "utf8"); - if (bytes.byteLength > maxEntryBytes) return bytes; - while (cachedBytes + bytes.byteLength > maxCachedBytes) { - const oldest = encoded.keys().next().value; - if (oldest === undefined) break; - cachedBytes -= encoded.get(oldest)!.byteLength; - encoded.delete(oldest); - } - encoded.set(value, bytes); - cachedBytes += bytes.byteLength; - return bytes; - }; - const hasLongString = (part: unknown): boolean => { - if (typeof part === "string") return part.length >= 1024; - if (part === null || typeof part !== "object") return false; - if (Array.isArray(part)) return part.some(hasLongString); - return Object.keys(part).some(key => key.length >= 1024 || hasLongString((part as JsonObject)[key])); - }; - return (value: unknown): string => { - const canonical = canonicalAuthorityJson(value); - const hash = createHash("sha256"); - // Ordinary record-rich states have nothing reusable in this cache. Keep - // their native JSON encoder instead of paying per-token JS traversal. - if (!hasLongString(canonical)) return hash.update(JSON.stringify(canonical)).digest("hex"); - // Batch small tokens: calling the native hash binding per punctuation/key - // regresses ordinary many-field Todo projections without large strings. - let pending = ""; - const flush = (): void => { - if (pending) { hash.update(pending); pending = ""; } - }; - const append = (chunk: string | Buffer): void => { - if (typeof chunk !== "string") { flush(); hash.update(chunk); return; } - pending += chunk; - if (pending.length >= 8192) flush(); - }; - const visit = (part: unknown): void => { - if (part === null) { append("null"); return; } - if (typeof part === "string") { append(stringBytes(part)); return; } - if (typeof part === "number" || typeof part === "boolean") { - append(JSON.stringify(part)); return; - } - if (Array.isArray(part)) { - append("["); - for (let index = 0; index < part.length; index++) { - if (index) append(","); - if (Object.hasOwn(part, index)) visit(part[index]); - else append("null"); // JSON.stringify represents sparse slots as null. - } - append("]"); - return; - } - append("{"); - let first = true; - for (const key of Object.keys(part as JsonObject)) { - if (!first) append(","); - first = false; - append(stringBytes(key)); - append(":"); - visit((part as JsonObject)[key]); - } - append("}"); - }; - visit(canonical); - flush(); - return hash.digest("hex"); - }; -} - export function parseAuthorityCursor(value: string | null): bigint { if (value === null) return 0n; if (typeof value !== "string" || !/^[1-9]\d*$/.test(value)) { diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index 9a168ad859..48091d9ad5 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -10,12 +10,11 @@ import type { AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult, AuthorityStoreIdentityResult, AuthorityStoreLoadResult, AuthorityStoreReadFailure, AuthorityStoreHead, AuthorityStoreReceiptResult, AuthorityStoreReceiptBatchResult, AuthorityStoreScanResult } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityBytes, canonicalAuthorityObject, - canonicalAuthorityObjectList, canonicalAuthoritySha256, createCanonicalAuthorityDigestWindow, - normalizeAuthorityStoreCommit, requireAuthorityStoreId, - type CanonicalAuthorityDigest } from "./authority_store_codec.ts"; + canonicalAuthorityObjectList, canonicalAuthoritySha256, + normalizeAuthorityStoreCommit, requireAuthorityStoreId } from "./authority_store_codec.ts"; import { AUTHORITY_STATE_CHECKPOINT_INTERVAL, - applyAuthorityStateDelta, + AuthorityStateReplay, authorityStateCheckpointCursor, authorityStateDelta, authorityStateDeltaReconstructs, @@ -105,17 +104,10 @@ interface SqliteCommitRow { receipts: JsonObject[]; } -interface SqliteVerifiedCommit { - state: SqliteStateCursor; - transaction: AuthorityStoreCommittedTransaction; -} - -/** One bounded window of retained transactions resumed from a checkpoint. */ -interface SqliteCommitWindow { - identity: string; - checkpoint: SqliteStateCursor; - transactions: readonly AuthorityStoreCommittedTransaction[]; - state: SqliteStateCursor; +interface SqliteReplayCursor { + cursor: bigint; + digest: string; + replay: AuthorityStateReplay; } export interface SqliteAuthorityBoundedProfile { @@ -258,53 +250,39 @@ export class SqliteAuthorityStore implements AuthorityStore { private verifyCommitRow( row: SqliteCommitRow, identity: string, - base: {kind: "predecessor" | "sealed"; state: SqliteStateCursor}, - digest: CanonicalAuthorityDigest, - ): SqliteVerifiedCommit { - let projection: JsonObject; + base: {kind: "predecessor" | "sealed"; state: SqliteReplayCursor}, + ): SqliteReplayCursor { + const replay = base.state.replay; if (base.kind === "sealed") { if (row.cursor !== base.state.cursor) protocol("SQLite authority state log window is invalid"); if (row.state_digest !== base.state.digest) protocol("SQLite authority checkpoint state digest mismatch"); - projection = base.state.projection; } else { if (row.cursor !== base.state.cursor + 1n) protocol("SQLite authority state log cursor lineage is invalid"); - // The first retained commit is the root of the chain and carries no - // parent digest; every later commit must name the state it extends. const parentMatches = row.cursor === 1n - ? row.parent_state_digest === null - : row.parent_state_digest === base.state.digest; + ? row.parent_state_digest === null : row.parent_state_digest === base.state.digest; if (!parentMatches) protocol("SQLite authority state log parent lineage is invalid"); - projection = applyAuthorityStateDelta(base.state.projection, row.delta); - // A later empty delta preserves the already-verified predecessor state. - // Keep the exact commit proof below, including its events and receipts, - // but avoid rehashing an unchanged large projection just to rediscover - // the same state digest. The root has no verified predecessor digest. - const stateDigest = row.cursor > 1n && row.delta.operations.length === 0 - ? base.state.digest : authorityStateDigest(projection, digest); - if (stateDigest !== row.state_digest) { - protocol("SQLite authority state log digest mismatch"); - } + replay.apply(row.delta); + // Empty deltas retain a proven predecessor, but the root must derive its + // own proof. Every row still checks its declared state and full commit. + const digest = row.cursor > 1n && row.delta.operations.length === 0 + ? base.state.digest : replay.stateDigest(); + if (digest !== row.state_digest) protocol("SQLite authority state log digest mismatch"); } - const expected = commitDigest(identity, row.cursor, row.operation_id, projection, row.events, row.receipts, digest); + const expected = replay.commitDigest(commitFields(identity, row.cursor, row.operation_id, row.events, row.receipts)); if (row.commit_digest !== expected) protocol("SQLite committed row digest mismatch"); - return { - state: {cursor: row.cursor, projection, digest: row.state_digest}, - transaction: {cursor: row.cursor.toString(), provider_revision: `${identity}:${row.cursor}`, - operation_id: row.operation_id, projection, events: row.events, receipts: row.receipts}, - }; + return {cursor: row.cursor, digest: row.state_digest, replay}; } - private loadCheckpoint(db: DatabaseSync, cursor: bigint): SqliteStateCursor { + private loadCheckpoint(db: DatabaseSync, cursor: bigint): SqliteReplayCursor { const expected = authorityStateCheckpointCursor(cursor); const row = db.prepare( "SELECT cursor, projection, projection_digest FROM checkpoints WHERE cursor = ?", ).get(expected.toString()); if (!row) protocol("SQLite authority checkpoint is missing for its window"); - const projection = canonicalAuthorityObject(this.parseJson(row.projection, "SQLite checkpoint projection"), - "SQLite checkpoint projection"); + const replay = new AuthorityStateReplay(this.parseJson(row.projection, "SQLite checkpoint projection")); const digest = this.requireDigest(row.projection_digest, "SQLite checkpoint digest"); - if (authorityStateDigest(projection) !== digest) protocol("SQLite authority checkpoint digest mismatch"); - return {cursor: expected, projection, digest}; + if (replay.stateDigest() !== digest) protocol("SQLite authority checkpoint digest mismatch"); + return {cursor: expected, replay, digest}; } private windowRows(db: DatabaseSync, from: bigint, to: bigint): Record[] { @@ -325,21 +303,18 @@ export class SqliteAuthorityStore implements AuthorityStore { * state; history outside the requested span is verified when it is read or * when `verifyAuthorityHistory` audits the complete archive. */ - private verifiedRange(db: DatabaseSync, from: bigint, to: bigint): SqliteCommitWindow { + private verifiedRange(db: DatabaseSync, from: bigint, to: bigint, + consume: (row: SqliteCommitRow, replay: AuthorityStateReplay, identity: string) => void): void { const identity = this.identity(db); const checkpoint = this.loadCheckpoint(db, from); const rows = this.windowRows(db, checkpoint.cursor, to); - const digest = createCanonicalAuthorityDigestWindow(); - let state: SqliteStateCursor = checkpoint; - const transactions: AuthorityStoreCommittedTransaction[] = []; + let state = checkpoint; for (const raw of rows) { const row = this.decodeCommitRow(raw); - const verified = this.verifyCommitRow(row, identity, - row.cursor === checkpoint.cursor ? {kind: "sealed", state} : {kind: "predecessor", state}, digest); - state = verified.state; - if (row.cursor >= from) transactions.push(verified.transaction); + state = this.verifyCommitRow(row, identity, + row.cursor === checkpoint.cursor ? {kind: "sealed", state} : {kind: "predecessor", state}); + if (row.cursor >= from) consume(row, state.replay, identity); } - return {identity, checkpoint, transactions, state}; } /** The live head, proven without materializing retained history. */ @@ -557,11 +532,10 @@ export class SqliteAuthorityStore implements AuthorityStore { const verified = new Map(); const wanted = new Set(operationIds); for (const {from, to} of ranges.values()) { - const window = this.verifiedRange(db, from, to); - for (const row of window.transactions) if (wanted.has(row.operation_id)) { - verified.set(row.operation_id, {status: "found", cursor: row.cursor, - provider_revision: row.provider_revision, receipts: row.receipts}); - } + this.verifiedRange(db, from, to, (row, _replay, identity) => { + if (wanted.has(row.operation_id)) verified.set(row.operation_id, {status: "found", cursor: row.cursor.toString(), + provider_revision: `${identity}:${row.cursor}`, receipts: row.receipts}); + }); } const results = operationIds.map((id, index): AuthorityStoreReceiptResult => { const cursor = selected[index]; @@ -593,16 +567,12 @@ export class SqliteAuthorityStore implements AuthorityStore { // One bounded pass verifies the whole returned page: recovery starts at // the checkpoint covering the first row, and the lookahead row both // proves has_more and closes the verified span. - const window = rows.length === 0 - ? null - : this.verifiedRange(db, BigInt(String(rows[0]!.sequence)), - BigInt(String(rows[rows.length - 1]!.sequence))); - const verified = rows.map(raw => { - const cursor = BigInt(String(raw.sequence)).toString(); - const transaction = window?.transactions.find(item => item.cursor === cursor); - if (!transaction) protocol("SQLite authority scan row is not part of its retained window"); - return transaction; - }); + const verified: AuthorityStoreCommittedTransaction[] = []; + if (rows.length) this.verifiedRange(db, BigInt(String(rows[0]!.sequence)), + BigInt(String(rows[rows.length - 1]!.sequence)), (row, replay, identity) => { + verified.push({cursor: row.cursor.toString(), provider_revision: `${identity}:${row.cursor}`, + operation_id: row.operation_id, projection: JSON.parse(replay.canonicalJson()) as JsonObject, events: row.events, receipts: row.receipts}); + }); return scan.page(verified, head === null ? null : {cursor: head.state.cursor.toString(), provider_revision: head.provider_revision, head: head.state.projection}); } catch (error) { return readFailure(error); } @@ -631,8 +601,7 @@ export class SqliteAuthorityStore implements AuthorityStore { return {schema_version: "loopx_sqlite_authority_history_audit_v0", status: "verified", commits: 0, checkpoints: 0}; } const counted = db.prepare("SELECT COUNT(*) AS count FROM checkpoints").get(); - const digest = createCanonicalAuthorityDigestWindow(); - let state: SqliteStateCursor | null = null; + let state: SqliteReplayCursor = {cursor: 0n, digest: "", replay: new AuthorityStateReplay({})}; let cursor = 0n; let commits = 0; for (;;) { @@ -643,14 +612,12 @@ export class SqliteAuthorityStore implements AuthorityStore { const row = this.decodeCommitRow(raw); // The exact delta chain is proved from the empty root, so a retained // delta can never diverge from the history it claims to extend. - const predecessor: SqliteStateCursor = state ?? {cursor: 0n, projection: {}, digest: ""}; - const replayed: SqliteStateCursor = this.verifyCommitRow(row, identity, - {kind: "predecessor", state: predecessor}, digest).state; + const replayed = this.verifyCommitRow(row, identity, {kind: "predecessor", state}); if (isAuthorityStateCheckpoint(row.cursor)) { const checkpoint = this.loadCheckpoint(db, row.cursor); - const sealed = this.verifyCommitRow(row, identity, {kind: "sealed", state: checkpoint}, digest).state; + const sealed = this.verifyCommitRow(row, identity, {kind: "sealed", state: checkpoint}); if (sealed.digest !== replayed.digest || - !canonicalAuthorityBytes(sealed.projection).equals(canonicalAuthorityBytes(replayed.projection))) { + sealed.replay.canonicalJson() !== replayed.replay.canonicalJson()) { protocol("SQLite authority checkpoint does not match retained history"); } } @@ -659,7 +626,7 @@ export class SqliteAuthorityStore implements AuthorityStore { commits += 1; } } - if (state === null || state.cursor !== head.state.cursor || state.digest !== head.state.digest) { + if (state.cursor !== head.state.cursor || state.digest !== head.state.digest) { protocol("SQLite authority history does not reach the head"); } return {schema_version: "loopx_sqlite_authority_history_audit_v0", status: "verified", commits, @@ -717,10 +684,12 @@ export function commitDigest( projection: JsonObject, events: readonly JsonObject[], receipts: readonly JsonObject[], - digest: CanonicalAuthorityDigest = canonicalAuthoritySha256, ): string { - return digest({ - expected_provider_revision: cursor === 1n ? null : `${identity}:${cursor - 1n}`, - operation_id: operationId, next_projection: projection, events, receipts, - }); + return canonicalAuthoritySha256({...commitFields(identity, cursor, operationId, events, receipts), next_projection: projection}); +} + +function commitFields(identity: string, cursor: bigint, operationId: string, + events: readonly JsonObject[], receipts: readonly JsonObject[]): Omit { + return {expected_provider_revision: cursor === 1n ? null : `${identity}:${cursor - 1n}`, + operation_id: operationId, events, receipts}; } diff --git a/tests/control_plane_ts/authority_state_log.test.ts b/tests/control_plane_ts/authority_state_log.test.ts index a3cc59e2ba..979c153533 100644 --- a/tests/control_plane_ts/authority_state_log.test.ts +++ b/tests/control_plane_ts/authority_state_log.test.ts @@ -181,3 +181,88 @@ test("authority state deltas keep every JSON key the stored projection could car assert.equal(text(created), text(JSON.parse(legacyText))); assert.equal(Object.getPrototypeOf(created), Object.prototype); }); + +test("delta batches preserve input ownership and apply operations in order", () => { + const previous = {nested: {metadata: {attempt: 1}}, list: [{value: 1}, {value: 2}]}; + const operations = [ + {op: "set", path: ["nested", "metadata", "attempt"], value: 2}, + {op: "splice", path: ["list"], index: 1, remove: 1, insert: [{value: 3}]}, + {op: "set", path: ["new"], value: {branch: 1}}, + {op: "remove", path: ["new", "branch"]}, + ]; + const delta = {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations} as never; + const result = applyAuthorityStateDelta(previous, delta); + assert.deepEqual(result, {nested: {metadata: {attempt: 2}}, list: [{value: 1}, {value: 3}], new: {}}); + assert.deepEqual(previous, {nested: {metadata: {attempt: 1}}, list: [{value: 1}, {value: 2}]}); + (result.nested as {metadata: {attempt: number}}).metadata.attempt = 99; + assert.equal(previous.nested.metadata.attempt, 1); + assert.throws(() => applyAuthorityStateDelta(previous, {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, + operations: [...operations, {op: "remove", path: ["absent"]}]} as never)); + assert.equal(previous.nested.metadata.attempt, 1); + assert.deepEqual(applyAuthorityStateDelta({negative: -0, list: Array(2)}, + {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: []}), {negative: -0, list: Array(2)}); +}); + +test("private replay keeps exact proofs while copying only changed paths", async () => { + const {AuthorityStateReplay} = await import("../../loopx/control_plane/coordination/authority_state_log.ts"); + const {canonicalAuthoritySha256} = await import("../../loopx/control_plane/coordination/authority_store_codec.ts"); + const initial = {todos: [{id: "a", metadata: {attempt: 1}}, {id: "b"}], nested: {value: 1}}; + const replay = new AuthorityStateReplay(initial); + initial.todos[0]!.metadata!.attempt = 100; // Caller never owns the replay's nodes. + const before = replay.canonicalJson(); + const beforeDigest = replay.stateDigest(); + assert.equal(before, '{"nested":{"value":1},"todos":[{"id":"a","metadata":{"attempt":1}},{"id":"b"}]}'); + const insert = {id: "a", metadata: {attempt: 2}}; + replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [ + {op: "splice", path: ["todos"], index: 0, remove: 1, insert: [insert]}, + {op: "set", path: ["nested", "value"], value: 2}, + ]}); + insert.metadata.attempt = 200; + const expected = {nested: {value: 2}, todos: [{id: "a", metadata: {attempt: 2}}, {id: "b"}]}; + assert.deepEqual(replay.snapshot(), expected); + assert.notEqual(replay.stateDigest(), beforeDigest); + assert.equal(replay.stateDigest(), canonicalAuthoritySha256(expected)); + const fields = {expected_provider_revision: "store:1", operation_id: "op-2", events: [{step: 2}], receipts: [{done: true}]}; + assert.equal(replay.commitDigest(fields), canonicalAuthoritySha256({...fields, next_projection: expected})); + fields.receipts[0]!.done = false; + assert.equal(replay.commitDigest(fields), canonicalAuthoritySha256({...fields, next_projection: expected})); + const copy = replay.snapshot(); (copy.nested as {value: number}).value = 999; + const durableCopy = JSON.parse(replay.canonicalJson()); durableCopy.todos[1].id = "changed"; + assert.deepEqual(replay.snapshot(), expected); + const stable = replay.canonicalJson(); + assert.throws(() => replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [ + {op: "set", path: ["nested", "value"], value: 3}, {op: "remove", path: ["missing"]}, + ]})); + assert.equal(replay.canonicalJson(), stable, "failed batch must not advance replay"); + replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: []}); + assert.equal(replay.canonicalJson(), stable); +}); + +test("replay encoding preserves canonical Unicode, sparse values, numeric keys and strict boundaries", async () => { + const {AuthorityStateReplay} = await import("../../loopx/control_plane/coordination/authority_state_log.ts"); + const {createHash} = await import("node:crypto"); + const hash = (bytes: string) => createHash("sha256").update(bytes).digest("hex"); + const value = JSON.parse('{"10":1,"2":2,"__proto__":{"own":true},"":0,"Ω":"🙂"}'); + const replay = new AuthorityStateReplay(value); + assert.equal(replay.stateDigest(), hash('{"2":2,"10":1,"":0,"__proto__":{"own":true},"Ω":"🙂"}')); + replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [ + {op: "set", path: ["1"], value: "first"}, {op: "remove", path: ["__proto__", "own"]}, + ]}); + assert.equal(replay.canonicalJson(), '{"1":"first","2":2,"10":1,"":0,"__proto__":{},"Ω":"🙂"}'); + const sparse = new AuthorityStateReplay({items: Array(2)}); + assert.equal(sparse.canonicalJson(), '{"items":[null,null]}'); + const cycle: unknown[] = []; cycle.push(cycle); + for (const bad of [undefined, NaN, Infinity, 1n, new Date(), {bad: undefined}, {cycle}]) { + assert.throws(() => new AuthorityStateReplay(bad)); + } + for (const delta of [ + {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: Array(1)}, + {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [{op: "set", path: Array(1), value: 1}]}, + {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [{op: "splice", path: ["items"], index: 0, remove: 0, insert: Array(1)}]}, + ]) assert.throws(() => sparse.apply(delta as never)); + for (const text of ['"\\\n'.repeat(1200), "\ud800".repeat(1200), "🙂".repeat(600), "x".repeat(3 * 1024 ** 2), + ...Array.from({length: 8}, (_, i) => String(i).repeat(1024 ** 2))]) { + replay.apply({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, operations: [{op: "set", path: ["payload"], value: text}]}); + assert.equal(replay.stateDigest(), hash(canonicalAuthorityBytes(replay.snapshot()).toString("utf8"))); + } +}); diff --git a/tests/control_plane_ts/authority_store_digest_window.test.ts b/tests/control_plane_ts/authority_store_digest_window.test.ts deleted file mode 100644 index fcf45639a1..0000000000 --- a/tests/control_plane_ts/authority_store_digest_window.test.ts +++ /dev/null @@ -1,51 +0,0 @@ -import assert from "node:assert/strict"; -import {createHash} from "node:crypto"; -import test from "node:test"; - -import {canonicalAuthoritySha256, createCanonicalAuthorityDigestWindow} from - "../../loopx/control_plane/coordination/authority_store_codec.ts"; - -test("window digest keeps the exact canonical v0 hash across JSON shapes and cache eviction", () => { - const digest = createCanonicalAuthorityDigestWindow(); - const sparse: unknown[] = []; - sparse.length = 3; - sparse[1] = "middle"; - const fixtures: unknown[] = [null, "quote\"slash\\\n", -0, 1.25e30, true, sparse, - JSON.parse('{"": {"marker": "empty"}, "__proto__": {"marker": "proto"}}'), - JSON.parse('{"10":1,"2":2,"Ω":"🙂","nested":{"b":1,"a":2}}')]; - for (const value of fixtures) assert.equal(digest(value), canonicalAuthoritySha256(value)); - for (let index = 0; index < 8; index++) { - const value = {payload: String(index).repeat(1024 * 1024), ordinal: index}; - assert.equal(digest(value), canonicalAuthoritySha256(value)); - } - assert.throws(() => digest({invalid: Number.NaN})); -}); - -test("one window reuses stable large JSON strings across many distinct proofs", () => { - const padding = "p".repeat(1024 * 1024); - const inputs = Array.from({length: 100}, (_, index) => ({padding, ordinal: index})); - const digest = createCanonicalAuthorityDigestWindow(); - const baseline = inputs.map(value => canonicalAuthoritySha256(value)); - assert.deepEqual(inputs.map(value => digest(value)), baseline); - // Relative speed belongs in the explicit experiment, not a timing-sensitive - // correctness test running beside unrelated CI workers. -}); - -test("window digest preserves strict input validation and exact UTF-8 bytes", () => { - const digest = createCanonicalAuthorityDigestWindow(); - const cycle: unknown[] = []; cycle.push(cycle); - for (const value of [undefined, NaN, Infinity, 1n, new Date(), {bad: undefined}, cycle]) { - assert.throws(() => digest(value)); - } - // Literal canonical bytes are independent of the implementation under test. - const plain = {b: [true, null, -0], a: "\u4e2d"}; - assert.equal(digest(plain), createHash("sha256").update('{"a":"中","b":[true,null,0]}').digest("hex")); - const strings = ["\ud800".repeat(1200), '"\\\n'.repeat(1200), "🙂".repeat(600), "x".repeat(3 * 1024 ** 2)]; - for (const value of strings) { - assert.equal(digest(value), createHash("sha256").update(JSON.stringify(value)).digest("hex")); - } - const changing = {padding: "p".repeat(4096), metadata: {attempt: 1}}; - const before = digest(changing); changing.metadata.attempt = 2; - assert.notEqual(digest(changing), before); - assert.equal(digest(changing), canonicalAuthoritySha256(changing)); -}); From 2dcc68910fc47cb066c86147c5b71329d747d767 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:47:15 +0800 Subject: [PATCH 08/10] perf(authority): stream canonical proofs and isolate archive replay ownership Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/authority_archive_read.ts | 12 ++-- .../coordination/authority_state_log.ts | 61 +++++++++++-------- .../coordination/sqlite_authority_store.ts | 3 +- .../authority_archive.test.ts | 22 +++++++ 4 files changed, 67 insertions(+), 31 deletions(-) diff --git a/loopx/control_plane/coordination/authority_archive_read.ts b/loopx/control_plane/coordination/authority_archive_read.ts index c5c0be042d..4b39d5065b 100644 --- a/loopx/control_plane/coordination/authority_archive_read.ts +++ b/loopx/control_plane/coordination/authority_archive_read.ts @@ -7,7 +7,7 @@ import type {JsonObject} from "../effect_program.ts"; import type {AuthorityStoreCommittedTransaction} from "./authority_store.ts"; import {AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthorityObjectList, canonicalAuthoritySha256, hasExactAuthorityKeys, requireAuthorityStoreId} from "./authority_store_codec.ts"; -import {applyAuthorityStateDelta, decodeAuthorityStateDelta} from "./authority_state_log.ts"; +import {AuthorityStateReplay, decodeAuthorityStateDelta} from "./authority_state_log.ts"; export const AUTHORITY_ARCHIVE_SCHEMA = "loopx_authority_archive_v0"; const MAX_LINE_BYTES = 64 * 1024 * 1024; @@ -87,7 +87,7 @@ async function* archiveRecords(path: string): AsyncGenerator { let header: ArchiveHeader | null = null; let digest: string | null = null; let cursor = 0n; - let state: JsonObject = {}; + const replay = new AuthorityStateReplay({}); let revision: string | null = null; let sealed = false; const operations = new Set(); @@ -116,9 +116,11 @@ async function* archiveRecords(path: string): AsyncGenerator { const operation = requireAuthorityStoreId(value.operation_id, "operation id"); if (operations.has(operation)) invalid("archive operation id is duplicated"); operations.add(operation); - state = applyAuthorityStateDelta(state, decodeAuthorityStateDelta(value.delta)); + replay.apply(decodeAuthorityStateDelta(value.delta)); + // Consumers own their returned row, never the decoder's replay state. + const state = JSON.parse(replay.canonicalJson()) as JsonObject; if (state.goal_id !== header.goal_id) invalid("archive transaction belongs to another goal"); - if (canonicalAuthoritySha256(state) !== archiveHash(value.projection_sha256)) invalid("archive state reconstruction mismatch"); + if (replay.stateDigest() !== archiveHash(value.projection_sha256)) invalid("archive state reconstruction mismatch"); revision = requireAuthorityStoreId(value.provider_revision, "provider revision"); cursor += 1n; yield {kind: "transaction", transaction: {cursor: cursor.toString(), provider_revision: revision, @@ -128,7 +130,7 @@ async function* archiveRecords(path: string): AsyncGenerator { exact(value, ["kind", "cursor", "projection_sha256", "previous_sha256"]); if (value.previous_sha256 !== digest || value.cursor !== header.cursor || cursor.toString() !== header.cursor || revision !== header.provider_revision || value.projection_sha256 !== header.projection_sha256 || - canonicalAuthoritySha256(state) !== header.projection_sha256) invalid("archive seal does not cover its captured head"); + replay.stateDigest() !== header.projection_sha256) invalid("archive seal does not cover its captured head"); sealed = true; verified = summary(header, archiveHash(raw.sha256)); } else invalid("unknown archive record kind"); diff --git a/loopx/control_plane/coordination/authority_state_log.ts b/loopx/control_plane/coordination/authority_state_log.ts index bba780cddd..df299d435e 100644 --- a/loopx/control_plane/coordination/authority_state_log.ts +++ b/loopx/control_plane/coordination/authority_state_log.ts @@ -13,7 +13,7 @@ * never decides Todo, lease, quota or promotion semantics. */ import type {JsonObject} from "../effect_program.ts"; -import {createHash} from "node:crypto"; +import {createHash, type Hash} from "node:crypto"; import type {AuthorityStoreCommit} from "./authority_store.ts"; import { AuthorityStoreProtocolError, @@ -196,9 +196,9 @@ export function authorityStateDeltaReconstructs( */ export class AuthorityStateReplay { #state: JsonObject; - #proofBytes: Buffer | undefined; + #byteObjects = new WeakMap(); #encoded = new WeakMap(); - #strings = new Map(); + #strings = new Map(); #stringBytes = 0; constructor(projection: unknown) { @@ -212,7 +212,6 @@ export class AuthorityStateReplay { for (const operation of decodeAuthorityStateDelta(delta).operations) { next = applyAuthorityStateOperation(next, operation, 0); } - if (this.#state !== next) this.#proofBytes = undefined; this.#state = next; } @@ -222,7 +221,7 @@ export class AuthorityStateReplay { canonicalJson(): string { return this.#encode(this.#state); } stateDigest(): string { - return createHash("sha256").update(this.#bytes()).digest("hex"); + return this.#updateProjection(createHash("sha256")).digest("hex"); } commitDigest(fields: Omit): string { @@ -233,32 +232,44 @@ export class AuthorityStateReplay { const keys = Object.keys(envelope), split = keys.indexOf("next_projection"); const field = (key: string) => JSON.stringify(key) + ":" + JSON.stringify(envelope[key]); const before = keys.slice(0, split).map(field), after = keys.slice(split + 1).map(field); - return createHash("sha256") - .update("{" + (before.length ? before.join(",") + "," : "") + '"next_projection":') - .update(this.#bytes()) - .update((after.length ? "," + after.join(",") : "") + "}").digest("hex"); + const hash = createHash("sha256").update("{" + (before.length ? before.join(",") + "," : "") + '"next_projection":'); + return this.#updateProjection(hash).update((after.length ? "," + after.join(",") : "") + "}").digest("hex"); } - #bytes(): Buffer { - // Both proofs consume identical UTF-8 projection bytes. Convert once for - // this frontier, avoiding another full string concatenation/UTF-8 pass. - return this.#proofBytes ??= Buffer.from(this.canonicalJson(), "utf8"); + #updateProjection(hash: Hash): Hash { + hash.update("{"); + let first = true; + for (const key of Object.keys(this.#state)) { + hash.update((first ? "" : ",") + JSON.stringify(key) + ":"); first = false; + const value = this.#state[key]; + if (value !== null && typeof value === "object") { + let bytes = this.#byteObjects.get(value); + if (!bytes) { bytes = Buffer.from(this.#encode(value), "utf8"); this.#byteObjects.set(value, bytes); } + hash.update(bytes); + } else if (typeof value === "string" && value.length >= 1024) hash.update(this.#longString(value)); + else hash.update(JSON.stringify(value)); + } + return hash.update("}"); + } + + #longString(value: string): Buffer { + const known = this.#strings.get(value); + if (known) return known; + const bytes = Buffer.from(JSON.stringify(value), "utf8"); + if (bytes.byteLength > 2 * 1024 ** 2) return bytes; + while (this.#stringBytes + bytes.byteLength > 4 * 1024 ** 2) { + const oldest = this.#strings.keys().next().value!; + this.#stringBytes -= this.#strings.get(oldest)!.byteLength; + this.#strings.delete(oldest); + } + this.#strings.set(value, bytes); this.#stringBytes += bytes.byteLength; + return bytes; } #encode(value: unknown): string { if (value === null || typeof value !== "object") { - if (typeof value !== "string" || value.length < 1024) return JSON.stringify(value); - const known = this.#strings.get(value); - if (known !== undefined) return known; - const encoded = JSON.stringify(value), bytes = Buffer.byteLength(encoded); - if (bytes > 2 * 1024 ** 2) return encoded; - while (this.#stringBytes + bytes > 4 * 1024 ** 2) { - const oldest = this.#strings.keys().next().value!; - this.#stringBytes -= Buffer.byteLength(this.#strings.get(oldest)!); - this.#strings.delete(oldest); - } - this.#strings.set(value, encoded); this.#stringBytes += bytes; - return encoded; + return typeof value === "string" && value.length >= 1024 + ? this.#longString(value).toString("utf8") : JSON.stringify(value); } const known = this.#encoded.get(value); if (known !== undefined) return known; diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index 48091d9ad5..7238b3e6f7 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -325,7 +325,8 @@ export class SqliteAuthorityStore implements AuthorityStore { // it was reached: the head row's digest covers the live projection, the // retained transaction at that cursor must carry the same state digest and // must reproduce its exact commit proof, and the cursor bounds must stay - // contiguous with the head. Neither cost grows with retained history. + // contiguous with the head. Projection decoding never replays history; + // the indexed continuity count still depends on retained cursor count. const bounds = db.prepare(`SELECT (SELECT CAST(MIN(cursor) AS TEXT) FROM commits) AS first, (SELECT CAST(MAX(cursor) AS TEXT) FROM commits) AS last, diff --git a/tests/control_plane_ts/authority_archive.test.ts b/tests/control_plane_ts/authority_archive.test.ts index d4efe2f636..79f5c15622 100644 --- a/tests/control_plane_ts/authority_archive.test.ts +++ b/tests/control_plane_ts/authority_archive.test.ts @@ -9,6 +9,7 @@ import test from "node:test"; import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; import {sqliteAuthorityRuntime} from "../../loopx/control_plane/coordination/sqlite_runtime.ts"; +import {withVerifiedAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive_read.ts"; import {exportAuthorityArchive, verifyAuthorityArchive, restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; @@ -319,3 +320,24 @@ test("restore consumes reviewed snapshot even when the original path is replaced if (receipt.status === "found") assert.deepEqual(receipt.receipts, [{decision: 1}]); } finally { await rm(root, {recursive: true, force: true}); } }); + + +test("archive consumer mutations cannot change subsequent replay or the terminal seal", async t => { + const root = await mkdtemp(join(tmpdir(), "authority-archive-owned-")); + t.after(() => rm(root, {recursive: true, force: true})); + const source = new FileAuthorityStore(join(root, "source"), "goal"); + await seed(source); + const path = join(root, "backup.ndjson"); + const summary = await exportAuthorityArchive(source, "goal", path); + await withVerifiedAuthorityArchive(path, summary.archive_sha256, async archive => { + let count = 0; + for await (const row of archive.transactions()) { + count++; + assert.deepEqual(row.projection, {goal_id: "goal", value: count}); + row.projection.goal_id = "consumer-local-edit"; + row.projection.value = 999; + } + assert.equal(count, 3); + }); + assert.deepEqual(await verifyAuthorityArchive(path), summary); +}); From 758db221e16b452d420d82fc377787979b0539a5 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:50:19 +0800 Subject: [PATCH 09/10] test(authority): make archive fixture commit result explicit Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- tests/control_plane_ts/authority_archive.test.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/control_plane_ts/authority_archive.test.ts b/tests/control_plane_ts/authority_archive.test.ts index 79f5c15622..08172c13b1 100644 --- a/tests/control_plane_ts/authority_archive.test.ts +++ b/tests/control_plane_ts/authority_archive.test.ts @@ -1,4 +1,4 @@ -import type {AuthorityStore, AuthorityStoreCommit} from "../../loopx/control_plane/coordination/authority_store.ts"; +import type {AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult} from "../../loopx/control_plane/coordination/authority_store.ts"; import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; import assert from "node:assert/strict"; @@ -30,7 +30,7 @@ for (const sourceKind of ["file", "sqlite"] as const) { : new SqliteAuthorityStore(join(root, "source"), "goal"); let revision: string | null = null; for (let i = 1; i <= 7; i++) { - const result = await source.commitAuthority({expected_provider_revision: revision, + const result: AuthorityStoreCommitResult = await source.commitAuthority({expected_provider_revision: revision, operation_id: `op-${i}`, events: [{kind: "change", i}], next_projection: {goal_id: "goal", i, archived: ["todo_old"], unicode: "复杂目标"}, receipts: [{request_sha256: `request-${i}`, changed: i % 2 === 0}]}); From a75e16849d3f046f061dce8bb462a9a5ad943f52 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 21:11:14 +0800 Subject: [PATCH 10/10] docs(authority): reconcile replay measurements and local default holds Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- ...shared-goal-authority-state-provider-v0.md | 12 +- ...-goal-authority-state-provider-v0.zh-CN.md | 6 +- docs/reference/sqlite-authority-store.md | 104 ++++++++++++------ 3 files changed, 83 insertions(+), 39 deletions(-) diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 544878ea0a..77cf4168b2 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3209,10 +3209,14 @@ candidate; File remains the real reference/explicit profile and migration rehearsal backend. Do not publish two ambiguous defaults, declare the current File history layout long-horizon-qualified, or silently fall back from a selected SQLite store. Release activation must cite its D2 evidence. The September 27 matched -short-history experiment also selects SQLite as the **short-term default -implementation target**: writes/head/restart beat current checkpoint/delta -File, while File retains faster warm history reads. Large-state receipt/scan -budgets remain unmet, so this is not permission to enable the default now. +short-history experiments also select SQLite as the **short-term default +implementation target**: mutation/restart costs beat current checkpoint/delta +File, while warm read tradeoffs depend on the projection. PR #4931 now shares +privately owned TS replay between SQLite proofs and provider-neutral archive +recovery, retaining exact byte proofs and isolating returned rows. Matched +Linux evidence reduces many-field receipt/scan p95 by 88%/68%, but large-state +budgets and sustained-memory qualification remain open. This is not permission +to enable the default now. [Measurements, reproduction and D2/D3/L9 dependencies](../../reference/sqlite-authority-store.md#short-term-default-decision-and-matched-experiment). PostgreSQL shares the TS semantic contracts but has independent service, tenant, restore and capacity qualification; its deployment must not delay the diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index 736aff15fd..66838a56df 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -2482,8 +2482,10 @@ provider 确认;权威空集合不回退到陈旧 Markdown。Legacy 与预览 对照、显式可选 profile 和迁移演练后端。不能发布两个含混的默认项,不能把现有 File 历史布局直接称为长程合格,也不能从选定 SQLite 静默回退。发布启用必须引用 D2 证据。9 月 27 日同负载短历史实验也选择 SQLite 作为**短期默认实现目标**:写入、 -当前 head 与重启快于现有 checkpoint/delta File,而 File 的热历史读取仍更快。 -大状态 receipt/scan 仍未达预算,因此不是立即启用默认的许可。 +重启成本优于现有 checkpoint/delta File;热读取取舍取决于数据形状。#4931 现让 +SQLite 证明与 provider-neutral 归档恢复共用持有私有状态的 TS 重放组件,保留完整 +摘要字节,隔离返回记录与内部状态。同机 Linux 对照中,多字段 receipt/scan p95 +降低 88%/68%,但大状态预算与持续内存资格仍未闭合,不能立即启用默认。 [实测、复现命令及 D2/D3/L9 依赖](../../reference/sqlite-authority-store.md#short-term-default-decision-and-matched-experiment)。PostgreSQL 复用 TS 语义合同,但 service、tenant、restore 和 capacity 单独 资格化;其部署不阻塞本地路线。 diff --git a/docs/reference/sqlite-authority-store.md b/docs/reference/sqlite-authority-store.md index c7e0661440..dff1b9ef2f 100644 --- a/docs/reference/sqlite-authority-store.md +++ b/docs/reference/sqlite-authority-store.md @@ -17,38 +17,58 @@ qualification**. File remains an explicit provider, a conformance reference and an export/recovery destination. Existing Goals retain their selected authority; a rejected SQLite runtime must never silently open File instead. -The decision uses current File checkpoint/delta storage, not its retired -full-projection-per-commit layout. On macOS arm64, Node 22.22.3 / SQLite 3.51.3, -measurement source `e8193ce83` produced the following p95 milliseconds. Arms ran -sequentially on one host; each uses 20 warm read samples, the final 100 writes, -and five fresh-process head reads. Cold-process timing includes module loading -and does not clear the OS page cache. These bounded observations are not a -population estimate or a formal capacity/soak result. - -| Projection / commits | Provider | Write | Warm head | Historical receipt | Scan 100 | Cold process head | -| --- | --- | ---: | ---: | ---: | ---: | ---: | -| Mixed 20 KiB / 128 | SQLite | 5.12 | 1.65 | 28.63 | 69.39 | 90.82 | -| Mixed 20 KiB / 128 | File | 25.23 | 6.22 | 4.62 | 32.82 | 140.14 | -| Mixed 20 KiB / 512 | SQLite | 4.76 | 1.62 | 28.36 | 66.78 | 77.69 | -| Mixed 20 KiB / 512 | File | 36.05 | 6.88 | 4.84 | 37.32 | 295.88 | -| Changing 1 MiB / 128 | SQLite | 48.50 | 9.87 | 153.46 | 377.60 | 91.62 | -| Changing 1 MiB / 128 | File | 262.20 | 34.10 | 53.80 | 137.00 | 763.60 | -| 464 Todos, 64 leases, 220 KiB / 128 | SQLite | 24.28 | 6.09 | 274.32 | 728.69 | 85.24 | -| 464 Todos, 64 leases, 220 KiB / 128 | File | 42.10 | 42.49 | 25.27 | 427.99 | 745.60 | - -The same runner on main `d23f1c87d` measured SQLite changing-1-MiB receipt/ -scan p95 at 458.25/1056.86 ms, versus 153.46/377.60 ms here. Ordinary mixed -receipt/scan was 26.00/68.29 ms versus 28.63/69.39 ms; the many-field fixture -was 270.74/690.99 ms versus 274.32/728.69 ms. The optimization benefits repeated -large strings; it does not establish a speedup for ordinary record-rich states, -and their small overhead remains visible. Do not generalize that speedup to -all Goal shapes. - -SQLite wins writes, current-state reads and process restart in these workloads. -File's verified warm history cache wins historical receipt and page reads. -The many-field fixture matters: optimizing repeated large strings alone does -not make a real Todo-rich projection cheap. Its SQLite receipt/scan timings -still exceed the RFC's 50/250 ms targets; no threshold has been increased. +The decision compares current File checkpoint/delta storage, not its retired +full-projection-per-commit layout. The final matched experiment ran sequentially +on Linux x86_64 (16 vCPUs), Node 22.22.3 / SQLite 3.51.3, with baseline +`96ce9efb3` and candidate `758db221e`. For each workload the baseline SQLite +arm ran immediately before its candidate arm; candidate File arms followed. +The same runner and qualified runtime served every arm. Each uses 20 warm +read samples, the final 100 writes and five fresh-process head reads. Values +below are p95 milliseconds. Cold process includes module loading and does not +clear the OS page cache. These are bounded observations, not population +estimates or formal capacity/soak qualification. + +| Projection / commits | SQLite receipt before | Receipt after | SQLite scan 100 before | Scan after | +| --- | ---: | ---: | ---: | ---: | +| Mixed 20 KiB / 128 | 67.69 | 16.04 | 184.69 | 71.37 | +| Mixed 20 KiB / 512 | 75.70 | 17.18 | 187.22 | 67.61 | +| 464 Todos, 64 leases, 220 KiB / 128 | 786.01 | 91.52 | 1964.03 | 623.49 | +| Changing 1 MiB / 128 | 456.97 | 288.76 | 1190.10 | 1114.94 | +| Fixed 64 KiB / 128 | 62.96 | 9.57 | 170.34 | 56.36 | + +The many-field receipt/scan costs fall by 88%/68%, but still exceed the +50/250 ms targets on this host. Changing-1-MiB scan improves only about 6%: +returning 100 complete large states remains expensive. Mixed and fixed-state +reads also improve; the optimization is no longer limited to large strings. +Write costs remain close to baseline because this is primarily a replay change. +No frozen qualification threshold has been increased. + +The candidate-provider comparison retains the tradeoffs instead of declaring +one provider faster for every operation: + +| Projection / commits | SQLite write / warm head / cold head | File write / warm head / cold head | File receipt / scan 100 | +| --- | ---: | ---: | ---: | +| Mixed 20 KiB / 128 | 10.26 / 2.11 / 146.14 | 21.01 / 4.81 / 296.79 | 3.28 / 97.72 | +| Mixed 20 KiB / 512 | 10.56 / 2.59 / 143.09 | 51.87 / 6.08 / 702.25 | 5.55 / 78.13 | +| 464 Todos, 64 leases, 220 KiB / 128 | 64.24 / 16.44 / 152.50 | 72.75 / 7.69 / 1494.86 | 3.48 / 939.29 | +| Changing 1 MiB / 128 | 97.50 / 19.51 / 164.87 | 761.97 / 115.46 / 1892.69 | 154.34 / 441.47 | +| Fixed 64 KiB / 128 | 12.27 / 2.84 / 144.58 | 14.85 / 3.19 / 329.43 | 3.22 / 85.83 | + +SQLite has lower write and restart costs in every measured shape. File has +cheaper warm receipts and, for the many-field fixture, a cheaper warm head. +SQLite now scans ordinary/many-field history faster, while File wins the +changing-1-MiB scan. This supports SQLite as the implementation target for +mutation/restart-heavy local operation, not a universal read-performance claim. + +The 1-MiB SQLite post-fill RSS rises from about 118 to 160 MiB; other measured +SQLite shapes remain near their baseline. This includes fixture and verification +allocations and is not a steady-state or leak measurement. Sustained-memory +qualification remains open. Baseline/candidate SQLite final store bytes match +in all five shapes; no persisted format is changed. Earlier macOS measurements +motivated the initial target, but a new local paired run encountered load above +40 and unstable unchanged-baseline timings. Its correctness checks were kept; +its timing samples are not used to claim improvement or qualification. + These tests change a deterministic observation field in a retained full state; they exercise storage, not the complete CLI/Turn workflow. In `changing-1m`, padding is resized to retain exactly 1 MiB, so the large string itself can @@ -156,10 +176,28 @@ layered so that each layer pays only for what it returns: | Layer | Proves | Cost | | --- | --- | --- | -| Live head (`loadAuthority`, `commitAuthority`) | Head row digest over the live projection, the retained transaction at that cursor reproducing its exact commit digest, parent linkage, `min=1`/`count=max=head` cursor continuity, and the presence of the checkpoint that covers the head | One head row, one retained row, one parent digest and index lookups; independent of retained history | +| Live head (`loadAuthority`, `commitAuthority`) | Head row digest over the live projection, the retained transaction at that cursor reproducing its exact commit digest, parent linkage, `min=1`/`count=max=head` cursor continuity, and the presence of the checkpoint that covers the head | One head row, one retained row and one parent digest; no projection replay. The indexed continuity count still depends on retained cursor count | | Materialized history (`scanCommitted`, `readReceipt`) | Every row from the covering checkpoint through the requested span, including each delta, state digest and parent lineage; paged scans also prove the lookahead row used for `has_more` | At most one checkpoint window plus the requested span | | Archive audit (`verifyAuthorityHistory`) | The complete delta chain from the empty root, every checkpoint against retained history, and the final state against the head | Linear in retained history; qualification and recovery only | +Historical reconstruction uses the shared TS `AuthorityStateReplay` owner. +It validates and privately copies the initial state and each delta, copies only +changed object paths, and reuses exact canonical encodings for unchanged +subtrees. Both proofs still hash the complete original v0 byte sequence; this +is not a new Merkle proof or a persisted digest format. Cached subtree keys are +weak, and encoded long-string buffers are capped separately; no cache survives +the provider read/audit call. A receipt query verifies its entire covering span +without materializing projections it does not return. Scans return independent +JSON objects, so editing one row cannot alter another row or a later read. + +The same owner now reconstructs provider-neutral archives. Previously a +consumer could edit a yielded transaction's projection and change the decoder's +next replay basis, causing a later digest or Goal-identity failure. The decoder +now retains private state; callers receive detached projections. Existing +File/SQLite/PostgreSQL archives and terminal seals keep the same bytes. Rejected +delta batches never advance the replay frontier, and sparse protocol arrays +are rejected rather than bypassing validation of their absent elements. + A missing head, rolled-back head, internal cursor gap, rewritten receipt/event, orphaned parent digest or mismatched state digest is rejected as `provider_protocol_violation` before returning authority or accepting a write.