From bb52457ee23bd76bd050fd4611e5261e2774d9e5 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:18:16 +0800 Subject: [PATCH 1/4] fix(authority): verify recovery snapshots and batch SQLite receipt proofs Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/authority_archive.py | 21 +- .../coordination/authority_archive.ts | 246 ++++-------------- .../coordination/authority_archive_audit.ts | 136 ++++++++++ .../coordination/authority_archive_read.ts | 176 +++++++++++++ .../coordination/authority_store.ts | 18 ++ .../coordination/local_authority_archive.ts | 29 ++- .../coordination/sqlite_authority_store.ts | 68 +++-- tests/control_plane/test_authority_archive.py | 11 + .../authority_archive.test.ts | 7 +- .../authority_archive_audit.test.ts | 194 ++++++++++++++ .../authority_archive_crash.test.ts | 70 +++++ ...ity_archive_postgresql.integration.test.ts | 3 + .../authority_archive_restore_process.ts | 23 ++ .../sqlite_authority_store.test.ts | 57 ++++ 14 files changed, 846 insertions(+), 213 deletions(-) create mode 100644 loopx/control_plane/coordination/authority_archive_audit.ts create mode 100644 loopx/control_plane/coordination/authority_archive_read.ts create mode 100644 tests/control_plane_ts/authority_archive_audit.test.ts create mode 100644 tests/control_plane_ts/authority_archive_crash.test.ts create mode 100644 tests/control_plane_ts/authority_archive_restore_process.ts diff --git a/loopx/cli_commands/authority_archive.py b/loopx/cli_commands/authority_archive.py index c2a1b5a0ae..49c379a1a0 100644 --- a/loopx/cli_commands/authority_archive.py +++ b/loopx/cli_commands/authority_archive.py @@ -16,7 +16,7 @@ def register_authority_archive_command( add_subcommand_format: Callable[[argparse.ArgumentParser], None], ) -> None: parser = subparsers.add_parser( - "authority-archive", help="Export, verify or restore an isolated canonical authority copy." + "authority-archive", help="Export, verify, audit or restore a canonical authority copy." ) add_subcommand_format(parser) actions = parser.add_subparsers(dest="authority_archive_action", required=True) @@ -27,11 +27,17 @@ def register_authority_archive_command( mode.add_argument("--execute", action="store_true") upgrade.add_argument("--all-known", action="store_true", help="Include runtime roots of registered projects.") mode.add_argument("--require-current", action="store_true", help="Fail if a format upgrade is needed; never write.") - for name in ("export", "verify", "restore"): + for name in ("export", "verify", "restore", "audit"): action = actions.add_parser(name) action.add_argument("--archive", type=Path, required=True) if name != "verify": action.add_argument("--goal-id", required=True) + if name == "audit": + action.add_argument("--archive-sha256", required=True) + action.add_argument("--destination", type=Path, + help="Audit an isolated restore directory; otherwise audit the selected runtime provider.") + action.add_argument("--allow-newer-head", action="store_true", + help="Verify only the retained archive prefix, permitting later target commits.") if name == "restore": action.add_argument("--destination", type=Path, required=True) action.add_argument("--provider", choices=("file", "sqlite"), required=True) @@ -65,6 +71,14 @@ def handle_authority_archive_command( elif args.authority_archive_action == "restore": request.update(goal_id=args.goal_id, destination=str(args.destination.expanduser().resolve()), provider=args.provider, archive_sha256=args.archive_sha256, execute=args.execute) + elif args.authority_archive_action == "audit": + request.update(goal_id=args.goal_id, archive_sha256=args.archive_sha256, + allow_newer_head=args.allow_newer_head) + if args.destination is not None: + request["destination"] = str(args.destination.expanduser().resolve()) + else: + request["runtime_root"] = str(resolve_runtime_root( + load_registry(registry_path), runtime_root_arg, registry_path=registry_path)) result = effect_runtime_result( "coordination.authority_archive.manage", request, timeout=300.0, retry_safe=False ) @@ -75,7 +89,8 @@ def handle_authority_archive_command( result.update(status="failed", reason="Authority format upgrade required before activating this runtime.") print_payload(result, output_format(args), lambda value: ( f"Authority archive: {value.get('status')}\n" - f"{value.get('reason', 'Active authority selection is unchanged.')}" + f"{value.get('reason', 'Active authority selection is unchanged.')}\n" + f"{value.get('audit', '')}" )) return 1 if result.get("status") == "failed" else 0 diff --git a/loopx/control_plane/coordination/authority_archive.ts b/loopx/control_plane/coordination/authority_archive.ts index d42c4d9de9..46fc2c35d0 100644 --- a/loopx/control_plane/coordination/authority_archive.ts +++ b/loopx/control_plane/coordination/authority_archive.ts @@ -1,152 +1,20 @@ -/** Portable retained-journal recovery. A restored store is an isolated copy, - * never an authority selection, writer-fence rollback or execution grant. */ +/** Portable retained-journal recovery; activation remains a separate operation. */ import {randomUUID} from "node:crypto"; -import {createReadStream} from "node:fs"; import {link, open, unlink} from "node:fs/promises"; import {dirname} from "node:path"; import type {JsonObject} from "../effect_program.ts"; import type {AuthorityStore, AuthorityStoreCommittedTransaction} from "./authority_store.ts"; -import {AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthorityObjectList, - canonicalAuthoritySha256, hasExactAuthorityKeys, requireAuthorityStoreId} from "./authority_store_codec.ts"; -import {applyAuthorityStateDelta, authorityStateDelta, decodeAuthorityStateDelta} from "./authority_state_log.ts"; +import {AuthorityStoreProtocolError, canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {authorityStateDelta} from "./authority_state_log.ts"; +import {AUTHORITY_ARCHIVE_SCHEMA, archiveCursor as positive, verifyAuthorityArchive, + withVerifiedAuthorityArchive, type ArchiveHeader, type AuthorityArchiveSummary} from "./authority_archive_read.ts"; +import {checkArchivePage, checkArchivePrefix} from "./authority_archive_audit.ts"; +export {verifyAuthorityArchive} from "./authority_archive_read.ts"; +export type {AuthorityArchiveSummary} from "./authority_archive_read.ts"; -const SCHEMA = "loopx_authority_archive_v0"; const MAX_LINE_BYTES = 64 * 1024 * 1024; -const HEX = /^[0-9a-f]{64}$/; -interface ArchiveHeader extends JsonObject { - kind: "header"; - schema_version: typeof SCHEMA; - goal_id: string; - source_provider: string; - store_identity: string; - cursor: string; - provider_revision: string; - projection_sha256: string; -} -export interface AuthorityArchiveSummary { - schema_version: typeof SCHEMA; - goal_id: string; - source_provider: string; - source_store_identity: string; - source_provider_revision: string; - commits: string; - projection_sha256: string; - archive_sha256: string; -} -type ArchiveRecord = {kind: "header"; header: ArchiveHeader} | - {kind: "transaction"; transaction: AuthorityStoreCommittedTransaction} | - {kind: "verified"; summary: AuthorityArchiveSummary}; - function invalid(reason: string): never { throw new AuthorityStoreProtocolError(reason); } -function positive(value: unknown): string { - if (typeof value !== "string" || !/^[1-9]\d*$/.test(value)) invalid("archive cursor must be a positive integer string"); - return value; -} -function hash(value: unknown): string { - if (typeof value !== "string" || !HEX.test(value)) invalid("archive digest must be SHA-256"); - return value; -} -function exact(value: JsonObject, keys: readonly string[]): void { - if (!hasExactAuthorityKeys(value, keys)) invalid("archive record has missing or unknown fields"); -} -function summary(header: ArchiveHeader, digest: string): AuthorityArchiveSummary { - return {schema_version: SCHEMA, goal_id: header.goal_id, source_provider: header.source_provider, - source_store_identity: header.store_identity, source_provider_revision: header.provider_revision, - commits: header.cursor, projection_sha256: header.projection_sha256, archive_sha256: digest}; -} -function signed(value: JsonObject): JsonObject { - return {...value, sha256: canonicalAuthoritySha256(value)}; -} -function unsigned(value: JsonObject): JsonObject { - const {sha256, ...body} = value; - if (hash(sha256) !== canonicalAuthoritySha256(body)) invalid("archive record digest mismatch"); - return body; -} - -/** Bound a single physical line before JSON parsing, without buffering history. */ -async function* lines(path: string): AsyncGenerator { - let chunks: Buffer[] = []; - let size = 0; - for await (const chunk of createReadStream(path)) { - const bytes = chunk as Buffer; - let start = 0; - while (start < bytes.length) { - const end = bytes.indexOf(10, start); - const piece = bytes.subarray(start, end === -1 ? bytes.length : end); - size += piece.length; - if (size > MAX_LINE_BYTES) invalid("archive record exceeds 64 MiB"); - chunks.push(piece); - if (end === -1) break; - // Fatal UTF-8 decoding prevents replacement characters silently changing data. - yield new TextDecoder("utf-8", {fatal: true}).decode(Buffer.concat(chunks, size)); - chunks = []; size = 0; start = end + 1; - } - } - if (size !== 0) invalid("archive is truncated: final newline missing"); -} - -/** One decoder owns ordering, state reconstruction and the terminal seal for - * verify and restore. Hashes detect corruption; they are not signatures. */ -async function* records(path: string): AsyncGenerator { - let header: ArchiveHeader | null = null; - let digest: string | null = null; - let cursor = 0n; - let state: JsonObject = {}; - let revision: string | null = null; - let sealed = false; - const operations = new Set(); - let verified: AuthorityArchiveSummary | null = null; - for await (const line of lines(path)) { - if (sealed) invalid("archive has content after its terminal seal"); - const raw = canonicalAuthorityObject(JSON.parse(line) as unknown, "archive record"); - const value = unsigned(raw); - if (header === null) { - exact(value, ["kind", "schema_version", "goal_id", "source_provider", "store_identity", - "cursor", "provider_revision", "projection_sha256"]); - if (value.kind !== "header" || value.schema_version !== SCHEMA || - !["file", "sqlite", "postgresql", "nokv"].includes(String(value.source_provider))) invalid("invalid archive header"); - header = {...value, kind: "header", schema_version: SCHEMA, - goal_id: requireAuthorityStoreId(value.goal_id, "archive goal id"), - source_provider: String(value.source_provider), store_identity: requireAuthorityStoreId(value.store_identity, "store identity"), - cursor: positive(value.cursor), provider_revision: requireAuthorityStoreId(value.provider_revision, "provider revision"), - projection_sha256: hash(value.projection_sha256)}; - yield {kind: "header", header}; - } else if (value.kind === "transaction") { - exact(value, ["kind", "cursor", "provider_revision", "operation_id", "events", "receipts", - "delta", "projection_sha256", "previous_sha256"]); - const nextCursor = positive(value.cursor); - if (value.previous_sha256 !== digest || BigInt(nextCursor) !== cursor + 1n || - BigInt(nextCursor) > BigInt(header.cursor)) invalid("archive transaction lineage mismatch"); - 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)); - if (state.goal_id !== header.goal_id) invalid("archive transaction belongs to another goal"); - if (canonicalAuthoritySha256(state) !== hash(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, - operation_id: operation, events: canonicalAuthorityObjectList(value.events, "events"), - receipts: canonicalAuthorityObjectList(value.receipts, "receipts"), projection: state}}; - } else if (value.kind === "seal") { - 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"); - sealed = true; - verified = summary(header, hash(raw.sha256)); - } else invalid("unknown archive record kind"); - digest = hash(raw.sha256); - } - if (!sealed || verified === null) invalid("archive is incomplete: terminal seal missing"); - // Verification is delivered only after EOF, including the no-trailing-data check. - yield {kind: "verified", summary: verified}; -} - -export async function verifyAuthorityArchive(path: string): Promise { - for await (const record of records(path)) if (record.kind === "verified") return record.summary; - return invalid("archive verification did not finish"); -} +function signed(value: JsonObject): JsonObject { return {...value, sha256: canonicalAuthoritySha256(value)}; } /** Pin one retained prefix. Concurrent appends are allowed; a changed store * identity, missing interval or rewritten captured head is not. Output is @@ -161,7 +29,7 @@ export async function exportAuthorityArchive(store: AuthorityStore, goalId: stri const identity = await store.storeIdentity(); if (identity.status !== "available") invalid("archive source identity is unavailable"); if (head.head.goal_id !== goalId) invalid("archive source goal mismatch"); - const header: ArchiveHeader = {kind: "header", schema_version: SCHEMA, goal_id: goalId, + const header: ArchiveHeader = {kind: "header", schema_version: AUTHORITY_ARCHIVE_SCHEMA, goal_id: goalId, source_provider: store.providerKind ?? "file", store_identity: identity.store_identity, cursor: positive(head.cursor), provider_revision: head.provider_revision, projection_sha256: canonicalAuthoritySha256(head.head)}; @@ -214,58 +82,58 @@ export async function exportAuthorityArchive(store: AuthorityStore, goalId: stri } } -function semanticTransaction(row: AuthorityStoreCommittedTransaction): JsonObject { - const {provider_revision: _revision, ...semantic} = row; - return semantic; -} - -/** Restore only into a caller-owned isolated store. CLI enforces a private - * destination manifest/lock. Existing prefixes must match every retained - * transaction; mismatches and extra writes reject without overwriting them. - * A crash leaves an explicitly incomplete copy, resumable with the same digest. */ +/** Restore a verified private snapshot into an isolated target. Validate every + * existing row and receipt before writing the suffix. Recovery never retries an + * ambiguous write: only exact journal and original-receipt readback proves it. */ export async function restoreAuthorityArchive(path: string, target: AuthorityStore, expectedArchiveSha256: string): Promise { - const verified = await verifyAuthorityArchive(path); - if (verified.archive_sha256 !== hash(expectedArchiveSha256)) invalid("restore archive differs from the reviewed digest"); - const identity = await target.storeIdentity(); - if (identity.status !== "available" || identity.store_identity === verified.source_store_identity) { - invalid("restore requires an independent available target lineage"); - } - const head = await target.loadAuthority(); - if (head.status !== "missing" && head.status !== "loaded") invalid("restore target is unavailable"); - if (head.status === "loaded" && BigInt(positive(head.cursor)) > BigInt(verified.commits)) invalid("restore target has extra commits"); - let previous: string | null = null; - let after: string | null = null; - let completed = false; - for await (const record of records(path)) { - if (record.kind === "header") continue; - if (record.kind === "verified") { - if (record.summary.archive_sha256 !== verified.archive_sha256) invalid("archive changed during restore; target remains isolated"); - completed = true; continue; + return withVerifiedAuthorityArchive(path, expectedArchiveSha256, async archive => { + const verified = archive.summary; + const identity = await target.storeIdentity(); + if (identity.status !== "available" || identity.store_identity === verified.source_store_identity) { + invalid("restore requires an independent available target lineage"); } - const row = record.transaction; - if (BigInt(row.cursor) > BigInt(head.status === "loaded" ? head.cursor : "0")) { - // The sole write attempt may commit and lose its response. Read the exact - // retained row below, never repeat the write or infer success from state alone. + const head = await target.loadAuthority(); + if (head.status !== "missing" && head.status !== "loaded") invalid("restore target is unavailable"); + const retained = head.status === "loaded" ? positive(head.cursor) : "0"; + if (BigInt(retained) > BigInt(verified.commits)) invalid("restore target has extra commits"); + // Prefix auditing is paged and read-only, including the receipt index. A + // divergent or unreadable prefix never receives a speculative new suffix. + const prefix = await checkArchivePrefix(archive, target, retained); + if (prefix.status !== "matched") invalid(`restore ${prefix.reason_code}; resume the same archive`); + if (head.status === "loaded" && (prefix.provider_revision !== head.provider_revision || + prefix.projection_sha256 !== canonicalAuthoritySha256(head.head))) invalid("restore target head differs from retained history"); + let previous = prefix.provider_revision; + let after: string | null = retained === "0" ? null : retained; + let pending: AuthorityStoreCommittedTransaction[] = []; + for await (const row of archive.transactions()) { + if (BigInt(row.cursor) <= BigInt(retained)) continue; + let acknowledged: Awaited> | undefined; try { - await target.commitAuthority({expected_provider_revision: previous, operation_id: row.operation_id, + acknowledged = await target.commitAuthority({expected_provider_revision: previous, operation_id: row.operation_id, events: row.events, next_projection: row.projection, receipts: row.receipts}); } catch { /* exact journal readback is the recovery proof */ } + pending.push(row); + if (acknowledged?.status === "applied") { + if (acknowledged.cursor !== row.cursor) invalid("restore commit acknowledgement cursor mismatch"); + previous = acknowledged.provider_revision; + } + // Normal writes share a bounded page proof. An uncertain/rejected write + // forces immediate readback before any later write; never issue it twice. + if (pending.length === 16 || row.cursor === verified.commits || acknowledged?.status !== "applied") { + const checked = await checkArchivePage(target, pending, after); + if (checked.status !== "matched") invalid(`restore ${checked.reason_code}; resume the same archive`); + if (acknowledged?.status === "applied" && checked.provider_revision !== previous) { + invalid("restore commit acknowledgement differs from durable revision"); + } + previous = checked.provider_revision; after = row.cursor; pending = []; + } } - const readback = await target.scanCommitted(after, 1); - if (readback.status !== "page" || readback.transactions.length !== 1 || - canonicalAuthoritySha256(semanticTransaction(readback.transactions[0])) !== - canonicalAuthoritySha256(semanticTransaction(row))) invalid("restore target transaction differs or is unavailable; resume the same archive"); - const durable = readback.transactions[0]; - const receipt = await target.readReceipt(row.operation_id); - if (receipt.status !== "found" || receipt.cursor !== row.cursor || receipt.provider_revision !== durable.provider_revision || - canonicalAuthoritySha256(receipt.receipts) !== canonicalAuthoritySha256(row.receipts)) invalid("restore receipt readback mismatch"); - previous = durable.provider_revision; after = row.cursor; - } - const final = await target.loadAuthority(); - const finalIdentity = await target.storeIdentity(); - if (!completed || finalIdentity.status !== "available" || finalIdentity.store_identity !== identity.store_identity || - final.status !== "loaded" || final.cursor !== verified.commits || final.provider_revision !== previous || - canonicalAuthoritySha256(final.head) !== verified.projection_sha256) invalid("restore final readback mismatch"); - return {...verified, status: "restored", target_store_identity: identity.store_identity, target_provider_revision: final.provider_revision}; + const final = await target.loadAuthority(); + const finalIdentity = await target.storeIdentity(); + if (finalIdentity.status !== "available" || finalIdentity.store_identity !== identity.store_identity || + final.status !== "loaded" || final.cursor !== verified.commits || final.provider_revision !== previous || + canonicalAuthoritySha256(final.head) !== verified.projection_sha256) invalid("restore final readback mismatch"); + return {...verified, status: "restored", target_store_identity: identity.store_identity, target_provider_revision: final.provider_revision}; + }); } diff --git a/loopx/control_plane/coordination/authority_archive_audit.ts b/loopx/control_plane/coordination/authority_archive_audit.ts new file mode 100644 index 0000000000..33e42fe0e0 --- /dev/null +++ b/loopx/control_plane/coordination/authority_archive_audit.ts @@ -0,0 +1,136 @@ +/** Provider-neutral recovery evidence. State equality alone cannot prove that + * historical decisions, no-change receipts, or receipt lookup survived. */ +import {readAuthorityReceipts, type AuthorityStore, type AuthorityStoreCommittedTransaction} from "./authority_store.ts"; +import {canonicalAuthoritySha256} from "./authority_store_codec.ts"; +import {withVerifiedAuthorityArchive, type VerifiedAuthorityArchive} from "./authority_archive_read.ts"; + +type AuditFailure = { + status: "mismatch" | "unavailable"; + reason_code: string; + /** Cursor only: private operation IDs, state and receipt bodies stay local. */ + cursor?: string; +}; +type PrefixProof = { + status: "matched"; + compared_commits: string; + provider_revision: string | null; + projection_sha256: string | null; +}; + +function semanticHash(row: AuthorityStoreCommittedTransaction): string { + // Physical CAS tokens change across providers; logical rows must not. + const {provider_revision: _physicalRevision, ...semantic} = row; + return canonicalAuthoritySha256(semantic); +} + +/** Compare one ordered page and prove lookup of all its original receipts. */ +export async function checkArchivePage(store: AuthorityStore, + expected: readonly AuthorityStoreCommittedTransaction[], after: string | null): Promise { + if (expected.length < 1 || expected.length > 16) throw new TypeError("archive page requires 1..16 rows"); + const scan = await store.scanCommitted(after, expected.length); + const cursor = expected[0].cursor; + if (scan.status !== "page") return {status: "unavailable", reason_code: "archive_history_unavailable", cursor}; + if (scan.transactions.length !== expected.length || scan.next_cursor !== expected.at(-1)!.cursor) { + return {status: "mismatch", reason_code: "archive_history_incomplete", cursor}; + } + for (let i = 0; i < expected.length; i++) { + if (semanticHash(expected[i]) !== semanticHash(scan.transactions[i])) { + return {status: "mismatch", reason_code: "archive_transaction_mismatch", cursor: expected[i].cursor}; + } + } + const batch = await readAuthorityReceipts(store, expected.map(row => row.operation_id)); + if (batch.status !== "receipts") return {status: "unavailable", reason_code: "archive_receipt_unavailable", cursor}; + if (batch.results.length !== expected.length) return {status: "mismatch", reason_code: "archive_receipt_mismatch", cursor}; + for (let i = 0; i < expected.length; i++) { + const receipt = batch.results[i], row = expected[i]; + if (receipt.status === "failed" || receipt.status === "unavailable") { + return {status: "unavailable", reason_code: "archive_receipt_unavailable", cursor: row.cursor}; + } + if (receipt.status !== "found" || receipt.cursor !== row.cursor || + receipt.provider_revision !== scan.transactions[i].provider_revision || + canonicalAuthoritySha256(receipt.receipts) !== canonicalAuthoritySha256(row.receipts)) { + return {status: "mismatch", reason_code: "archive_receipt_mismatch", cursor: row.cursor}; + } + } + const last = scan.transactions.at(-1)!; + return {status: "matched", compared_commits: last.cursor, provider_revision: last.provider_revision, + projection_sha256: canonicalAuthoritySha256(last.projection)}; +} + +/** The source is fully verified before this cursor-bounded comparison starts. + * Reaching a partial target never accepts an unchecked archive suffix. */ +export async function checkArchivePrefix(archive: VerifiedAuthorityArchive, + store: AuthorityStore, count: string): Promise { + if (!/^(0|[1-9]\d*)$/.test(count) || BigInt(count) > BigInt(archive.summary.commits)) { + throw new TypeError("archive prefix count is outside verified history"); + } + let proof: PrefixProof = {status: "matched", compared_commits: "0", + provider_revision: null, projection_sha256: null}; + if (count === "0") return proof; + let page: AuthorityStoreCommittedTransaction[] = []; + for await (const expected of archive.transactions()) { + page.push(expected); + if (page.length < 16 && expected.cursor !== count) continue; + const checked = await checkArchivePage(store, page, proof.compared_commits === "0" ? null : proof.compared_commits); + if (checked.status !== "matched") return checked; + proof = checked; page = []; + if (proof.compared_commits === count) return proof; + } + throw new Error("verified archive prefix is incomplete"); +} + +export type AuthorityArchiveAuditResult = { + schema_version: "loopx_authority_archive_audit_v0"; + scope: "exact" | "retained_prefix"; + archive_sha256: string; + execution_authority_granted: false; +} & (AuditFailure | { + status: "matched"; + compared_commits: string; + target_store_identity: string; + captured_target_cursor: string; + captured_target_provider_revision: string; + matched_prefix_provider_revision: string; +}); + +export async function auditAuthorityArchive(path: string, target: AuthorityStore, + expectedSha256: string, scope: "exact" | "retained_prefix" = "exact"): Promise { + if (scope !== "exact" && scope !== "retained_prefix") throw new TypeError("invalid archive audit scope"); + return withVerifiedAuthorityArchive(path, expectedSha256, async archive => { + const base = {schema_version: "loopx_authority_archive_audit_v0" as const, scope, + archive_sha256: archive.summary.archive_sha256, execution_authority_granted: false as const}; + const identity = await target.storeIdentity(); + const head = await target.loadAuthority(); + if (identity.status !== "available" || head.status === "failed" || head.status === "unavailable") { + return {...base, status: "unavailable", reason_code: "archive_target_unavailable"}; + } + if (head.status === "missing") return {...base, status: "mismatch", reason_code: "archive_target_missing"}; + if (head.head.goal_id !== archive.summary.goal_id) return {...base, status: "mismatch", reason_code: "archive_goal_mismatch"}; + if (BigInt(head.cursor) < BigInt(archive.summary.commits)) { + return {...base, status: "mismatch", reason_code: "archive_target_incomplete"}; + } + if (scope === "exact" && head.cursor !== archive.summary.commits) { + return {...base, status: "mismatch", reason_code: "archive_target_has_newer_commits"}; + } + const proof = await checkArchivePrefix(archive, target, archive.summary.commits); + if (proof.status !== "matched") return {...base, ...proof}; + const finalIdentity = await target.storeIdentity(); + const final = await target.loadAuthority(); + if (finalIdentity.status !== "available" || final.status !== "loaded") { + return {...base, status: "unavailable", reason_code: "archive_target_readback_unavailable"}; + } + if (finalIdentity.store_identity !== identity.store_identity || final.head.goal_id !== archive.summary.goal_id || + BigInt(final.cursor) < BigInt(head.cursor) || + (final.cursor === head.cursor && final.provider_revision !== head.provider_revision)) { + return {...base, status: "mismatch", reason_code: "archive_target_lineage_changed"}; + } + if (scope === "exact" && (final.cursor !== head.cursor || final.provider_revision !== proof.provider_revision || + canonicalAuthoritySha256(final.head) !== archive.summary.projection_sha256)) { + return {...base, status: "mismatch", reason_code: "archive_target_changed"}; + } + return {...base, status: "matched", compared_commits: proof.compared_commits, + target_store_identity: identity.store_identity, captured_target_cursor: head.cursor, + captured_target_provider_revision: head.provider_revision, + matched_prefix_provider_revision: proof.provider_revision!}; + }); +} diff --git a/loopx/control_plane/coordination/authority_archive_read.ts b/loopx/control_plane/coordination/authority_archive_read.ts new file mode 100644 index 0000000000..c5c0be042d --- /dev/null +++ b/loopx/control_plane/coordination/authority_archive_read.ts @@ -0,0 +1,176 @@ +/** Streaming archive decoding and a private, verified input snapshot. */ +import {constants, createReadStream} from "node:fs"; +import {chmod, copyFile, mkdtemp, rm} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +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"; + +export const AUTHORITY_ARCHIVE_SCHEMA = "loopx_authority_archive_v0"; +const MAX_LINE_BYTES = 64 * 1024 * 1024; +const HEX = /^[0-9a-f]{64}$/; +export interface ArchiveHeader extends JsonObject { + kind: "header"; + schema_version: typeof AUTHORITY_ARCHIVE_SCHEMA; + goal_id: string; + source_provider: string; + store_identity: string; + cursor: string; + provider_revision: string; + projection_sha256: string; +} +export interface AuthorityArchiveSummary { + schema_version: typeof AUTHORITY_ARCHIVE_SCHEMA; + goal_id: string; + source_provider: string; + source_store_identity: string; + source_provider_revision: string; + commits: string; + projection_sha256: string; + archive_sha256: string; +} +type ArchiveRecord = {kind: "header"; header: ArchiveHeader} | + {kind: "transaction"; transaction: AuthorityStoreCommittedTransaction} | + {kind: "verified"; summary: AuthorityArchiveSummary}; + +function invalid(reason: string): never { throw new AuthorityStoreProtocolError(reason); } +export function archiveCursor(value: unknown): string { + if (typeof value !== "string" || !/^[1-9]\d*$/.test(value)) invalid("archive cursor must be a positive integer string"); + return value; +} +export function archiveHash(value: unknown): string { + if (typeof value !== "string" || !HEX.test(value)) invalid("archive digest must be SHA-256"); + return value; +} +function exact(value: JsonObject, keys: readonly string[]): void { + if (!hasExactAuthorityKeys(value, keys)) invalid("archive record has missing or unknown fields"); +} +function summary(header: ArchiveHeader, digest: string): AuthorityArchiveSummary { + return {schema_version: AUTHORITY_ARCHIVE_SCHEMA, goal_id: header.goal_id, source_provider: header.source_provider, + source_store_identity: header.store_identity, source_provider_revision: header.provider_revision, + commits: header.cursor, projection_sha256: header.projection_sha256, archive_sha256: digest}; +} +function unsigned(value: JsonObject): JsonObject { + const {sha256, ...body} = value; + if (archiveHash(sha256) !== canonicalAuthoritySha256(body)) invalid("archive record digest mismatch"); + return body; +} + +/** Bound a single physical line before JSON parsing, without buffering history. */ +async function* lines(path: string): AsyncGenerator { + let chunks: Buffer[] = []; + let size = 0; + for await (const chunk of createReadStream(path)) { + const bytes = chunk as Buffer; + let start = 0; + while (start < bytes.length) { + const end = bytes.indexOf(10, start); + const piece = bytes.subarray(start, end === -1 ? bytes.length : end); + size += piece.length; + if (size > MAX_LINE_BYTES) invalid("archive record exceeds 64 MiB"); + chunks.push(piece); + if (end === -1) break; + // Fatal UTF-8 decoding prevents replacement characters silently changing data. + yield new TextDecoder("utf-8", {fatal: true}).decode(Buffer.concat(chunks, size)); + chunks = []; size = 0; start = end + 1; + } + } + if (size !== 0) invalid("archive is truncated: final newline missing"); +} + +/** One decoder owns ordering, state reconstruction and the terminal seal for + * verify and restore. Hashes detect corruption; they are not signatures. */ +async function* archiveRecords(path: string): AsyncGenerator { + let header: ArchiveHeader | null = null; + let digest: string | null = null; + let cursor = 0n; + let state: JsonObject = {}; + let revision: string | null = null; + let sealed = false; + const operations = new Set(); + let verified: AuthorityArchiveSummary | null = null; + for await (const line of lines(path)) { + if (sealed) invalid("archive has content after its terminal seal"); + const raw = canonicalAuthorityObject(JSON.parse(line) as unknown, "archive record"); + const value = unsigned(raw); + if (header === null) { + exact(value, ["kind", "schema_version", "goal_id", "source_provider", "store_identity", + "cursor", "provider_revision", "projection_sha256"]); + if (value.kind !== "header" || value.schema_version !== AUTHORITY_ARCHIVE_SCHEMA || + !["file", "sqlite", "postgresql", "nokv"].includes(String(value.source_provider))) invalid("invalid archive header"); + header = {...value, kind: "header", schema_version: AUTHORITY_ARCHIVE_SCHEMA, + goal_id: requireAuthorityStoreId(value.goal_id, "archive goal id"), + source_provider: String(value.source_provider), store_identity: requireAuthorityStoreId(value.store_identity, "store identity"), + cursor: archiveCursor(value.cursor), provider_revision: requireAuthorityStoreId(value.provider_revision, "provider revision"), + projection_sha256: archiveHash(value.projection_sha256)}; + yield {kind: "header", header}; + } else if (value.kind === "transaction") { + exact(value, ["kind", "cursor", "provider_revision", "operation_id", "events", "receipts", + "delta", "projection_sha256", "previous_sha256"]); + const nextCursor = archiveCursor(value.cursor); + if (value.previous_sha256 !== digest || BigInt(nextCursor) !== cursor + 1n || + BigInt(nextCursor) > BigInt(header.cursor)) invalid("archive transaction lineage mismatch"); + 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)); + 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"); + revision = requireAuthorityStoreId(value.provider_revision, "provider revision"); + cursor += 1n; + yield {kind: "transaction", transaction: {cursor: cursor.toString(), provider_revision: revision, + operation_id: operation, events: canonicalAuthorityObjectList(value.events, "events"), + receipts: canonicalAuthorityObjectList(value.receipts, "receipts"), projection: state}}; + } else if (value.kind === "seal") { + 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"); + sealed = true; + verified = summary(header, archiveHash(raw.sha256)); + } else invalid("unknown archive record kind"); + digest = archiveHash(raw.sha256); + } + if (!sealed || verified === null) invalid("archive is incomplete: terminal seal missing"); + // Verification is delivered only after EOF, including the no-trailing-data check. + yield {kind: "verified", summary: verified}; +} + +export async function verifyAuthorityArchive(path: string): Promise { + for await (const record of archiveRecords(path)) if (record.kind === "verified") return record.summary; + return invalid("archive verification did not finish"); +} + + +/** The path stays private to this scope. TypeScript readonly protects callers; + * private copied bytes, rather than a type assertion or an open original inode, + * protect a restore against replacement and in-place edits of the input. */ +export interface VerifiedAuthorityArchive { + readonly summary: Readonly; + transactions(): AsyncGenerator; +} + +export async function withVerifiedAuthorityArchive(path: string, expectedSha256: string, + consume: (archive: VerifiedAuthorityArchive) => Promise): Promise { + archiveHash(expectedSha256); + const directory = await mkdtemp(join(tmpdir(), "loopx-authority-archive-")); + const snapshot = join(directory, "input.ndjson"); + try { + // Copy may race a writer. Full verification and the reviewed digest must + // pass on the resulting private bytes before any target is even opened. + await copyFile(path, snapshot, constants.COPYFILE_EXCL); + await chmod(snapshot, 0o600); + const verified = await verifyAuthorityArchive(snapshot); + if (verified.archive_sha256 !== expectedSha256) invalid("restore archive differs from the reviewed digest"); + return await consume({summary: Object.freeze(verified), async *transactions() { + for await (const record of archiveRecords(snapshot)) { + if (record.kind === "transaction") yield record.transaction; + } + }}); + } finally { + await rm(directory, {recursive: true, force: true}); + } +} diff --git a/loopx/control_plane/coordination/authority_store.ts b/loopx/control_plane/coordination/authority_store.ts index 46ca8cf2a6..fcf45fa85c 100644 --- a/loopx/control_plane/coordination/authority_store.ts +++ b/loopx/control_plane/coordination/authority_store.ts @@ -156,6 +156,13 @@ export type AuthorityStoreReceiptResult = | { status: "missing" } | AuthorityStoreReadFailure; +/** Results retain caller order, including duplicate IDs and missing receipts. + * Each item has exactly the same proof contract as readReceipt. No global + * snapshot promise is made by the fallback; native providers may strengthen it. */ +export type AuthorityStoreReceiptBatchResult = + | {status: "receipts"; results: readonly AuthorityStoreReceiptResult[]} + | AuthorityStoreReadFailure; + export type AuthorityStoreScanResult = | { status: "page"; @@ -177,9 +184,20 @@ export interface AuthorityStore { loadAuthority(): Promise; commitAuthority(commit: AuthorityStoreCommit): Promise; readReceipt(operationId: string): Promise; + /** Optional bounded acceleration; never a weaker receipt verification path. */ + readReceipts?(operationIds: readonly string[]): Promise; scanCommitted(afterCursor: string | null, limit: number): Promise; } +export async function readAuthorityReceipts(store: AuthorityStore, + operationIds: readonly string[]): Promise { + if (operationIds.length < 1 || operationIds.length > 64) throw new TypeError("receipt batch requires 1..64 operations"); + if (store.readReceipts !== undefined) return store.readReceipts(operationIds); + const results: AuthorityStoreReceiptResult[] = []; + for (const id of operationIds) results.push(await store.readReceipt(id)); + return {status: "receipts", results}; +} + /** Map a storage implementation to the public source label used by adapters. */ export function authorityStoreSourceAuthority(store: AuthorityStore): AuthorityStoreSourceAuthority { const kind = store.providerKind ?? "file"; diff --git a/loopx/control_plane/coordination/local_authority_archive.ts b/loopx/control_plane/coordination/local_authority_archive.ts index 444670ed67..b51d9977e5 100644 --- a/loopx/control_plane/coordination/local_authority_archive.ts +++ b/loopx/control_plane/coordination/local_authority_archive.ts @@ -9,6 +9,7 @@ import {durableWriteJson, withFileMutationLock} from "../effect_runtime_io.ts"; import {requireJsonObject} from "../runtime_decode.ts"; import {canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; import {exportAuthorityArchive, restoreAuthorityArchive, verifyAuthorityArchive} from "./authority_archive.ts"; +import {auditAuthorityArchive} from "./authority_archive_audit.ts"; import {FileAuthorityStore} from "./file_authority_store.ts"; import {SqliteAuthorityStore} from "./sqlite_authority_store.ts"; import {openRuntimeAuthorityStore, requireLocalAuthorityRuntimeRoot, @@ -38,9 +39,35 @@ export async function manageLocalAuthorityArchive(value: unknown, if (request.action === "verify") return {...base, status: "verified", archive: await verifyAuthorityArchive(archive)}; const goalId = requireAuthorityStoreId(request.goal_id, "goal id"); if (request.action === "export") { - const store = await openRuntimeAuthorityStore(requireLocalAuthorityRuntimeRoot(request.runtime_root), goalId, dependencies); + const store = await openRuntimeAuthorityStore(requireLocalAuthorityRuntimeRoot(request.runtime_root), goalId, dependencies, {existingOnly: true}); return {...base, status: "exported", archive: await exportAuthorityArchive(store, goalId, archive)}; } + if (request.action === "audit") { + if (typeof request.archive_sha256 !== "string" || + (request.allow_newer_head !== undefined && typeof request.allow_newer_head !== "boolean")) { + throw new Error("audit requires the reviewed archive digest and a boolean prefix policy"); + } + const inspected = await verifyAuthorityArchive(archive); + if (inspected.goal_id !== goalId) throw new Error("audit archive goal mismatch"); + let store; + if (request.destination !== undefined) { + if (request.runtime_root !== undefined) throw new Error("audit selects a runtime or an isolated destination, not both"); + const destination = path(request.destination, "audit destination"); + const binding = requireJsonObject(JSON.parse(await readFile(join(destination, "restore-binding.json"), "utf8")), "restore binding"); + if (binding.schema_version !== "loopx_authority_restore_destination_v0" || + binding.goal_id !== goalId || binding.archive_sha256 !== request.archive_sha256 || + (binding.provider !== "file" && binding.provider !== "sqlite")) { + throw new Error("audit destination belongs to a different archive or provider"); + } + store = binding.provider === "file" ? new FileAuthorityStore(join(destination, "store"), goalId, {existingOnly: true}) + : new SqliteAuthorityStore(join(destination, "store"), goalId, {existingOnly: true}); + } else { + store = await openRuntimeAuthorityStore(requireLocalAuthorityRuntimeRoot(request.runtime_root), goalId, dependencies, {existingOnly: true}); + } + const audit = await auditAuthorityArchive(archive, store, request.archive_sha256, + request.allow_newer_head === true ? "retained_prefix" : "exact"); + return {...base, status: audit.status === "matched" ? "audited" : "failed", audit}; + } if (request.action !== "restore") throw new Error("unknown authority archive action"); const inspected = await verifyAuthorityArchive(archive); if (inspected.goal_id !== goalId || inspected.archive_sha256 !== request.archive_sha256) { diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index 79571dc61b..ceabac52f2 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -8,7 +8,7 @@ import type { DatabaseSync } from "node:sqlite"; import type { JsonObject } from "../effect_program.ts"; import type { AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult, AuthorityStoreCommittedTransaction, AuthorityStoreIdentityResult, AuthorityStoreLoadResult, AuthorityStoreReadFailure, AuthorityStoreHead, - AuthorityStoreReceiptResult, AuthorityStoreScanResult } from "./authority_store.ts"; + AuthorityStoreReceiptResult, AuthorityStoreReceiptBatchResult, AuthorityStoreScanResult } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityBytes, canonicalAuthorityObject, canonicalAuthorityObjectList, canonicalAuthoritySha256, normalizeAuthorityStoreCommit, requireAuthorityStoreId } from "./authority_store_codec.ts"; @@ -333,10 +333,6 @@ export class SqliteAuthorityStore implements AuthorityStore { return {identity, checkpoint, transactions, state}; } - private verifiedWindow(db: DatabaseSync, target: bigint): SqliteCommitWindow { - return this.verifiedRange(db, target, target); - } - /** The live head, proven without materializing retained history. */ private current(db: DatabaseSync): {state: SqliteStateCursor; provider_revision: string; identity: string} | null { @@ -514,24 +510,60 @@ export class SqliteAuthorityStore implements AuthorityStore { } async readReceipt(operationId: string): Promise { + const batch = await this.readReceipts([operationId]); + return batch.status === "receipts" ? batch.results[0]! : batch; + } + + /** One snapshot and one proof per touched checkpoint window. The operation + * index still resolves each requested ID; scanned data never substitutes for + * lookup, and no verified state escapes this transaction as a cache. */ + async readReceipts(operationIds: readonly string[]): Promise { let db: DatabaseSync | null = null; try { - requireAuthorityStoreId(operationId, "operation id"); + if (!Array.isArray(operationIds) || operationIds.length < 1 || operationIds.length > 64) { + protocol("receipt batch requires 1..64 operations"); + } + for (const id of operationIds) requireAuthorityStoreId(id, "operation id"); + const missing = (): AuthorityStoreReceiptBatchResult => ({status: "receipts", + results: operationIds.map(() => ({status: "missing"}))}); db = this.open(false); - if (!db) return {status: "missing"}; + if (!db) return missing(); db.exec("BEGIN"); const head = this.current(db); - if (head === null) return {status: "missing"}; - const row = db.prepare(`SELECT ${COMMIT_COLUMNS} FROM commits WHERE operation_id = ?`) - .get(operationId) as unknown as Record | undefined; - if (!row) return {status: "missing"}; - // The selected row is revalidated with the bounded window that produced - // it, so a receipt cannot be read without its own proof. - const window = this.verifiedWindow(db, this.decodeCommitRow(row).cursor); - const transaction = window.transactions.find(item => item.operation_id === operationId); - if (!transaction) protocol("SQLite authority receipt is not part of its retained window"); - return {status: "found", cursor: transaction.cursor, - provider_revision: transaction.provider_revision, receipts: transaction.receipts}; + if (head === null) return missing(); + const ranges = new Map(); + const selected = operationIds.map(id => { + const row = db!.prepare(`SELECT ${COMMIT_COLUMNS} FROM commits WHERE operation_id = ?`) + .get(id) as unknown as Record | undefined; + if (!row) return null; + const cursor = this.decodeCommitRow(row).cursor; + const checkpoint = authorityStateCheckpointCursor(cursor); + const range = ranges.get(checkpoint); + ranges.set(checkpoint, {from: range && range.from < cursor ? range.from : cursor, + to: range && range.to > cursor ? range.to : cursor}); + return cursor; + }); + // Retain only requested receipts, not all reconstructed projections from + // every window. Even sparse requests never scan the gap between windows. + 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}); + } + } + const results = operationIds.map((id, index): AuthorityStoreReceiptResult => { + const cursor = selected[index]; + if (cursor === null) return {status: "missing"}; + const result = verified.get(id); + if (result?.status !== "found" || result.cursor !== cursor!.toString()) { + protocol("SQLite authority receipt is not part of its retained window"); + } + return result; + }); + return {status: "receipts", results}; } catch (error) { return readFailure(error); } finally { db?.close(); } } diff --git a/tests/control_plane/test_authority_archive.py b/tests/control_plane/test_authority_archive.py index 809801656f..06891ba843 100644 --- a/tests/control_plane/test_authority_archive.py +++ b/tests/control_plane/test_authority_archive.py @@ -45,6 +45,13 @@ def cli(*args, exit_code=0): verified = cli("verify", "--archive", str(archive)) assert verified["archive"] == exported["archive"] assert verified["archive"]["commits"] == "1" + audit_args = ("audit", "--goal-id", goal, "--archive", str(archive), + "--archive-sha256", verified["archive"]["archive_sha256"]) + audited = cli(*audit_args) + assert audited["status"] == "audited" + assert audited["audit"]["scope"] == "exact" + assert audited["audit"]["compared_commits"] == "1" + assert not audited["authority_changed"] occupied = cli("export", "--goal-id", goal, "--archive", str(archive), exit_code=1) assert occupied["status"] == "failed" for target_provider in ("file", "sqlite"): @@ -62,6 +69,10 @@ def cli(*args, exit_code=0): proof = json.loads((destination / "verified-restore.json").read_text()) assert proof["archive_sha256"] == verified["archive"]["archive_sha256"] assert proof["target_store_identity"] != proof["source_store_identity"] + audited = cli(*audit_args, "--destination", str(destination)) + assert audited["audit"]["status"] == "matched" + assert audited["audit"]["target_store_identity"] == proof["target_store_identity"] + assert cli(*audit_args, "--allow-newer-head")["audit"]["scope"] == "retained_prefix" # Reopen the actual restored backend in a separate process. module = REPO / f"loopx/control_plane/coordination/{target_provider}_authority_store.ts" class_name = "FileAuthorityStore" if target_provider == "file" else "SqliteAuthorityStore" diff --git a/tests/control_plane_ts/authority_archive.test.ts b/tests/control_plane_ts/authority_archive.test.ts index b07d2c3187..d4efe2f636 100644 --- a/tests/control_plane_ts/authority_archive.test.ts +++ b/tests/control_plane_ts/authority_archive.test.ts @@ -294,7 +294,7 @@ test("SQLite retained history crosses checkpoint windows and preserves an old sa } finally { await rm(root, {recursive: true, force: true}); } }); -test("a changed archive during the second pass cannot claim a verified recovery", {skip: sqliteSkip}, async () => { +test("restore consumes reviewed snapshot even when the original path is replaced", {skip: sqliteSkip}, async () => { const root = await mkdtemp(join(tmpdir(), "authority-archive-change-")); try { const source = new FileAuthorityStore(join(root, "source"), "goal"); @@ -309,7 +309,10 @@ test("a changed archive during the second pass cannot claim a verified recovery" await writeFile(archive, resign(rows)); return target.storeIdentity(); }}); - await assert.rejects(restoreAuthorityArchive(archive, wrapped, report.archive_sha256), /archive changed/); + assert.equal((await restoreAuthorityArchive(archive, wrapped, report.archive_sha256)).status, "restored"); + const restoredReceipt = await target.readReceipt("op-1"); + assert.equal(restoredReceipt.status, "found"); + if (restoredReceipt.status === "found") assert.deepEqual(restoredReceipt.receipts, [{decision: 1}]); // It remains an isolated recovery store; no source transaction was rewritten. const receipt = await source.readReceipt("op-1"); assert.equal(receipt.status, "found"); diff --git a/tests/control_plane_ts/authority_archive_audit.test.ts b/tests/control_plane_ts/authority_archive_audit.test.ts new file mode 100644 index 0000000000..121ffd1362 --- /dev/null +++ b/tests/control_plane_ts/authority_archive_audit.test.ts @@ -0,0 +1,194 @@ +import assert from "node:assert/strict"; +import {mkdtemp, readFile, rm, writeFile} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import test from "node:test"; +import type {AuthorityStore} 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 {exportAuthorityArchive, restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; +import {auditAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive_audit.ts"; +import {manageLocalAuthorityArchive} from "../../loopx/control_plane/coordination/local_authority_archive.ts"; +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; + +function observe(store: AuthorityStore, overrides: Partial): AuthorityStore { + return {providerKind: store.providerKind, storeIdentity: () => store.storeIdentity(), + loadAuthority: () => store.loadAuthority(), commitAuthority: c => store.commitAuthority(c), + readReceipt: id => store.readReceipt(id), scanCommitted: (cursor, limit) => store.scanCommitted(cursor, limit), + ...overrides}; +} +async function seed(store: AuthorityStore, count = 35, divergent = false) { + let previous: string | null = null; + for (let i = 1; i <= count; i++) { + const committed = await store.commitAuthority({expected_provider_revision: previous, + operation_id: `operation-${i}`, events: [{kind: "observed", round: i}], + next_projection: {goal_id: "goal", value: i}, + // A different old no-change receipt is invisible in the final head. + receipts: [{changed: false, decision: divergent && i === 2 ? "other" : "same"}]}); + assert.equal(committed.status, "applied"); + if (committed.status === "applied") previous = committed.provider_revision; + } + return previous; +} + +for (const provider of ["file", "sqlite"] as const) { + test(`${provider}: audit is read-only, paged, checks every receipt; resume reuses prefix`, async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + const source = new SqliteAuthorityStore(join(root, "source"), "goal"); + await seed(source); + const path = join(root, "archive"); + const archive = await exportAuthorityArchive(source, "goal", path); + const target = provider === "file" ? new FileAuthorityStore(join(root, "target"), "goal") + : new SqliteAuthorityStore(join(root, "target"), "goal"); + await restoreAuthorityArchive(path, target, archive.archive_sha256); + const before = await target.loadAuthority(); + let scans = 0, receipts = 0; + const readOnly = observe(target, { + commitAuthority: async () => { throw new Error("audit must never write"); }, + scanCommitted: (after, limit) => { scans++; assert.ok(limit <= 16); return target.scanCommitted(after, limit); }, + readReceipt: id => { receipts++; return target.readReceipt(id); }, + }); + const report = await auditAuthorityArchive(path, readOnly, archive.archive_sha256); + assert.equal(report.status, "matched"); + assert.equal(scans, 3); assert.equal(receipts, 35); + scans = 0; receipts = 0; + assert.equal((await restoreAuthorityArchive(path, readOnly, archive.archive_sha256)).status, "restored"); + assert.equal(scans, 3); assert.equal(receipts, 35); + assert.deepEqual(await target.loadAuthority(), before); + // Missing original lookup is not redeemed by identical head/scan data. + const broken = observe(target, {readReceipt: id => id === "operation-2" + ? Promise.resolve({status: "missing"}) : target.readReceipt(id)}); + const rejected = await auditAuthorityArchive(path, broken, archive.archive_sha256); + assert.equal(rejected.status, "mismatch"); + if (rejected.status !== "matched") { + assert.equal(rejected.reason_code, "archive_receipt_mismatch"); assert.equal(rejected.cursor, "2"); + } + assert.ok(!JSON.stringify(rejected).includes("operation-2")); + } finally { await rm(root, {recursive: true, force: true}); } + }); +} + +test("matching head cannot hide divergent retained decisions; no suffix is written", async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + const source = new FileAuthorityStore(join(root, "source"), "goal"); + const target = new SqliteAuthorityStore(join(root, "target"), "goal"); + await seed(source, 5); await seed(target, 3, true); + const path = join(root, "archive"); + const archive = await exportAuthorityArchive(source, "goal", path); + const before = await target.loadAuthority(); + await assert.rejects(restoreAuthorityArchive(path, target, archive.archive_sha256), /transaction_mismatch/); + assert.deepEqual(await target.loadAuthority(), before); + await seed(new FileAuthorityStore(join(root, "divergent"), "goal"), 5, true); + const other = new FileAuthorityStore(join(root, "divergent"), "goal", {existingOnly: true}); + const mismatch = await auditAuthorityArchive(path, other, archive.archive_sha256); + assert.equal(mismatch.status, "mismatch"); + if (mismatch.status !== "matched") assert.equal(mismatch.reason_code, "archive_transaction_mismatch"); + } finally { await rm(root, {recursive: true, force: true}); } +}); + +for (const scope of ["exact", "retained_prefix"] as const) { + test(`${scope}: concurrent append has explicit audit semantics`, async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + const store = new SqliteAuthorityStore(join(root, "store"), "goal"); + const previous = await seed(store, 3); + const path = join(root, "archive"); + const archive = await exportAuthorityArchive(store, "goal", path); + let appended = false; + const target = observe(store, {scanCommitted: async (after, limit) => { + if (!appended) { + appended = true; + assert.equal((await store.commitAuthority({expected_provider_revision: previous, operation_id: "later", + events: [], receipts: [], next_projection: {goal_id: "goal", value: 4}})).status, "applied"); + } + return store.scanCommitted(after, limit); + }}); + const report = await auditAuthorityArchive(path, target, archive.archive_sha256, scope); + assert.equal(report.status, scope === "exact" ? "mismatch" : "matched"); + if (report.status === "matched") assert.equal(report.compared_commits, "3"); + const exact = await auditAuthorityArchive(path, store, archive.archive_sha256); + assert.equal(exact.status, "mismatch"); + if (exact.status !== "matched") assert.equal(exact.reason_code, "archive_target_has_newer_commits"); + } finally { await rm(root, {recursive: true, force: true}); } + }); +} + +for (const fault of ["identity", "history", "receipt", "gap"] as const) { + test(`audit reports ${fault} failure without authority writes`, async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + const store = new FileAuthorityStore(join(root, "store"), "goal"); + await seed(store, 3); + const path = join(root, "archive"); + const archive = await exportAuthorityArchive(store, "goal", path); + let identityReads = 0; + const broken = observe(store, { + commitAuthority: async () => { assert.fail("audit invoked a writer"); }, + storeIdentity: () => fault === "identity" && identityReads++ > 0 + ? Promise.resolve({status: "available", store_identity: "file:" + "f".repeat(32)}) : store.storeIdentity(), + scanCommitted: async (after, limit) => { + if (fault === "history") return {status: "unavailable", reason_code: "offline", reason: "private cause"}; + const page = await store.scanCommitted(after, limit); + if (page.status === "page" && fault === "gap") page.next_cursor = "99"; + return page; + }, + readReceipt: id => fault === "receipt" + ? Promise.resolve({status: "unavailable", reason_code: "offline", reason: "private cause"}) : store.readReceipt(id), + }); + const report = await auditAuthorityArchive(path, broken, archive.archive_sha256); + assert.equal(report.status, ["receipt", "history"].includes(fault) ? "unavailable" : "mismatch"); + assert.ok(!JSON.stringify(report).includes("private cause")); + } finally { await rm(root, {recursive: true, force: true}); } + }); +} + +test("missing stores remain absent; admin rejects mismatched destination binding and digest", async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + const store = new FileAuthorityStore(join(root, "source"), "goal"); + await seed(store, 1); + const archive = join(root, "archive"); + const summary = await exportAuthorityArchive(store, "goal", archive); + const request = {schema_version: "loopx_authority_archive_admin_request_v0", action: "audit", goal_id: "goal", + archive, archive_sha256: summary.archive_sha256, runtime_root: join(root, "absent")}; + assert.equal((await manageLocalAuthorityArchive(request)).status, "failed"); + await assert.rejects(readFile(join(root, "absent", "authority", "file-v0", "store-identity")), {code: "ENOENT"}); + const destination = join(root, "restored"); + assert.equal((await manageLocalAuthorityArchive({...request, action: "restore", destination, + provider: "sqlite", execute: true})).status, "restored"); + const {runtime_root: _root, ...isolated} = request; + assert.equal((await manageLocalAuthorityArchive({...isolated, destination})).status, "audited"); + const binding = join(destination, "restore-binding.json"); + const body = JSON.parse(await readFile(binding, "utf8")); + await writeFile(binding, JSON.stringify({...body, goal_id: "other"})); + assert.equal((await manageLocalAuthorityArchive({...isolated, destination})).status, "failed"); + await assert.rejects(auditAuthorityArchive(archive, store, "0".repeat(64)), /reviewed digest/); + } finally { await rm(root, {recursive: true, force: true}); } +}); + +test("native and imported complete graphs audit all history across SQLite checkpoint boundaries", async () => { + const root = await mkdtemp(join(tmpdir(), "archive-audit-")); + try { + for (const mode of ["native", "legacy"] as const) { + const source = new SqliteAuthorityStore(join(root, mode, "source"), "goal"); + const fixture = productionScaleCoordinationFixture("goal", mode); + let previous: string | null = null; + for (let i = 1; i <= 66; i++) { + const committed = await source.commitAuthority({expected_provider_revision: previous, + operation_id: `step-${i}`, events: [{kind: "checkpoint_observation", round: i}], + next_projection: {...fixture.projection, round: i}, receipts: [{changed: false, round: i}]}); + assert.equal(committed.status, "applied"); + if (committed.status === "applied") previous = committed.provider_revision; + } + const archive = join(root, mode, "archive"); + const summary = await exportAuthorityArchive(source, "goal", archive); + // Use a reopened real SQLite store, including its retained receipt index. + const reopened = new SqliteAuthorityStore(join(root, mode, "source"), "goal", {existingOnly: true}); + const proof = await auditAuthorityArchive(archive, reopened, summary.archive_sha256); + assert.equal(proof.status, "matched"); + if (proof.status === "matched") assert.equal(proof.compared_commits, "66"); + } + } finally { await rm(root, {recursive: true, force: true}); } +}); diff --git a/tests/control_plane_ts/authority_archive_crash.test.ts b/tests/control_plane_ts/authority_archive_crash.test.ts new file mode 100644 index 0000000000..8f3938bc73 --- /dev/null +++ b/tests/control_plane_ts/authority_archive_crash.test.ts @@ -0,0 +1,70 @@ +import assert from "node:assert/strict"; +import {fork} from "node:child_process"; +import {once} from "node:events"; +import {mkdtemp, rm} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +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 {exportAuthorityArchive, restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; +import {auditAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive_audit.ts"; + +for (const provider of ["file", "sqlite"] as const) { + test(`${provider}: killed after checkpoint commit; reopen, resume, audit and continue CAS`, {timeout: 60000}, async () => { + const root = await mkdtemp(join(tmpdir(), "archive-crash-")); + const source = new SqliteAuthorityStore(join(root, "source"), "goal"); + let worker: ReturnType | undefined; + try { + let previous: string | null = null; + for (let i = 1; i <= 67; i++) { + const result = await source.commitAuthority({expected_provider_revision: previous, + operation_id: `step-${i}`, events: [{round: i}], + receipts: [{accepted: true, round: i}], next_projection: {goal_id: "goal", round: i}}); + assert.equal(result.status, "applied"); + if (result.status === "applied") previous = result.provider_revision; + } + const archive = join(root, "archive"); + const summary = await exportAuthorityArchive(source, "goal", archive); + const destination = join(root, "target"); + worker = fork(new URL("./authority_archive_restore_process.ts", import.meta.url), + [archive, destination, summary.archive_sha256, provider, "65"], + {execArgv: ["--no-warnings", "--experimental-strip-types", "--experimental-sqlite"], + env: {...process.env, TMPDIR: root, TEMP: root, TMP: root}, + stdio: ["ignore", "ignore", "pipe", "ipc"]}); + let stderr = ""; + worker.stderr!.on("data", chunk => { stderr += String(chunk); }); + const boundary = once(worker, "message"); + const ended = once(worker, "exit"); + const notification = await Promise.race([boundary, ended.then(() => { throw new Error(`worker exited before crash: ${stderr}`); })]); + assert.deepEqual(notification[0], {status: "durable", cursor: "65"}); + worker.kill("SIGKILL"); + const exit = await ended; + assert.equal(exit[1], "SIGKILL"); + const reopened = provider === "sqlite" ? new SqliteAuthorityStore(destination, "goal", {existingOnly: true}) + : new FileAuthorityStore(destination, "goal", {existingOnly: true}); + const partial = await reopened.loadAuthority(); + assert.equal(partial.status, "loaded"); + if (partial.status === "loaded") assert.equal(partial.cursor, "65"); + const restored = await restoreAuthorityArchive(archive, reopened, summary.archive_sha256); + const audited = await auditAuthorityArchive(archive, reopened, summary.archive_sha256); + assert.equal(audited.status, "matched"); + const original = await reopened.readReceipt("step-1"); + assert.equal(original.status, "found"); + if (original.status === "found") assert.deepEqual(original.receipts, [{accepted: true, round: 1}]); + assert.equal((await reopened.commitAuthority({expected_provider_revision: restored.target_provider_revision, + operation_id: "after-recovery", events: [], receipts: [{accepted: true}], + next_projection: {goal_id: "goal", round: 68}})).status, "applied"); + assert.equal((await auditAuthorityArchive(archive, reopened, summary.archive_sha256, "retained_prefix")).status, "matched"); + assert.equal((await auditAuthorityArchive(archive, reopened, summary.archive_sha256)).status, "mismatch"); + const untouched = await source.loadAuthority(); + assert.equal(untouched.status, "loaded"); + if (untouched.status === "loaded") assert.equal(untouched.cursor, "67"); + } finally { + if (worker && worker.exitCode === null && worker.signalCode === null) { + const exited = once(worker, "exit"); worker.kill("SIGKILL"); await exited; + } + await rm(root, {recursive: true, force: true}); + } + }); +} diff --git a/tests/control_plane_ts/authority_archive_postgresql.integration.test.ts b/tests/control_plane_ts/authority_archive_postgresql.integration.test.ts index f9aacc1e45..593d09ef58 100644 --- a/tests/control_plane_ts/authority_archive_postgresql.integration.test.ts +++ b/tests/control_plane_ts/authority_archive_postgresql.integration.test.ts @@ -11,6 +11,7 @@ import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlit import {installPostgreSqlAuthorityStoreSchema, PostgreSqlAuthorityStore, type PostgreSqlAuthorityDatabase} from "../../loopx/control_plane/coordination/postgresql_authority_store.ts"; import {exportAuthorityArchive, restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; +import {auditAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive_audit.ts"; import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; const url = process.env.LOOPX_TEST_POSTGRES_URL; @@ -44,6 +45,7 @@ for (const local of ["file", "sqlite"] as const) { const restored = await restoreAuthorityArchive(archive, target, exported.archive_sha256); assert.equal(restored.status, "restored"); const reopened = direction === "to-postgresql" ? new PostgreSqlAuthorityStore(database, options) : target; + assert.equal((await auditAuthorityArchive(archive, reopened, exported.archive_sha256)).status, "matched"); const rows = await reopened.scanCommitted(null, 10); assert.equal(rows.status, "page"); if (rows.status !== "page") throw new Error("readback failed"); @@ -60,6 +62,7 @@ for (const local of ["file", "sqlite"] as const) { operation_id: "isolated-recovery-probe", events: [], receipts: [{probe: true}], next_projection: {...projection, observation: 4}})).status, "applied"); await assert.rejects(restoreAuthorityArchive(archive, reopened, exported.archive_sha256), /extra commits/); + assert.equal((await auditAuthorityArchive(archive, reopened, exported.archive_sha256, "retained_prefix")).status, "matched"); const unchanged = await source.loadAuthority(); assert.equal(unchanged.status, "loaded"); if (unchanged.status === "loaded") assert.equal(unchanged.cursor, "3"); diff --git a/tests/control_plane_ts/authority_archive_restore_process.ts b/tests/control_plane_ts/authority_archive_restore_process.ts new file mode 100644 index 0000000000..6226fa43c6 --- /dev/null +++ b/tests/control_plane_ts/authority_archive_restore_process.ts @@ -0,0 +1,23 @@ +/** A real process death after durable SQLite/File commit and before readback. */ +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; +import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; +const [archive, target, digest, kind, stopAt] = process.argv.slice(2); +const store = kind === "sqlite" ? new SqliteAuthorityStore(target, "goal") : new FileAuthorityStore(target, "goal"); +const interrupted: AuthorityStore = { + providerKind: store.providerKind, + storeIdentity: () => store.storeIdentity(), loadAuthority: () => store.loadAuthority(), + readReceipt: id => store.readReceipt(id), scanCommitted: (after, limit) => store.scanCommitted(after, limit), + commitAuthority: async request => { + const committed = await store.commitAuthority(request); + if (committed.status === "applied" && committed.cursor === stopAt) { + process.send!({status: "durable", cursor: committed.cursor}); + // IPC keeps the worker alive until the parent sends SIGKILL. + await new Promise(() => {}); + } + return committed; + }, +}; +await restoreAuthorityArchive(archive, interrupted, digest); +process.exitCode = 2; // The parent requires the named crash boundary, never this path. diff --git a/tests/control_plane_ts/sqlite_authority_store.test.ts b/tests/control_plane_ts/sqlite_authority_store.test.ts index 9631010443..cb603deff5 100644 --- a/tests/control_plane_ts/sqlite_authority_store.test.ts +++ b/tests/control_plane_ts/sqlite_authority_store.test.ts @@ -361,3 +361,60 @@ test("SQLite real processes serialize CAS and preserve a lost-response receipt", assert.equal(reopened.status, "loaded"); if (reopened.status === "loaded") assert.equal(reopened.cursor, "2"); }); + + +test("SQLite receipt batches preserve scalar proofs and order across checkpoints without repeated replay", async t => { + const {store} = await fixture(t); + let revision: string | null = null; + for (let i = 1; i <= 130; i++) { + const result = await store.commitAuthority(authorityStoreCommitFixture(revision, `batch-${i}`, i, i)); + assert.equal(result.status, "applied"); if (result.status !== "applied") return; + revision = result.provider_revision; + } + const ids = ["batch-63", "batch-2", "missing", "batch-65", "batch-64", "batch-2", "batch-129"]; + const expected = await Promise.all(ids.map(id => store.readReceipt(id))); + assert.deepEqual(await store.readReceipts(ids), {status: "receipts", results: expected}); + const {DatabaseSync} = createRequire(import.meta.url)("node:sqlite"); + const prepare = DatabaseSync.prototype.prepare; + let rowsRead = 0; + DatabaseSync.prototype.prepare = function(this: import("node:sqlite").DatabaseSync, sql: string) { + const statement = prepare.call(this, sql); + if (!/FROM commits\b/i.test(sql)) return statement; + return new Proxy(statement, {get(target, property) { + const value = Reflect.get(target, property); + if (typeof value !== "function") return value; + if (property !== "all" && property !== "get") return value.bind(target); + return (...args: unknown[]) => { + const result = value.apply(target, args); + rowsRead += Array.isArray(result) ? result.length : result === undefined ? 0 : 1; + return result; + }; + }}); + }; + try { + const grouped = await store.readReceipts(Array.from({length: 16}, (_, i) => `batch-${i + 48}`)); + assert.equal(grouped.status, "receipts"); + // Sixteen indexed lookups, one <=64-row checkpoint proof, and bounded head proof. + // Replaying the same window for every receipt would exceed this by an order of magnitude. + assert.ok(rowsRead <= 90, `batch materialized ${rowsRead} retained rows`); + } finally { DatabaseSync.prototype.prepare = prepare; } + // A successful earlier batch must not cache a proof over a later disk mutation. + const db = new DatabaseSync(store.path); + db.prepare("UPDATE commits SET receipts=? WHERE cursor=2").run(JSON.stringify([{forged: true}])); + db.close(); + assert.equal((await store.loadAuthority()).status, "loaded"); + const failed = await store.readReceipts(["batch-129", "missing", "batch-2"]); + assert.equal(failed.status, "failed"); + if (failed.status === "failed") assert.equal(failed.reason_code, "provider_protocol_violation"); +}); + +test("SQLite receipt batch bounds are checked before opening storage", async t => { + const {store} = await fixture(t); + for (const ids of [[], Array(65).fill("operation"), [""]]) { + assert.equal((await store.readReceipts(ids)).status, "failed"); + } + assert.deepEqual(await store.readReceipts(["missing", "missing"]), + {status: "receipts", results: [{status: "missing"}, {status: "missing"}]}); + const {existsSync} = await import("node:fs"); + assert.equal(existsSync(store.path), false); +}); From 079ce47e3dff187a96a8708d8dc67737beb87e9f Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:18:16 +0800 Subject: [PATCH 2/4] docs(authority): reconcile recovery evidence and default cutover scopes Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../2026-09-27-recovery-audit.md | 98 +++++++++++++++++++ .../2026-09-27-recovery-audit.zh-CN.md | 76 ++++++++++++++ ...shared-goal-authority-state-provider-v0.md | 21 ++-- ...-goal-authority-state-provider-v0.zh-CN.md | 15 +-- .../typescript-control-plane-migration-v0.md | 21 ++-- ...script-control-plane-migration-v0.zh-CN.md | 15 +-- docs/reference/file-authority-state-log.md | 75 ++++++++++++++ 7 files changed, 287 insertions(+), 34 deletions(-) create mode 100644 docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md create mode 100644 docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md new file mode 100644 index 0000000000..3944614ac6 --- /dev/null +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md @@ -0,0 +1,98 @@ +# Local default cutover: recovery audit and remaining delivery scopes + +- Audited baseline: `157ab7b11`, 2026-09-27, plus this delivery. +- Owner: overall roadmap #4574 R5/G2; shared authority D2/D3; TS T3/T4. +- Supersedes the **current count**, not historical evidence, in the + [September 24 reconciliation](2026-09-24-default-cutover-reconciliation.md). + +## What is already delivered + +Complete source transport/assembly, transaction outbox capture, reviewed +promotion, canonical pagination, File checkpoint/delta format upgrade and +bounded Python prototype retirement are on main. In particular #5013, #5063, +#5102 and #5105 must not be commissioned again. Two running promoted Goals do +not prove every supported source, consumer, rollback or execution lifecycle. + +The older three-row plan grouped migration/rollback too broadly to be three +reviewable PR commitments. This audit splits its recovery prerequisite from +activation. The reason is concrete: restore verified an archive and then +reopened its mutable source pathname; a replacement could enter the isolated +target before the final digest mismatch stopped recovery. There was also no +independent read-only CLI proof of a restored store's complete retained history +and receipt lookup. These are recovery gaps, not missing capture writers. + +## Four scoped deliveries starting with this PR + +| Delivery | Observable result and remaining boundary | +| --- | --- | +| **1. This PR: reviewed restore and independent history audit** | A private verified input is consumed throughout restore. File/SQLite roundtrips preserve every logical row and original receipt. Audit checks historical transactions and receipt lookup independently, with explicit exact-head versus retained-prefix semantics. Real process death at a checkpoint can resume without rerunning acknowledged commits. SQLite batches receipt proofs in one read transaction without weakening scalar verification. No live selector/fence changes. | +| **2. External execution interval protection** | Existing lease owners supervise real Host execution, renewal, authority loss, cancellation and uncertain effects. Test expiry/reclaim while the old executor is still running. Post-execution rejection alone is insufficient. Attached Hosts without cancellation need an explicit supported boundary. | +| **3. Whole-Goal activation and rollback integration** | Reconcile #5054's retained-source inventory, then exercise source drain, saved reviewed cutover, all retained command consumers and fenced recovery/rollback together. Bind a recovered copy through an explicit transition; do not revive a source lease or overwrite later writes. Delete only Python decisions whose callers have actually moved. | +| **4. Default entrypoints and final bounded retirement** | New Goal creation, settings, installation, packaged frontend/Lark/CLI consistently use the qualified local profile. Existing Goals have explicit migration and disable/recovery paths. Remove last legacy business writers after their caller inventory and rollback constraints pass; retain rendering and Host IO. | + +This is **four planned new delivery PRs including this one, three afterwards**, +not a guarantee that no acceptance defect will require another PR. The original +three *architectural packages* are not a decrementing PR counter. This PR closes +one named recovery slice inside package 2; it does not close all of package 2. +Future checkpoints must identify which row actually completed rather than +repeating a range such as “5–8”. + +Existing PRs are separate: #5054 retires the old Todo event path and isolates +supervisor logging; #4931 optimizes SQLite retained proof reads. They were open +at the audited baseline. Do not duplicate them or request a new capture writer +for a source being retired. #4915 is filesystem placement, not authority default +selection. Further SQLite work should reuse #4224's evidence/contract and +coordinate any overlapping implementation with #4931. + +## Evidence gates are not PR allocations + +The latest #4224 formal 1 MiB report still has receipt p95 269.03 ms versus 50 ms +and scan-100 p95 801.81 ms versus 250 ms. No exact-head formal rerun on #4931 or +complete passing D2 report was present at this audit. The planned soak end date +is not an observed pass. Domain workload, steady-state RSS, large-history +recovery, consumer lag, upgrade/rollback and platform coverage remain distinct +rows. This PR's small checkpoint/crash matrix does not qualify the formal +100k/300k workload or replace ten days of natural elapsed soak. + +D1 consumer parity, D3 reviewed cohort activation and maintainer default choice +also require real evidence. A File opt-in, qualified SQLite default and migration +of all existing Goals are distinct claims. It is therefore not honest to give +an unconditional total PR count or a calendar deadline today. + +PostgreSQL reuses the same archive auditor and logical transactions. Its real +isolated store integration is exercised here; authenticated transport, tenant +policy, restore-incarnation operations, failover/pooling and capacity remain its +separate medium-term path. Local default does not require that service deployment. + +## Delivery contract + +The existing authority-archive CLI owns this administrative journey. TS owns +format verification, snapshot lifetime, historical comparison, resume decisions +and provider readback; Python only projects CLI input/output. This introduces no +capability/provider registration, storage format or new settings. Frontend and +Lark business readers continue through the same provider interfaces; no +companion configuration editor is needed. + +The negative matrix covers replaced and damaged input, old-history divergence +behind a matching head, missing or unavailable receipt lookup, incomplete pages, +incarnation changes and concurrent appends. File/SQLite process-death tests and +real PostgreSQL cross-provider tests retain actual storage. The local-source +rehearsal uses a detached byte-verified copy, never an active Goal mutation. +Public evidence excludes private Goal state and raw logs. + +[Commands, semantic changes and operational limits](../../../../reference/file-authority-state-log.md#provider-migration-and-recovery). + +The detached real-source rehearsal exposed a 300-second restore RPC timeout: +per-row receipt readback repeatedly replayed the same SQLite checkpoint window. +The bounded batch receipt method addresses that redundancy at the existing +provider port; the scalar method delegates to the same proof owner. Archive +restore and audit share page comparison, while uncertain writes force immediate +proof. This complements, rather than replaces, #4931's proof-encoding work and +neither raises the RPC budget nor qualifies D2's separate performance gates. + +Validation on this delivery: 336 native archive/SQLite conformance and crash +checks, four CLI checks, and four cross-provider checks against an isolated +PostgreSQL 16 server passed, with no skipped checks in these suites. A detached +129-commit real-source snapshot passed File and SQLite restore and independent +audit. Recovery of the earlier timed-out SQLite destination passed without +reissuing its committed operations. These checks do not claim active cutover. diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md new file mode 100644 index 0000000000..8964d3291a --- /dev/null +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md @@ -0,0 +1,76 @@ +# 本地默认切换:恢复审计与剩余交付范围 + +- 核对基线:2026-09-27 `157ab7b11`,加本次交付。 +- 归属:总目标 #4574 R5/G2;shared authority D2/D3;TS T3/T4。 +- 取代[九月二十四日清单](2026-09-24-default-cutover-reconciliation.zh-CN.md)的 + **当前数量口径**,不覆盖历史证据。 + +## 已交付的部分 + +完整来源传输/组装、事务 outbox 捕获、reviewed promotion、canonical 分页、File +checkpoint/delta 格式升级和有界 Python 原型退役已在 main。尤其 #5013、#5063、 +#5102、#5105 不能重复安排。两个 Goal 已晋升,不代表所有保留来源、消费者、回退 +和执行生命周期都通过验收。 + +此前三行计划把“迁移/回退”合得太宽,不能据此承诺三个可审查 PR。本次将恢复前置 +与正式激活分开,依据是实际缺陷:restore 验证归档后重新打开可变路径,被替换的 +内容可能先写入隔离目标,直到最终摘要不符才失败。此外缺少独立只读 CLI,证明恢复 +目标的完整历史及回执查询仍正确。这是恢复缺口,不是尚未实现事件捕获。 + +## 从本次开始的四个交付范围 + +| 交付 | 可观察结果及剩余边界 | +| --- | --- | +| **1. 本次:审核输入恢复与独立历史审计** | 恢复始终消费私有、已验证的输入副本。File/SQLite 往返保留逐笔逻辑状态和原回执。审计独立核对历史与回执查询,明确 exact 与 retained-prefix 两种语义。检查点提交后进程被杀可续传,不重复执行已提交的命令。不修改线上 selector/fence。 | +| **2. 外部执行区间保护** | 既有 lease owner 覆盖真实 Host 执行、续约、权限丢失、取消及不确定副作用。必须测试旧 executor 尚在运行时的过期/接管;结束后拒绝写回不够。不可取消的 attached Host 需明确支持边界。 | +| **3. 整 Goal 激活及回退集成** | 对齐 #5054 的保留来源清单,联合验证来源 drain、保存的 reviewed cutover、全部保留命令消费者及 fenced recovery/rollback。通过明确转换绑定恢复副本,不复活旧租约,不覆盖后来写入。仅删除 caller 已迁走的 Python 决策。 | +| **4. 默认入口与最后一批有界退役** | 新建 Goal、settings、安装及打包 frontend/Lark/CLI 一致使用合格本地 profile;存量 Goal 有明确迁移及停用/恢复路径。caller 清单与回退约束通过后,删除最后的旧业务 writer,保留渲染及 Host IO。 | + +这是**包含本次在内四个规划新 PR,本次交付后剩三个**;不保证验收不会再发现需要 +修复的缺陷。原来的三个“架构工作包”不是倒计时 PR 数。本次关闭第 2 包中的一个 +具名恢复切片,没有把整个第 2 包标为完成。后续必须指出哪行真正交付,不能再重复 +一个不变的“5–8”。 + +已有 PR 单列:#5054 退役旧 Todo event 路径并隔离 supervisor 日志,#4931 优化 +SQLite retained proof 读取。二者在本次核对时仍开放,不重复实现,也不为将退役的 +来源新增捕获 writer。#4915 属于目录布局,不能算 authority 默认切换。后续 SQLite 工作复用 #4224 的资格合同及证据,与 #4931 的重叠实现协调。 + +## 证据门不是 PR 配额 + +#4224 最新正式 1 MiB 报告仍有 receipt p95 269.03 ms / 50 ms 和 scan-100 p95 +801.81 ms / 250 ms 两项失败。本次核对未发现 #4931 精确 head 的正式复测或完整 +D2 通过报告。计划 soak 结束日期不等于实测通过。domain workload、steady-state +RSS、大历史恢复、consumer lag、升级/回退及平台覆盖是各自独立的项目。本次小型 +检查点/进程崩溃矩阵不证明正式 100k/300k 负载,也不替代十天自然时间 soak。 + +D1 消费者一致性、D3 经审核的 cohort 激活及维护者默认值选择同样需要实证。 +File opt-in、合格 SQLite 默认和全部存量 Goal 迁移是不同主张;目前不能诚实地给出 +无条件的 PR 总数或完成日期。 + +PostgreSQL 复用同一归档审计和逻辑事务,本次运行真实隔离 store 集成;认证传输、 +tenant 策略、恢复 incarnation、failover/pool 及容量仍属于独立中期路线。本地默认 +不必等待 PostgreSQL 服务部署。 + +## 本次交付合同 + +现有 authority-archive CLI 拥有此管理旅程;TS 负责格式验证、副本生命周期、历史 +比对、续传判断及 provider 回读,Python 仅适配 CLI 输入输出。不新增 capability、 +provider 注册、磁盘格式或 settings。frontend/Lark 业务读取仍走同一 provider 接口, +无需新增配置编辑器。 + +负例覆盖输入替换/损坏、最终状态相同但旧历史不同、回执查询丢失/不可用、分页不 +完整、incarnation 改变及并发追加。File/SQLite 进程崩溃测试与真实 PostgreSQL +跨 provider 测试使用实际存储。本机来源演练使用独立且核对字节的快照,不修改活跃 +Goal。公开材料排除私有 Goal 内容及原始日志。 + +[操作、语义变化及限制](../../../../reference/file-authority-state-log.md#provider-migration-and-recovery)。 + +隔离真实快照演练暴露了恢复 RPC 的 300 秒超时:逐笔回执回读反复重放同一个 SQLite +检查点窗口。本次在现有 provider 接口增加有界批量查询,标量查询委托给同一个证明 +实现;恢复和审计共享分页比对,遇到不确定写入立即回读。此处减少重复重放,与 +#4931 的证明编码优化互补;不提高 RPC 预算,也不替代 D2 的独立性能验收。 + +本次交付验证:336 项原生归档/SQLite 一致性及崩溃检查、4 项 CLI 检查、真实隔离 +PostgreSQL 16 的 4 项跨 provider 检查全部通过,这些套件没有跳过项。129 笔历史的 +独立真实来源快照通过 File、SQLite 恢复及独立审计;此前超时的 SQLite 目标也成功 +恢复,没有重发已提交操作。这些结果不表示已经完成活跃 Goal 切换。 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 e1ca420061..bd99602e12 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -24,16 +24,17 @@ [Chinese version](./shared-goal-authority-state-provider-v0.zh-CN.md) and this English version are semantic mirrors. A difference between them is a defect. -## Current delivery frontier (2026-09-25) - -Audit `37bbaec79` and current PR states: complete-source transport, transaction -capture, source assembly and the five previously open caller/event fixes are -merged, not future implementation. After the current promotion-admission repair, -three named code boundaries remain planned: external-effect execution fencing; -event-writer binding plus whole-Goal migration/rollback; default onboarding plus -bounded Python retirement. #4931 and outstanding D2 evidence are tracked -separately. Three is a delivery plan, not a guaranteed total PR count. -[Current inventory and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.md). +## Current delivery frontier (2026-09-27) + +Audit `157ab7b11` and current PR states: source capture, pagination, File format +upgrade and Python prototype retirement are delivered. This delivery repairs +reviewed-input recovery and adds independent retained-history audit. Plan four +scoped PRs starting here: this recovery slice, external execution interval +protection, whole-Goal activation/rollback integration, and default entrypoints +with final bounded Python retirement. Three planned scopes follow this PR; +existing #5054/#4931 and D2/D3 evidence remain separate. This is not a guaranteed +count of future defect repairs. +[Current inventory, rationale and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md). File retained-state storage now reuses the existing TS checkpoint/delta codec, stacked on #5063's verified read cache and RPC budgets. Original revisions, 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 3535c6da99..5b9a17dbf3 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 @@ -21,13 +21,14 @@ - 语言说明:[英文版](./shared-goal-authority-state-provider-v0.md)与本中文版互为 语义镜像;两者不一致属于缺陷 -## 当前交付边界(2026-09-25) - -按 `37bbaec79` 与当前 PR 状态核对:完整来源传输、事务捕获、来源组装及此前五个 -在途 caller/event 修复都已合入,不再计入待开发。当前晋升准入修复之后,规划三个 -明确代码边界:外部动作执行区间保护、事件 writer 绑定与整 Goal 迁移/回退闭环、 -默认启用与最后一批有界 Python 退役。#4931 与 D2 的剩余资格证据单列;三个是 -可命名的开发批次,不是保证总 PR 数。[唯一当前清单与退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.zh-CN.md)。 +## 当前交付边界(2026-09-27) + +按 `157ab7b11` 与当前 PR 核对,来源捕获、分页、File 格式升级及 Python 原型 +退役已交付。本次修复审核输入恢复并增加独立历史审计;从本次开始规划四个交付 +PR:本次恢复切片、外部执行区间保护、整 Goal 激活/回退集成、默认入口及最后 +一批有界 Python 退役。本次之后剩后三个规划范围;#5054/#4931 已有 PR,D2/D3 +缺失证据另列,不能保证最终缺陷修复数量。 +[当前清单、依据及退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md)。 File 历史存储在 #5063 的读取缓存和 RPC 预算之上,复用现有 TS checkpoint/delta 编码;物理格式升级保留原版本、回执和每条完整历史投影。正常读写只接受 v1, diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 58cb3bc440..b9b5451d05 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -29,16 +29,17 @@ required for the first App outcome. These are planned product consumers of T0–T4, not additional provider promotion or completed migration claims. -## Current delivery frontier (2026-09-25) - -Audit `37bbaec79` and current PR states: complete-source transport, transaction -capture, source assembly and the five previously open caller/event fixes are -merged, not future implementation. After the current promotion-admission repair, -three named code boundaries remain planned: external-effect execution fencing; -event-writer binding plus whole-Goal migration/rollback; default onboarding plus -bounded Python retirement. #4931 and outstanding D2 evidence are tracked -separately. Three is a delivery plan, not a guaranteed total PR count. -[Current inventory and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.md). +## Current delivery frontier (2026-09-27) + +Audit `157ab7b11` and current PR states: source capture, pagination, File format +upgrade and Python prototype retirement are delivered. This delivery repairs +reviewed-input recovery and adds independent retained-history audit. Plan four +scoped PRs starting here: this recovery slice, external execution interval +protection, whole-Goal activation/rollback integration, and default entrypoints +with final bounded Python retirement. Three planned scopes follow this PR; +existing #5054/#4931 and D2/D3 evidence remain separate. This is not a guaranteed +count of future defect repairs. +[Current inventory, rationale and exits](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md). ## Native authority qualification and prototype retirement (2026-09-26) diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index 06c222c481..dc6c7dbb6e 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -27,13 +27,14 @@ R1–R3 的 TS 消费者包括 App 产品路径,不只 CLI 结算。 这里是 T0–T4 的产品消费计划,不新增 provider promotion,也不声称迁移完成。 -## 当前交付边界(2026-09-25) - -按 `37bbaec79` 与当前 PR 状态核对:完整来源传输、事务捕获、来源组装及此前五个 -在途 caller/event 修复都已合入,不再计入待开发。当前晋升准入修复之后,规划三个 -明确代码边界:外部动作执行区间保护、事件 writer 绑定与整 Goal 迁移/回退闭环、 -默认启用与最后一批有界 Python 退役。#4931 与 D2 的剩余资格证据单列;三个是 -可命名的开发批次,不是保证总 PR 数。[唯一当前清单与退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.zh-CN.md)。 +## 当前交付边界(2026-09-27) + +按 `157ab7b11` 与当前 PR 核对,来源捕获、分页、File 格式升级及 Python 原型 +退役已交付。本次修复审核输入恢复并增加独立历史审计;从本次开始规划四个交付 +PR:本次恢复切片、外部执行区间保护、整 Goal 激活/回退集成、默认入口及最后 +一批有界 Python 退役。本次之后剩后三个规划范围;#5054/#4931 已有 PR,D2/D3 +缺失证据另列,不能保证最终缺陷修复数量。 +[当前清单、依据及退出条件](ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md)。 ## 旧观测写入退役(2026-09-24) diff --git a/docs/reference/file-authority-state-log.md b/docs/reference/file-authority-state-log.md index 7cbea5b1b3..670899bda6 100644 --- a/docs/reference/file-authority-state-log.md +++ b/docs/reference/file-authority-state-log.md @@ -111,6 +111,12 @@ loopx --format json authority-archive export --goal-id example --archive /absolu loopx --format json authority-archive verify --archive /absolute/history.ndjson loopx --format json authority-archive restore --goal-id example --archive /absolute/history.ndjson \ --destination /absolute/new-isolated-store --provider sqlite --archive-sha256 DIGEST --execute +# Independently audit the actual restored provider, not just its old success file. +loopx --format json authority-archive audit --goal-id example --archive /absolute/history.ndjson \ + --archive-sha256 DIGEST --destination /absolute/new-isolated-store +# Audit the selected runtime's retained prefix after it has received later writes. +loopx --format json authority-archive audit --goal-id example --archive /absolute/history.ndjson \ + --archive-sha256 DIGEST --allow-newer-head ``` File/SQLite exports restore into either File or SQLite. Restore creates a new @@ -120,6 +126,52 @@ converter for every pair of storage formats. PostgreSQL's existing archive source contract remains unchanged; authenticated service activation, cutover and PostgreSQL destination administration are separate work. +Restore first copies the input into a private temporary directory and verifies +that copy against the reviewed digest. All subsequent reads use that same copy. +Previously, replacing the original path between verification and restore could +write unreviewed rows into the isolated destination before the final seal check +failed. Replacement or in-place edits to the original now cannot change the +accepted restore input. A corrupt or wrong-digest copy fails before target +access. Normal completion/failure removes the temporary directory; abrupt process +death can leave private temporary files for normal system/operator cleanup. +Budget additional local disk space approximately equal to the archive size. + +Resume checks the target's entire existing prefix before appending its missing +suffix. Prefix scans use at most 16 reconstructed transactions per page; each +original operation is also looked up through the provider's receipt index. The +same TS comparator serves restore and independent audit. Acknowledged new writes +also share a page proof; an uncertain write forces immediate readback before any +later write, without retrying the uncertain operation. Logical commits remain +individual CAS operations. SQLite resolves all requested operation IDs in one +read transaction and verifies each touched checkpoint window once per batch; +other providers use the same contract with scalar receipt lookup as a fallback. +This removes redundant replay without skipping old receipts, caching proof +across calls, or claiming constant memory independent of state size. +Input decoding retains one reconstructed state plus operation-ID uniqueness +tracking; target pages retain up to 16 full states. This is count-bounded paging, +not a new provider byte-budget guarantee. + +`audit` is read-only and never creates a missing authority. It compares every +historical projection, event, operation ID and original receipt, including +receipt lookup cursor/version. Physical provider revisions may differ across +providers. A matching final head alone cannot pass. Reports expose a typed +reason and first failing cursor, not private state or receipt bodies. + +- Default `exact` requires the target to contain exactly the archive's commits + and remain at that head through readback. Concurrent advancement asks for a + retry; it is not silently classified as equality. +- Explicit `--allow-newer-head` verifies only the archive's retained prefix and + tolerates later append-only commits in the same store lineage. The report says + `retained_prefix`; it does not certify those later commits against the backup. +- `--destination` reads the isolated restore binding to choose File or SQLite; + without it, audit uses the selected runtime provider. A wrong binding, source + digest, missing historical receipt or unavailable provider fails closed. + +An audit proves retained data, not active execution safety. It does not transfer +leases, select a provider, remove a writer fence, or authorize rollback over +newer writes. `verified-restore.json` remains a historical completion receipt; +run `audit` for current evidence rather than trusting the file's presence. + Old raw backups can be copied into an **isolated** provider directory, with their original identity and canonical filename, then upgraded and exported. Never rewrite a schema label or restore an old backup over newer acknowledged writes. @@ -127,6 +179,29 @@ Binary rollback requires the target runtime to pass `--require-current`; an old binary without that gate is not automatically activated. Data rollback and provider cutover require their own reviewed, fenced recovery operation. +## 恢复与审计 + +恢复先复制到权限受限的临时目录,再核对审核摘要,后续只读这一份副本。旧实现 +在校验后重开原路径,可能先将被替换的内容写入隔离目标、最后才报错;现在原路径 +被替换或原地改写不会改变已接受的恢复输入。正常结束会清理副本;进程被强制杀掉 +可能留下私有临时文件。需预留约一份归档大小的额外磁盘空间。 + +续传先核对已有前缀,再追加缺失后缀;TS 的同一比对规则同时服务恢复与独立审计。 +最多每页 16 条完整历史投影,逐笔核对原始回执。新写入仍是独立 CAS,但明确成功的 +提交可共享分页回读;遇到不确定写入,立即回读证明后才可继续,绝不盲目重发。 +SQLite 在一个读事务内按操作 ID 查询回执,同一批只重放一次每个涉及的检查点窗口; +其他 provider 通过同一接口回退到逐笔查询。不跨请求缓存证明,不跳过旧回执,也不 +声称内存与状态大小无关。 + +`authority-archive audit` 不写业务状态、不创建缺失存储。默认 exact 要求当前目标 +与归档完整一致;显式 `--allow-newer-head` 只证明归档对应的保留历史前缀,允许随后 +追加,但不把之后的提交算成已审核。`--destination` 审计隔离恢复目录,否则审计 +当前选择的 runtime provider。逐笔状态、事件、操作身份及回执查询都必须相符; +最终 head 一样不足以通过。不同 provider 的物理 CAS 版本无需相同。 + +审计不选择 authority、不转移执行租约、不撤销 writer fence,也不授权覆盖新写入。 +`verified-restore.json` 是过去的完成回执,当前完整性应重新 audit。 + ## Qualification and limits Tests cover legacy rejection, original historical identity/receipts, checkpoint From aa8833a3a79c3e8e2277ca1060761a5a0358ad9c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 04:05:11 +0800 Subject: [PATCH 3/4] test(authority): keep the crash worker alive until SIGKILL The restore crash worker awaited a bare promise after its durable boundary. A pending promise does not hold the event loop, so the child could exit 13 ("unfinished top-level await") before the parent delivered SIGKILL, and the File case then asserted null instead of a signal. Pin the IPC channel with the same message-listener barrier the task-lease crash worker uses, assert that SIGKILL was actually delivered, and hold the parent 250ms past the durable message so the late-signal race stays covered. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../control_plane_ts/authority_archive_crash.test.ts | 8 +++++++- .../authority_archive_restore_process.ts | 11 +++++++++-- 2 files changed, 16 insertions(+), 3 deletions(-) diff --git a/tests/control_plane_ts/authority_archive_crash.test.ts b/tests/control_plane_ts/authority_archive_crash.test.ts index 8f3938bc73..288891365c 100644 --- a/tests/control_plane_ts/authority_archive_crash.test.ts +++ b/tests/control_plane_ts/authority_archive_crash.test.ts @@ -10,6 +10,11 @@ import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlit import {exportAuthorityArchive, restoreAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive.ts"; import {auditAuthorityArchive} from "../../loopx/control_plane/coordination/authority_archive_audit.ts"; +// The worker must survive the durable boundary until this late SIGKILL. An +// immediate signal hides a worker that already exited 13 ("unfinished +// top-level await") because its IPC channel stopped holding the event loop. +const crashSignalDelayMs = 250; + for (const provider of ["file", "sqlite"] as const) { test(`${provider}: killed after checkpoint commit; reopen, resume, audit and continue CAS`, {timeout: 60000}, async () => { const root = await mkdtemp(join(tmpdir(), "archive-crash-")); @@ -38,7 +43,8 @@ for (const provider of ["file", "sqlite"] as const) { const ended = once(worker, "exit"); const notification = await Promise.race([boundary, ended.then(() => { throw new Error(`worker exited before crash: ${stderr}`); })]); assert.deepEqual(notification[0], {status: "durable", cursor: "65"}); - worker.kill("SIGKILL"); + await new Promise(resolve => setTimeout(resolve, crashSignalDelayMs)); + assert.equal(worker.kill("SIGKILL"), true, "the crash worker was still alive to receive SIGKILL"); const exit = await ended; assert.equal(exit[1], "SIGKILL"); const reopened = provider === "sqlite" ? new SqliteAuthorityStore(destination, "goal", {existingOnly: true}) diff --git a/tests/control_plane_ts/authority_archive_restore_process.ts b/tests/control_plane_ts/authority_archive_restore_process.ts index 6226fa43c6..b6290fb026 100644 --- a/tests/control_plane_ts/authority_archive_restore_process.ts +++ b/tests/control_plane_ts/authority_archive_restore_process.ts @@ -5,6 +5,13 @@ import {restoreAuthorityArchive} from "../../loopx/control_plane/coordination/au import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; const [archive, target, digest, kind, stopAt] = process.argv.slice(2); const store = kind === "sqlite" ? new SqliteAuthorityStore(target, "goal") : new FileAuthorityStore(target, "goal"); +// Hold the IPC channel across the durable boundary. A pending promise alone +// leaves the event loop empty, so the child exits 13 ("unfinished top-level +// await") before the parent can deliver SIGKILL. The message listener pins the +// channel, the same barrier the task-lease crash worker uses. +let resume: () => void = () => {}; +const crashBarrier = new Promise(resolve => {resume = resolve;}); +process.on("message", () => resume()); const interrupted: AuthorityStore = { providerKind: store.providerKind, storeIdentity: () => store.storeIdentity(), loadAuthority: () => store.loadAuthority(), @@ -13,8 +20,8 @@ const interrupted: AuthorityStore = { const committed = await store.commitAuthority(request); if (committed.status === "applied" && committed.cursor === stopAt) { process.send!({status: "durable", cursor: committed.cursor}); - // IPC keeps the worker alive until the parent sends SIGKILL. - await new Promise(() => {}); + // The parent requires the named crash boundary and signals SIGKILL here. + await crashBarrier; } return committed; }, From f59651df79aabfd5e1624dd77e308af1c8ceb3e7 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 05:34:33 +0800 Subject: [PATCH 4/4] docs(authority): regenerate the RFC status index for the recovery ledger entry The generator cannot run while the App conversation RFC header is malformed; after that header was repaired on main, regenerating the index records the shared-goal-authority ledger entry added by this branch (19 -> 20 entries). Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/architecture/rfcs/STATUS.md | 2 +- docs/architecture/rfcs/STATUS.zh-CN.md | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/architecture/rfcs/STATUS.md b/docs/architecture/rfcs/STATUS.md index b4594cd7fc..9b35c7c9d3 100644 --- a/docs/architecture/rfcs/STATUS.md +++ b/docs/architecture/rfcs/STATUS.md @@ -57,7 +57,7 @@ appendix may keep dated history, but no dated log heading may precede it. | [RFC: Research Exploration Control Plane v0](research-exploration-control-plane-v0.md) | Accepted | none | — | | [RFC: Semantic Vocabulary Convergence and Commit-Time Drift Checks (v0)](semantic-vocabulary-convergence-v0.md) | Accepted | none | [5 entries](ledger/semantic-vocabulary-convergence-v0/) | | [RFC: Shared Goal Alignment and Governed Amendment Protocol (v0)](shared-goal-alignment-and-governed-amendment-v0.md) | Accepted | none | [2 entries](ledger/shared-goal-alignment-and-governed-amendment-v0/) | -| [RFC: LoopX Shared Control-Plane Authority and Pluggable State Providers (v0)](shared-goal-authority-state-provider-v0.md) | Accepted | none | [19 entries](ledger/shared-goal-authority-state-provider-v0/) | +| [RFC: LoopX Shared Control-Plane Authority and Pluggable State Providers (v0)](shared-goal-authority-state-provider-v0.md) | Accepted | none | [20 entries](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | Accepted | none | — | | [RFC: TypeScript Control-Plane Migration Direction v0](typescript-control-plane-migration-v0.md) | Accepted | none | [12 entries](ledger/typescript-control-plane-migration-v0/) | diff --git a/docs/architecture/rfcs/STATUS.zh-CN.md b/docs/architecture/rfcs/STATUS.zh-CN.md index acd3b20f59..4812ff4404 100644 --- a/docs/architecture/rfcs/STATUS.zh-CN.md +++ b/docs/architecture/rfcs/STATUS.zh-CN.md @@ -54,7 +54,7 @@ | [RFC:研究型探索控制面 v0](research-exploration-control-plane-v0.zh-CN.md) | 已接受 | 无 | — | | [RFC:语义词表收敛与提交期漂移检查(v0)](semantic-vocabulary-convergence-v0.zh-CN.md) | 已接受 | 无 | [5 条](ledger/semantic-vocabulary-convergence-v0/) | | [RFC:共享 Goal 对齐与受治理 Amendment 协议(v0)](shared-goal-alignment-and-governed-amendment-v0.zh-CN.md) | 已接受 | 无 | [2 条](ledger/shared-goal-alignment-and-governed-amendment-v0/) | -| [RFC:LoopX 共享控制面权威与可插拔状态 Provider(v0)](shared-goal-authority-state-provider-v0.zh-CN.md) | 已接受 | 无 | [19 条](ledger/shared-goal-authority-state-provider-v0/) | +| [RFC:LoopX 共享控制面权威与可插拔状态 Provider(v0)](shared-goal-authority-state-provider-v0.zh-CN.md) | 已接受 | 无 | [20 条](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | 已接受 | none | — | | [RFC:LoopX 控制面 TypeScript 渐进迁移方向 v0](typescript-control-plane-migration-v0.zh-CN.md) | 已接受 | 无 | [12 条](ledger/typescript-control-plane-migration-v0/) |