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
17 changes: 14 additions & 3 deletions docs/reference/sqlite-authority-store.md
Original file line number Diff line number Diff line change
Expand Up @@ -310,14 +310,21 @@ The shipped version-1 database keeps one full projection per retained row, so it
cannot be read by the version-2 provider. `sqlite_authority_migration.ts`
migrates one Goal database in place: it reads the frozen version-1 rows, proves
every stored commit digest, writes the checkpoint/delta log, proves that each
written delta reconstructs its projection, requires the commit count to match,
swaps tables and updates the schema version inside a single
written delta reconstructs its projection, reads the not-yet-swapped tables back
and replays them through the store's own delta decoder, requires the commit
count to match, swaps tables and updates the schema version inside a single
`BEGIN IMMEDIATE` transaction. Any failure rolls back and leaves version 1
untouched; a second run reports `already_current`; a rewritten proof, a
mismatched goal/incarnation or an existing swap target fails closed. Cursors,
operation IDs, commit digests, provider revisions, receipts, events and scan
pages are byte-identical after the migration.

Retained projections keep every JSON object key the version-1 provider accepted,
including an empty key and a `__proto__` key: a database the previous provider
could read must not become one the version-2 provider cannot. The replay proof
is what keeps that promise honest, because identical identity and digest columns
alone would not show that a migrated state log is unreadable.

A version-1 database that published only its schema and metadata — the state a
goal leaves behind when it selected the provider and never committed — migrates
to an equally empty version-2 database instead of failing, so the operator is
Expand Down Expand Up @@ -390,4 +397,8 @@ delta”:活跃头读取只用自己的行、对应提交和游标连续性自
才写入,`--expected-identity` 可拒绝并非操作者所指的 incarnation,失败保持 v1
原样)。只发布过 schema 与 metadata、从未提交的 v1 库会迁移成同样为空的 v2 库,
不会让操作者落在两个 provider 都不接受的状态。迁移不改 cursor、operation id、
commit digest、provider revision、receipt、event 或 scan 页面字节。
commit digest、provider revision、receipt、event 或 scan 页面字节。迁移在提交前
还会把刚写入的表读回来、用 store 自己的 delta 解码器重放一遍:只核对搬过去的
标识与摘要无法证明新的状态日志可读。v1 能接受的 JSON key(包括空字符串和
`__proto__`)在 v2 中保持同样的数据语义,迁移不会把原本可读的库变成读不出来的
状态。
24 changes: 22 additions & 2 deletions loopx/control_plane/coordination/authority_state_log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,10 @@ export const AUTHORITY_STATE_CHECKPOINT_INTERVAL = 64;
/**
* Object path inside one projection. Arrays are addressed as a whole value;
* element identity is owned by the domain, not by this codec.
*
* Segments are JSON object keys, and strict JSON allows any string as a key:
* the empty string is a key, and so is `__proto__`. V1 retained whole
* projections, so both are legal retained data and stay legal here.
*/
export type AuthorityStatePath = readonly string[];

Expand Down Expand Up @@ -183,7 +187,22 @@ function applyAuthorityStateOperation(root: JsonObject, operation: AuthorityStat
if (!present) protocol("authority state delta removed a value that was never stored");
delete container[last];
}
else container[last] = structuredClone(operation.value);
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,
});
}

