diff --git a/loopx/control_plane/coordination/authority_state_log.ts b/loopx/control_plane/coordination/authority_state_log.ts index 1fc4e97bb1..12e656e899 100644 --- a/loopx/control_plane/coordination/authority_state_log.ts +++ b/loopx/control_plane/coordination/authority_state_log.ts @@ -163,6 +163,29 @@ export function applyAuthorityStateDelta( return canonicalAuthorityObject(result, "reconstructed authority state"); } +/** + * Does this delta rebuild exactly this projection from `previous`? + * + * The live SQLite writer and the V1 migration both have to prove that the delta + * they are about to persist really reconstructs the projection they are + * publishing, so the rule has one owner instead of one inline comparison per + * caller: the caller only ever needs to know whether the log it is writing is + * readable by its own read path. A delta that cannot be decoded, or that does + * not apply to `previous`, is a failed reconstruction rather than a different + * outcome, because a failed proof is what makes those callers fail closed. + */ +export function authorityStateDeltaReconstructs( + previous: JsonObject, + delta: AuthorityStateDelta, + projection: JsonObject, +): boolean { + try { + return canonicalBytesEqual(applyAuthorityStateDelta(previous, delta), projection); + } catch { + return false; + } +} + function applyAuthorityStateOperation(root: JsonObject, operation: AuthorityStateOperation): void { const segments = operation.path; if (segments.length === 0) protocol("authority state delta cannot target the root state"); diff --git a/loopx/control_plane/coordination/sqlite_authority_migration.ts b/loopx/control_plane/coordination/sqlite_authority_migration.ts index e053876ef6..edc12a7f6a 100644 --- a/loopx/control_plane/coordination/sqlite_authority_migration.ts +++ b/loopx/control_plane/coordination/sqlite_authority_migration.ts @@ -23,7 +23,7 @@ import type { JsonObject } from "../effect_program.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthorityObjectList, requireAuthorityStoreId } from "./authority_store_codec.ts"; import { applyAuthorityStateDelta, authorityStateCheckpointCursor, authorityStateDelta, - authorityStateDigest, decodeAuthorityStateDelta, + authorityStateDeltaReconstructs, authorityStateDigest, decodeAuthorityStateDelta, isAuthorityStateCheckpoint } from "./authority_state_log.ts"; import { SQLITE_AUTHORITY_STORE_SCHEMA, @@ -197,7 +197,7 @@ function executeSqliteAuthorityMigration( // commit: replaying the stored delta must reproduce this projection // byte for byte. A retained projection the new format cannot read // would otherwise be copied into V2 and only fail on a later read. - if (authorityStateDigest(applyAuthorityStateDelta(previous ?? {}, delta)) !== stateDigest) { + if (!authorityStateDeltaReconstructs(previous ?? {}, delta, projection)) { throw new AuthorityStoreProtocolError("V1 authority state delta does not reconstruct its commit"); } if (isAuthorityStateCheckpoint(cursor)) { diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index c206ad2c59..ca54a225a3 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -17,6 +17,7 @@ import { applyAuthorityStateDelta, authorityStateCheckpointCursor, authorityStateDelta, + authorityStateDeltaReconstructs, authorityStateDigest, authorityStateReplayBudget, decodeAuthorityStateDelta, @@ -461,8 +462,7 @@ export class SqliteAuthorityStore implements AuthorityStore { const delta = authorityStateDelta(before.projection, nextState); // The encoder is proved before it is persisted: replaying the stored // delta must reproduce the committed projection byte for byte. - if (!canonicalAuthorityBytes(applyAuthorityStateDelta(before.projection, delta)) - .equals(canonicalAuthorityBytes(nextState))) { + if (!authorityStateDeltaReconstructs(before.projection, delta, nextState)) { protocol("SQLite authority state delta does not reconstruct its commit"); } const stateDigest = authorityStateDigest(nextState); diff --git a/tests/control_plane_ts/authority_state_log.test.ts b/tests/control_plane_ts/authority_state_log.test.ts index 05e475f133..a3cc59e2ba 100644 --- a/tests/control_plane_ts/authority_state_log.test.ts +++ b/tests/control_plane_ts/authority_state_log.test.ts @@ -7,6 +7,7 @@ import { applyAuthorityStateDelta, authorityStateCheckpointCursor, authorityStateDelta, + authorityStateDeltaReconstructs, authorityStateDigest, authorityStateReplayBudget, decodeAuthorityStateDelta, @@ -92,6 +93,36 @@ test("authority state delta decoding fails closed at the storage boundary", () = } }); +test("the reconstruction rule answers for every JSON object key and a broken delta", () => { + // One owner decides "this delta rebuilds exactly this projection" for both + // the live writer and the V1 migration, so this test is about the rule's own + // contract: it answers for a projection keyed with `""` or `__proto__`, and a + // delta that cannot be decoded or applied is a failed reconstruction rather + // than a thrown error or a partial state. + const special = JSON.parse( + '{"": {"marker": "empty"}, "__proto__": {"marker": "proto"}, "todos": [{"id": "a"}]}', + ) as Record; + const nested = JSON.parse('{"scope": {"": {"__proto__": {"depth": 1}}}}') as Record; + for (const projection of [{}, special, nested]) { + assert.equal(authorityStateDeltaReconstructs({}, authorityStateDelta({}, projection), projection), true); + } + // A delta that decodes but describes a different projection is a failed + // reconstruction, not a different outcome the caller has to interpret. + const mismatched = decodeAuthorityStateDelta({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, + operations: [{op: "set", path: ["a"], value: 2}]}); + assert.equal(authorityStateDeltaReconstructs({}, mismatched, {a: 1}), false); + assert.equal(authorityStateDeltaReconstructs({}, mismatched, {a: 2}), true); + assert.equal(authorityStateDeltaReconstructs({}, mismatched, {}), false); + // An undecodable delta, and a delta whose path leaves the previous state, + // both fail closed through the same answer instead of propagating. + const undecodable = {schema_version: AUTHORITY_STATE_DELTA_SCHEMA, + operations: [{op: "set", path: [], value: 1}]} as never; + assert.equal(authorityStateDeltaReconstructs({}, undecodable, {}), false); + const missingPath = decodeAuthorityStateDelta({schema_version: AUTHORITY_STATE_DELTA_SCHEMA, + operations: [{op: "splice", path: ["absent"], index: 0, remove: 0, insert: []}]}); + assert.equal(authorityStateDeltaReconstructs({}, missingPath, {}), false); +}); + test("authority state digests and checkpoint windows are stable and bounded", () => { const left = {b: 2, a: [1, {z: 1, y: 2}]}; const right = {a: [1, {y: 2, z: 1}], b: 2}; diff --git a/tests/control_plane_ts/sqlite_authority_store.test.ts b/tests/control_plane_ts/sqlite_authority_store.test.ts index 67deb001ab..9631010443 100644 --- a/tests/control_plane_ts/sqlite_authority_store.test.ts +++ b/tests/control_plane_ts/sqlite_authority_store.test.ts @@ -8,6 +8,7 @@ import { spawn, spawnSync } from "node:child_process"; import { fileURLToPath } from "node:url"; import { SqliteAuthorityStore } from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; import { AUTHORITY_STATE_CHECKPOINT_INTERVAL } from "../../loopx/control_plane/coordination/authority_state_log.ts"; +import { canonicalAuthorityBytes } from "../../loopx/control_plane/coordination/authority_store_codec.ts"; import { authorityStoreCommitFixture, registerAuthorityStoreConformance } from "./authority_store_conformance.ts"; async function fixture(t: test.TestContext) { @@ -17,6 +18,50 @@ async function fixture(t: test.TestContext) { } registerAuthorityStoreConformance("SQLite", fixture); +test("SQLite commits and reads back every JSON object key", {timeout: 30000}, async t => { + const {store} = await fixture(t); + // The live writer must accept the same key space the migration has to carry: + // a projection may key an object with `""` or `__proto__`, and both must + // survive the stored delta, the head row and the retained history. + const projections: Record[] = [ + {}, + JSON.parse('{"": {"marker": "empty"}, "__proto__": {"marker": "proto"}}') as Record, + JSON.parse('{"nested": {"": [{"__proto__": "leaf"}]}}') as Record, + JSON.parse('{}') as Record, + ]; + let revision: string | null = null; + for (const [index, projection] of projections.entries()) { + const receipt = await store.commitAuthority({expected_provider_revision: revision, + operation_id: `key-op-${String(index).padStart(3, "0")}`, next_projection: projection, + events: [], receipts: []}); + assert.equal(receipt.status, "applied", JSON.stringify(receipt)); + revision = receipt.status === "applied" ? receipt.provider_revision : null; + } + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status === "loaded") { + assert.equal(canonicalAuthorityBytes(head.head).toString("utf8"), + canonicalAuthorityBytes(projections[projections.length - 1]!).toString("utf8")); + } + assert.equal((await store.verifyAuthorityHistory()).status, "verified"); + const read: Record[] = []; + let after: string | null = null; + for (;;) { + const page = await store.scanCommitted(after, 4); + assert.equal(page.status, "page", JSON.stringify(page)); + if (page.status !== "page" || page.transactions.length === 0) break; + for (const transaction of page.transactions) { + read.push(transaction.projection as Record); + } + after = page.transactions[page.transactions.length - 1]!.cursor; + } + for (const [index, projection] of projections.entries()) { + assert.equal(canonicalAuthorityBytes(read[index]!).toString("utf8"), + canonicalAuthorityBytes(projection).toString("utf8"), `projection ${index}`); + } + assert.equal(({} as Record).marker, undefined); +}); + 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");