Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions loopx/control_plane/coordination/authority_state_log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)) {
Expand Down
4 changes: 2 additions & 2 deletions loopx/control_plane/coordination/sqlite_authority_store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import {
applyAuthorityStateDelta,
authorityStateCheckpointCursor,
authorityStateDelta,
authorityStateDeltaReconstructs,
authorityStateDigest,
authorityStateReplayBudget,
decodeAuthorityStateDelta,
Expand Down Expand Up @@ -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);
Expand Down
31 changes: 31 additions & 0 deletions tests/control_plane_ts/authority_state_log.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
applyAuthorityStateDelta,
authorityStateCheckpointCursor,
authorityStateDelta,
authorityStateDeltaReconstructs,
authorityStateDigest,
authorityStateReplayBudget,
decodeAuthorityStateDelta,
Expand Down Expand Up @@ -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<string, unknown>;
const nested = JSON.parse('{"scope": {"": {"__proto__": {"depth": 1}}}}') as Record<string, unknown>;
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};
Expand Down
45 changes: 45 additions & 0 deletions tests/control_plane_ts/sqlite_authority_store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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<string, unknown>[] = [
{},
JSON.parse('{"": {"marker": "empty"}, "__proto__": {"marker": "proto"}}') as Record<string, unknown>,
JSON.parse('{"nested": {"": [{"__proto__": "leaf"}]}}') as Record<string, unknown>,
JSON.parse('{}') as Record<string, unknown>,
];
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<string, unknown>[] = [];
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<string, unknown>);
}
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<string, unknown>).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");
Expand Down
Loading