function descend(container: unknown, segment: string): unknown {
Expand Down Expand Up @@ -233,7 +252,8 @@ function decodeAuthorityStateOperation(value: unknown, index: number): Authority
function decodeAuthorityStatePath(value: unknown, label: string): AuthorityStatePath {
if (!Array.isArray(value) || value.length === 0) protocol(`${label} path is invalid`);
return value.map(segment => {
if (typeof segment === "string" && segment.length > 0) return 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`);
});
}
Expand Down
89 changes: 87 additions & 2 deletions loopx/control_plane/coordination/sqlite_authority_migration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,9 @@
import type { JsonObject } from "../effect_program.ts";
import { AuthorityStoreProtocolError, canonicalAuthorityObject,
canonicalAuthorityObjectList, requireAuthorityStoreId } from "./authority_store_codec.ts";
import { authorityStateCheckpointCursor, authorityStateDelta,
authorityStateDigest, isAuthorityStateCheckpoint } from "./authority_state_log.ts";
import { applyAuthorityStateDelta, authorityStateCheckpointCursor, authorityStateDelta,
authorityStateDigest, decodeAuthorityStateDelta,
isAuthorityStateCheckpoint } from "./authority_state_log.ts";
import {
SQLITE_AUTHORITY_STORE_SCHEMA,
commitDigest,
Expand Down Expand Up @@ -192,6 +193,13 @@
}
const stateDigest = authorityStateDigest(projection);
const delta = authorityStateDelta(previous ?? {}, projection);
// Prove the encoder before anything is written, exactly like a live
// 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) {
throw new AuthorityStoreProtocolError("V1 authority state delta does not reconstruct its commit");
}
if (isAuthorityStateCheckpoint(cursor)) {
insertCheckpoint.run(cursor.toString(), JSON.stringify(projection), stateDigest);
checkpoints += 1;
Expand All @@ -214,6 +222,11 @@
db.prepare("INSERT INTO head_v2 VALUES (1, ?, ?, ?)").run(cursor.toString(),
JSON.stringify(previous), previousDigest);
}
// Prove the new format is readable before it is adopted: every retained
// row is read back from its stored text and replayed through the codec the
// store reads with, so the swap cannot publish a log only its writer can
// decode.
verifyWrittenStateLogReadable(db, commits);
// Swap only after every retained row was proved and re-published.
db.exec("DROP TABLE head");
db.exec("DROP TABLE commits");
Expand Down Expand Up @@ -243,6 +256,78 @@
} finally { db?.close(); }
}

/**
* Read the not-yet-swapped V2 tables back and prove they reconstruct.
*
* The migration already proved each delta against its in-memory predecessor.
* This reads the rows it actually wrote, decodes them with the store's own
* delta decoder, replays them from the empty state, and requires every payload
* to match the digest the row published. Identical identity and digest columns
* alone would not catch a state log the new format cannot read.
*/
function verifyWrittenStateLogReadable(db: DatabaseSync, commits: number): void {

Check failure on line 268 in loopx/control_plane/coordination/sqlite_authority_migration.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 28 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=huangruiteng_loopx&issues=AaCoWDrvh7xehv9OKwwt&open=AaCoWDrvh7xehv9OKwwt&pullRequest=4497
const page = db.prepare(`SELECT CAST(cursor AS TEXT) AS sequence, state_digest, parent_state_digest,
delta FROM commits_v2 WHERE cursor > ? ORDER BY cursor LIMIT ?`);
let state: JsonObject = {};
let digest: string | null = null;
let cursor = 0n;
let read = 0;
for (;;) {
const rows = page.all(cursor.toString(), BigInt(MIGRATION_PAGE)) as unknown as
Record<string, unknown>[];
if (rows.length === 0) break;
for (const row of rows) {
const sequence = row.sequence;
if (typeof sequence !== "string" || !/^[1-9]\d*$/.test(sequence)) {
throw new AuthorityStoreProtocolError("migrated authority store cursor is invalid");
}
cursor = BigInt(sequence);
read += 1;
if ((row.parent_state_digest === null ? null : String(row.parent_state_digest)) !== digest) {

Check warning on line 286 in loopx/control_plane/coordination/sqlite_authority_migration.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

'row.parent_state_digest' may use Object's default stringification format ('[object Object]') when stringified.

See more on https://sonarcloud.io/project/issues?id=huangruiteng_loopx&issues=AaCoWDrvh7xehv9OKwwu&open=AaCoWDrvh7xehv9OKwwu&pullRequest=4497
throw new AuthorityStoreProtocolError("migrated authority state log parent lineage is invalid");
}
state = applyAuthorityStateDelta(state,
decodeAuthorityStateDelta(JSON.parse(String(row.delta)) as unknown));
if (authorityStateDigest(state) !== row.state_digest) {
throw new AuthorityStoreProtocolError("migrated authority state log digest mismatch");
}
digest = String(row.state_digest);
}
}
if (read !== commits) {
throw new AuthorityStoreProtocolError("migrated authority state log lost retained transactions");
}
const retainedDigest = db.prepare("SELECT state_digest FROM commits_v2 WHERE cursor = ?");
for (const row of db.prepare(`SELECT CAST(cursor AS TEXT) AS sequence, projection, projection_digest
FROM checkpoints_v2 ORDER BY cursor`).all() as unknown as Record<string, unknown>[]) {
const projection = canonicalAuthorityObject(JSON.parse(String(row.projection)) as unknown,
"migrated authority checkpoint projection");
if (authorityStateDigest(projection) !== row.projection_digest) {
throw new AuthorityStoreProtocolError("migrated authority checkpoint is not readable");
}
if (retainedDigest.get(String(row.sequence))?.state_digest !== row.projection_digest) {
throw new AuthorityStoreProtocolError("migrated authority checkpoint does not cover its cursor");
}
}
const head = db.prepare("SELECT CAST(cursor AS TEXT) AS sequence, projection, state_digest FROM head_v2 WHERE singleton = 1")
.get() as Record<string, unknown> | undefined;
if (read === 0) {
if (head !== undefined) {
throw new AuthorityStoreProtocolError(
"migrated authority head was published without a retained transaction");
}
return;
}
if (head === undefined || head.sequence !== cursor.toString() || head.state_digest !== digest) {

Check warning on line 321 in loopx/control_plane/coordination/sqlite_authority_migration.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer using an optional chain expression instead, as it's more concise and easier to read.

See more on https://sonarcloud.io/project/issues?id=huangruiteng_loopx&issues=AaCoWDrvh7xehv9OKwwv&open=AaCoWDrvh7xehv9OKwwv&pullRequest=4497
throw new AuthorityStoreProtocolError("migrated authority head does not match its state log");
}
const projection = canonicalAuthorityObject(JSON.parse(String(head.projection)) as unknown,
"migrated authority head projection");
if (authorityStateDigest(projection) !== digest) {
throw new AuthorityStoreProtocolError("migrated authority head is not readable");
}
}

/**
* Read the migrated store back and prove that its retained transaction
* identity sequence and authority lineage are exactly the pre-migration ones.
Expand Down
42 changes: 42 additions & 0 deletions tests/control_plane_ts/authority_state_log.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ import {
decodeAuthorityStateDelta,
isAuthorityStateCheckpoint,
} from "../../loopx/control_plane/coordination/authority_state_log.ts";
import {canonicalAuthorityBytes} from
"../../loopx/control_plane/coordination/authority_store_codec.ts";

test("authority state deltas reconstruct every committed projection exactly", () => {
const previous = {
Expand Down Expand Up @@ -108,3 +110,43 @@ test("authority state digests and checkpoint windows are stable and bounded", ()
assert.equal(isAuthorityStateCheckpoint(64n), false);
assert.throws(() => authorityStateCheckpointCursor(0n), /checkpoint cursor/u);
});

test("authority state deltas keep every JSON key the stored projection could carry", () => {
// V1 retained a whole projection per commit, so every JSON object key was
// legal legacy input. The delta codec inherits that contract: an empty key
// is an ordinary key, and a `__proto__` key is stored data instead of the
// inherited accessor. Keys that only existed in V1 data would otherwise be
// copied into the new format and fail on the first read.
const text = (value: unknown): string => canonicalAuthorityBytes(value).toString("utf8");
const stored = (previous: unknown, next: unknown): unknown => JSON.parse(
JSON.stringify(authorityStateDelta(previous as never, next as never))) as unknown;

const previous = JSON.parse('{"": {"empty": true}, "__proto__": {"own": 1}, "todos": []}') as
Record<string, unknown>;
const next = JSON.parse('{"": {"empty": false}, "__proto__": {"own": 2},' +
' "nested": {"__proto__": {"deep": true}}, "todos": [{"todo_id": "todo-0"}]}') as
Record<string, unknown>;
const applied = applyAuthorityStateDelta(previous, decodeAuthorityStateDelta(stored(previous, next)));
assert.equal(text(applied), text(next));
// The reconstruction is an ordinary JSON object: `__proto__` stayed a key.
assert.equal(Object.getPrototypeOf(applied), Object.prototype);
assert.deepEqual(Object.keys(applied).sort(), ["", "__proto__", "nested", "todos"]);
assert.deepEqual(applied[""], {empty: false});
assert.deepEqual(applied["__proto__"], {own: 2});
assert.deepEqual(applied["nested"], {["__proto__"]: {deep: true}});

// Removing a legacy key is legal too, including the empty key and a key that
// the previous state only carried inside a nested object.
const trimmed = JSON.parse('{"nested": {}, "todos": []}') as Record<string, unknown>;
const removal = applyAuthorityStateDelta(next, decodeAuthorityStateDelta(stored(next, trimmed)));
assert.equal(text(removal), text(trimmed));
assert.deepEqual(Object.keys(removal).sort(), ["nested", "todos"]);
assert.equal(Object.hasOwn(removal, "__proto__"), false);

// A legacy key the delta has to create is an own data property as well.
const legacyText = '{"__proto__": {"own": 1}, "": 0}';
const created = applyAuthorityStateDelta({},
decodeAuthorityStateDelta(stored({}, JSON.parse(legacyText) as Record<string, unknown>)));
assert.equal(text(created), text(JSON.parse(legacyText)));
assert.equal(Object.getPrototypeOf(created), Object.prototype);
});
68 changes: 68 additions & 0 deletions tests/control_plane_ts/sqlite_authority_migration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,74 @@ function runMigrationCli(args: readonly string[]): {status: number | null; stdou
"--experimental-strip-types", script, ...args], {encoding: "utf8"});
}

/**
* Retained V1 history that used only JSON object keys V1 never restricted.
*
* V1 stored a whole projection per commit, so `""` was an ordinary key and
* `__proto__` was data. The migration has to carry both into the new format and
* prove they still read back before it commits.
*/
function legacyKeySeeds(): SqliteAuthorityV1Seed[] {
return [
{operation_id: "legacy-op-001",
projection: {authority_revision: 1, "": {empty: true}, ["__proto__"]: {own: 1}, todos: []}},
{operation_id: "legacy-op-002",
projection: {authority_revision: 2, "": {empty: true}, ["__proto__"]: {own: 1},
nested: {["__proto__"]: {deep: true}}, todos: [{todo_id: "todo-0", status: "open"}]}},
{operation_id: "legacy-op-003",
projection: {authority_revision: 3, "": {empty: false}, nested: {},
todos: [{todo_id: "todo-0", status: "done"}]}},
];
}

test("SQLite V1 migration keeps every legacy JSON key readable", async t => {
const root = await directory(t);
const v1 = createSqliteAuthorityStoreV1(root, GOAL_ID, legacyKeySeeds());
const {DatabaseSync: Reader} = createRequire(import.meta.url)("node:sqlite");
const legacy = new Reader(v1.path, {readOnly: true});
try {
// The frozen V1 rows really carry the keys this test is about.
const text = String(legacy.prepare("SELECT projection FROM commits WHERE cursor = 2").get()?.projection);
assert.match(text, /"__proto__"/u);
assert.match(text, /""\s*:/u);
} finally { legacy.close(); }
const migrated = migrateSqliteAuthorityStoreV1ToV2(root, GOAL_ID, {execute: true});
assert.equal(migrated.status, "migrated", JSON.stringify(migrated));
assert.equal(migrated.commits, 3);
const store = new SqliteAuthorityStore(root, GOAL_ID, {existingOnly: true, expectedIdentity: v1.identity});
// The live head reads back as the exact projection the last V1 commit published.
const head = await store.loadAuthority();
assert.equal(head.status, "loaded", JSON.stringify(head));
if (head.status === "loaded") {
assert.equal(canonicalAuthorityBytes(head.head).toString("utf8"),
JSON.stringify(v1.rows[2]!.operation_receipt.projection));
// The empty key survives the migration, and the last commit really did
// drop `__proto__` rather than the codec silently dropping it for us.
assert.equal(Object.hasOwn(head.head, ""), true);
assert.equal(Object.hasOwn(head.head, "__proto__"), false);
assert.equal(Object.getPrototypeOf(head.head), Object.prototype);
}
// A retained projection keeps the legacy `__proto__` key as stored data.
const page = await store.scanCommitted(null, 6);
assert.equal(page.status, "page", JSON.stringify(page));
if (page.status === "page") {
const retained = page.transactions[1]!.projection;
assert.equal(Object.hasOwn(retained, "__proto__"), true);
assert.equal(Object.hasOwn(retained, ""), true);
assert.equal(Object.getPrototypeOf(retained), Object.prototype);
assert.deepEqual(retained["__proto__"], {own: 1});
}
// Every retained projection reads back byte-identically, not only the head.
const history = await logicalHistory(store);
for (const [index, row] of v1.rows.entries()) {
const parsed = JSON.parse(history[2 + index]!) as Record<string, unknown>;
assert.equal(parsed.cursor, row.cursor);
assert.equal(parsed.operation_id, row.operation_id);
assert.deepEqual(parsed.projection, JSON.stringify(row.operation_receipt.projection));
}
assert.equal((await store.verifyAuthorityHistory()).status, "verified");
});

test("SQLite V1 migration refuses a rewritten proof and leaves the database intact", async t => {
const root = await directory(t);
const v1 = createSqliteAuthorityStoreV1(root, GOAL_ID, seeds());
Expand Down