From ebe29c5afeed0b4c563d024c6c3d9665a921192f Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:21:23 +0800 Subject: [PATCH 1/4] feat(authority): add reviewed File and SQLite provider cutover Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/authority_archive.py | 22 +- .../authority_format_inspection.ts | 5 +- .../coordination/file_authority_store.ts | 16 +- .../coordination/local_authority_archive.ts | 2 + .../coordination/local_authority_migration.ts | 225 +++++++++++++++ .../coordination/local_authority_provider.ts | 81 +++--- tests/control_plane/test_authority_archive.py | 34 +++ .../authority_format_upgrade.test.ts | 10 +- .../local_authority_migration.test.ts | 256 ++++++++++++++++++ .../local_authority_migration_process.ts | 38 +++ 10 files changed, 645 insertions(+), 44 deletions(-) create mode 100644 loopx/control_plane/coordination/local_authority_migration.ts create mode 100644 tests/control_plane_ts/local_authority_migration.test.ts create mode 100644 tests/control_plane_ts/local_authority_migration_process.ts diff --git a/loopx/cli_commands/authority_archive.py b/loopx/cli_commands/authority_archive.py index 49c379a1a0..48d19f9efb 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, audit or restore a canonical authority copy." + "authority-archive", help="Back up, audit, restore or migrate canonical local authority." ) add_subcommand_format(parser) actions = parser.add_subparsers(dest="authority_archive_action", required=True) @@ -27,6 +27,15 @@ 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 ("plan-migration", "migrate"): + action = actions.add_parser(name, help="Review or execute a quiescent File/SQLite provider migration.") + action.add_argument("--goal-id", required=True) + action.add_argument("--plan", type=Path, required=True) + if name == "plan-migration": + action.add_argument("--provider", choices=("file", "sqlite"), required=True) + else: + action.add_argument("--plan-sha256", required=True) + action.add_argument("--execute", action="store_true", help="Publish the verified provider; otherwise preview.") for name in ("export", "verify", "restore", "audit"): action = actions.add_parser(name) action.add_argument("--archive", type=Path, required=True) @@ -63,6 +72,14 @@ def handle_authority_archive_command( elif args.authority_archive_action == "upgrade": request.update(runtime_roots=authority_upgrade_roots( registry_path, runtime_root_arg, all_known=args.all_known), execute=args.execute) + elif args.authority_archive_action in {"plan-migration", "migrate"}: + request.update(goal_id=args.goal_id, plan=str(args.plan.expanduser().resolve()), + runtime_root=str(resolve_runtime_root(load_registry(registry_path), runtime_root_arg, + registry_path=registry_path))) + if args.authority_archive_action == "plan-migration": + request["provider"] = args.provider + else: + request.update(plan_sha256=args.plan_sha256, execute=args.execute) else: request["archive"] = str(args.archive.expanduser().resolve()) if args.authority_archive_action == "export": @@ -89,7 +106,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.')}\n" + f"{value.get('reason', 'Authority changed: ' + str(value.get('authority_changed')))}\n" + f"Plan digest: {value.get('plan_sha256', 'not applicable')}\n" f"{value.get('audit', '')}" )) return 1 if result.get("status") == "failed" else 0 diff --git a/loopx/control_plane/coordination/authority_format_inspection.ts b/loopx/control_plane/coordination/authority_format_inspection.ts index ee04a16476..d750dff82a 100644 --- a/loopx/control_plane/coordination/authority_format_inspection.ts +++ b/loopx/control_plane/coordination/authority_format_inspection.ts @@ -10,6 +10,7 @@ import {FILE_AUTHORITY_JOURNAL_SCHEMA} from "./file_authority_journal.ts"; import {sqliteAuthorityRuntime} from "./sqlite_runtime.ts"; import {SQLITE_AUTHORITY_STORE_SCHEMA, sqliteAuthorityPath} from "./sqlite_authority_store.ts"; import {SQLITE_AUTHORITY_STORE_V1_SCHEMA} from "./sqlite_authority_migration.ts"; +import {decodeLocalAuthoritySelection} from "./local_authority_provider.ts"; import {verifyAuthorityArchive} from "./authority_archive.ts"; type StoreInspection = { @@ -87,9 +88,7 @@ export async function inspectAuthorityFormat(path: string): Promise { try { const identity = await readFile(this.identityPath, "utf8"); - if (!STORE_IDENTITY_PATTERN.test(identity)) { - throw new AuthorityStoreProtocolError("store identity does not match file:<32 lowercase hex>"); + if (!STORE_IDENTITY_PATTERN.test(identity) || + (this.expectedIdentity !== undefined && identity !== this.expectedIdentity)) { + throw new AuthorityStoreProtocolError("store identity is invalid or differs from the expected File lineage"); } if (createIfMissing) await syncAuthorityDirectory(this.directory); return identity; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } - if (!createIfMissing) throw new FileStoreUnavailableError("existing store identity is missing"); + if (!createIfMissing || this.expectedIdentity !== undefined) throw new FileStoreUnavailableError("existing store identity is missing"); return await withFileMutationLock(this.identityPath, async () => { try { const identity = await readFile(this.identityPath, "utf8"); - if (!STORE_IDENTITY_PATTERN.test(identity)) { - throw new AuthorityStoreProtocolError("store identity does not match file:<32 lowercase hex>"); + if (!STORE_IDENTITY_PATTERN.test(identity) || + (this.expectedIdentity !== undefined && identity !== this.expectedIdentity)) { + throw new AuthorityStoreProtocolError("store identity is invalid or differs from the expected File lineage"); } await syncAuthorityDirectory(this.directory); return identity; diff --git a/loopx/control_plane/coordination/local_authority_archive.ts b/loopx/control_plane/coordination/local_authority_archive.ts index b51d9977e5..5678b08960 100644 --- a/loopx/control_plane/coordination/local_authority_archive.ts +++ b/loopx/control_plane/coordination/local_authority_archive.ts @@ -1,5 +1,6 @@ /** Administrative archive transport. Large private state stays in local files; * the managed effect runtime returns only compact integrity/readback facts. */ +import {manageLocalAuthorityMigration} from "./local_authority_migration.ts"; import {inspectAuthorityFormat} from "./authority_format_inspection.ts"; import {upgradeAuthorityFormats} from "./authority_format_upgrade.ts"; import {mkdir, readFile} from "node:fs/promises"; @@ -35,6 +36,7 @@ export async function manageLocalAuthorityArchive(value: unknown, } return {...base, ...await upgradeAuthorityFormats(request.runtime_roots as string[], request.execute === true)}; } + if (request.action === "plan-migration" || request.action === "migrate") return await manageLocalAuthorityMigration(request); const archive = path(request.archive, "archive path"); if (request.action === "verify") return {...base, status: "verified", archive: await verifyAuthorityArchive(archive)}; const goalId = requireAuthorityStoreId(request.goal_id, "goal id"); diff --git a/loopx/control_plane/coordination/local_authority_migration.ts b/loopx/control_plane/coordination/local_authority_migration.ts new file mode 100644 index 0000000000..26478b3a2a --- /dev/null +++ b/loopx/control_plane/coordination/local_authority_migration.ts @@ -0,0 +1,225 @@ +/** Reviewed local-provider cutover. The canonical writer guard serializes the + * copy/audit/publication boundary; the archive preserves domain history. */ +import {randomUUID} from "node:crypto"; +import {mkdir, open, readFile, realpath} from "node:fs/promises"; +import {dirname, isAbsolute, join} from "node:path"; +import type {JsonObject} from "../effect_program.ts"; +import {durableWriteJson} from "../effect_runtime_io.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import type {AuthorityStore} from "./authority_store.ts"; +import {canonicalAuthoritySha256 as sha256, hasExactAuthorityKeys, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {exportAuthorityArchive, restoreAuthorityArchive, verifyAuthorityArchive, type AuthorityArchiveSummary} from "./authority_archive.ts"; +import {auditAuthorityArchive} from "./authority_archive_audit.ts"; +import {indexCoordinationProjection} from "./coordination_projection.ts"; +import {canonicalTaskLease} from "./task_lease_state.ts"; +import {loadLegacyCoordinationWriterFence} from "./legacy_writer_fence.ts"; +import {withCanonicalWriter} from "./local_authority_write.ts"; +import {FileAuthorityStore, syncAuthorityDirectory} from "./file_authority_store.ts"; +import {SqliteAuthorityStore} from "./sqlite_authority_store.ts"; +import {localAuthorityProviderPaths, openLocalAuthorityStoreHandle, publishLocalAuthoritySelection, + requireLocalAuthorityRuntimeRoot} from "./local_authority_provider.ts"; + +type Provider = "file" | "sqlite"; +interface Source extends JsonObject { + provider: Provider; + store_identity: string; + provider_revision: string; + cursor: string; + projection_sha256: string; + fence_sha256: string; +} +interface Plan extends JsonObject { + schema_version: "loopx_local_authority_migration_plan_v0"; + plan_id: string; + runtime_root: string; + goal_id: string; + target_provider: Provider; + source: Source; +} +interface Recovery extends JsonObject { + schema_version: "loopx_local_authority_migration_recovery_v0"; + plan_sha256: string; + phase: "prepared" | "completed"; + target_store_identity: string; + archive_sha256: string; +} +const HEX = /^[0-9a-f]{64}$/; +function provider(value: unknown): Provider { + if (value !== "file" && value !== "sqlite") throw new Error("Local migration requires file or sqlite"); + return value; +} +function absolute(value: unknown): string { + if (typeof value !== "string" || !isAbsolute(value)) throw new Error("Migration plan path must be absolute"); + return value; +} +async function optionalJson(path: string): Promise { + try { return JSON.parse(await readFile(path, "utf8")); } + catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; throw error; } +} +async function writePlan(path: string, plan: Plan): Promise { + // Never replace the artifact an operator already reviewed. + const file = await open(path, "wx", 0o600); + try { await file.writeFile(JSON.stringify(plan) + "\n"); await file.sync(); } + finally { await file.close(); } + await syncAuthorityDirectory(dirname(path)); +} +function decodePlan(value: unknown): Plan { + const plan = requireJsonObject(value, "migration plan"); + if (!hasExactAuthorityKeys(plan, ["schema_version", "plan_id", "runtime_root", "goal_id", "target_provider", "source"]) || + plan.schema_version !== "loopx_local_authority_migration_plan_v0" || + typeof plan.plan_id !== "string" || !/^[0-9a-f-]{36}$/.test(plan.plan_id)) throw new Error("Invalid migration plan"); + const source = requireJsonObject(plan.source, "migration source"); + const sourceProvider = provider(source.provider); + const target = provider(plan.target_provider); + if (!hasExactAuthorityKeys(source, ["provider", "store_identity", "provider_revision", "cursor", "projection_sha256", "fence_sha256"]) || + sourceProvider === target || typeof source.store_identity !== "string" || + !new RegExp(`^${sourceProvider}:[0-9a-f]{32}$`).test(source.store_identity) || + typeof source.provider_revision !== "string" || !source.provider_revision || + typeof source.cursor !== "string" || !/^[1-9]\d*$/.test(source.cursor) || + typeof source.projection_sha256 !== "string" || !HEX.test(source.projection_sha256) || + typeof source.fence_sha256 !== "string" || !HEX.test(source.fence_sha256)) throw new Error("Invalid migration source binding"); + absolute(plan.runtime_root); requireAuthorityStoreId(plan.goal_id, "goal id"); + return plan as Plan; +} +function decodeRecovery(value: unknown, plan: Plan, digest: string): Recovery { + const record = requireJsonObject(value, "migration recovery"); + if (!hasExactAuthorityKeys(record, ["schema_version", "plan_sha256", "phase", "target_store_identity", "archive_sha256"]) || + record.schema_version !== "loopx_local_authority_migration_recovery_v0" || record.plan_sha256 !== digest || + (record.phase !== "prepared" && record.phase !== "completed") || + typeof record.target_store_identity !== "string" || + !new RegExp(`^${plan.target_provider}:[0-9a-f]{32}$`).test(record.target_store_identity) || + typeof record.archive_sha256 !== "string" || !HEX.test(record.archive_sha256)) throw new Error("Invalid migration recovery binding"); + return record as Recovery; +} +async function observe(root: string, goalId: string): Promise<{source: Source; store: AuthorityStore}> { + const fence = await loadLegacyCoordinationWriterFence(root, goalId); + if (fence.status !== "loaded") throw new Error("Migration requires an engaged legacy writer fence"); + const handle = await openLocalAuthorityStoreHandle(root, goalId, {}, {existingOnly: true}); + const selected = provider(handle.provider); + const identity = await handle.store.storeIdentity(); + const head = await handle.store.loadAuthority(); + if (identity.status !== "available" || head.status !== "loaded") throw new Error("Migration requires an available canonical head"); + const index = indexCoordinationProjection(head.head, goalId); + for (const [todoId, raw] of index.leases) { + // Expiry permits another lease decision, but does not prove a Host stopped. + if (canonicalTaskLease(raw, goalId, todoId).status === "active") { + throw new Error("Migration requires settled task leases, including expired active leases; stop writers and settle leases first"); + } + } + return {store: handle.store, source: {provider: selected, store_identity: identity.store_identity, + provider_revision: head.provider_revision, cursor: head.cursor, projection_sha256: sha256(head.head), fence_sha256: sha256(fence.fence)}}; +} +function checkArchive(archive: AuthorityArchiveSummary, plan: Plan): void { + const source = plan.source; + if (archive.goal_id !== plan.goal_id || archive.source_provider !== source.provider || + archive.source_store_identity !== source.store_identity || archive.source_provider_revision !== source.provider_revision || + archive.commits !== source.cursor || archive.projection_sha256 !== source.projection_sha256) { + throw new Error("Migration archive does not match the reviewed source"); + } +} +function targetStore(root: string, plan: Plan, identity?: string): AuthorityStore { + const paths = localAuthorityProviderPaths(root, plan.goal_id); + const options = identity === undefined ? {} : {existingOnly: true, expectedIdentity: identity}; + return plan.target_provider === "file" ? new FileAuthorityStore(paths.file, plan.goal_id, identity === undefined ? {} : {expectedIdentity: identity}) + : new SqliteAuthorityStore(paths.sqlite, plan.goal_id, options); +} + +/** No provider factory injection: this administrative operation explicitly owns + * the built-in local stores, and never silently adopts a PostgreSQL binding. */ +export async function manageLocalAuthorityMigration(request: JsonObject): Promise { + const base = {schema_version: "loopx_local_authority_migration_result_v0", legacy_fallback_used: false, + execution_authority_granted: false}; + // Publication may have succeeded even if fsync/readback throws. Report unknown, + // not a false assertion that no authority changed; retry the same plan. + let publicationAttempted = false; + try { + const root = await realpath(requireLocalAuthorityRuntimeRoot(request.runtime_root)); + const goalId = requireAuthorityStoreId(request.goal_id, "goal id"); + const planPath = absolute(request.plan); + return await withCanonicalWriter(root, goalId, false, async () => { + if (request.action === "plan-migration") { + const target = provider(request.provider); + const {source} = await observe(root, goalId); + if (source.provider === target) throw new Error("Requested provider is already selected"); + const plan: Plan = {schema_version: "loopx_local_authority_migration_plan_v0", plan_id: randomUUID(), + runtime_root: root, goal_id: goalId, target_provider: target, source}; + await writePlan(planPath, plan); + return {...base, status: "planned", authority_changed: false, plan_sha256: sha256(plan), plan, + requires_execute: true}; + } + if (request.action !== "migrate" || typeof request.execute !== "boolean") throw new Error("Invalid migration action"); + const plan = decodePlan(await optionalJson(planPath)); + const digest = sha256(plan); + if (request.plan_sha256 !== digest || plan.runtime_root !== root || plan.goal_id !== goalId) { + throw new Error("Migration plan digest, runtime or Goal mismatch"); + } + const directory = join(root, "authority-transition", "local-provider", digest); + const recoveryPath = join(directory, "recovery.json"); + const archivePath = join(directory, "source.archive.jsonl"); + const raw = await optionalJson(recoveryPath); + let recovery = raw === undefined ? undefined : decodeRecovery(raw, plan, digest); + const current = await openLocalAuthorityStoreHandle(root, goalId, {}, {existingOnly: true}); + if (current.provider === plan.target_provider && recovery) { + const archive = await verifyAuthorityArchive(archivePath); + checkArchive(archive, plan); + const proof = await auditAuthorityArchive(archivePath, current.store, recovery.archive_sha256, "retained_prefix"); + if (proof.status !== "matched" || proof.target_store_identity !== recovery.target_store_identity) { + throw new Error("Published migration target no longer matches its retained history or identity"); + } + if (request.execute) await durableWriteJson(recoveryPath, {...recovery, phase: "completed"}); + return {...base, status: "already_applied", authority_changed: false, plan_sha256: digest, + selected_provider: current.provider, audit: proof}; + } + if (recovery?.phase === "completed") throw new Error("Completed migration was superseded; create a fresh plan"); + const source = await observe(root, goalId); + if (sha256(source.source) !== sha256(plan.source)) throw new Error("Migration source changed; create and review a fresh plan"); + if (!request.execute) return {...base, status: "planned", authority_changed: false, plan_sha256: digest, + selected_provider: source.source.provider, target_provider: plan.target_provider, requires_execute: true}; + await mkdir(directory, {recursive: true, mode: 0o700}); + // Bind implicit File selection before preparing another provider. Otherwise + // SQLite existence would intentionally trip the missing-selector guard. + await publishLocalAuthoritySelection(root, goalId, source.source.provider, source.source.store_identity); + let archive: AuthorityArchiveSummary; + try { archive = await verifyAuthorityArchive(archivePath); } + catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT" || recovery) throw error; + archive = await exportAuthorityArchive(source.store, goalId, archivePath); + } + checkArchive(archive, plan); + if (!recovery) { + const identity = await targetStore(root, plan).storeIdentity(); + if (identity.status !== "available") throw new Error("Migration target identity unavailable"); + recovery = {schema_version: "loopx_local_authority_migration_recovery_v0", plan_sha256: digest, + phase: "prepared", target_store_identity: identity.store_identity, archive_sha256: archive.archive_sha256}; + await durableWriteJson(recoveryPath, recovery); + } + if (archive.archive_sha256 !== recovery.archive_sha256) throw new Error("Migration backup changed"); + const target = targetStore(root, plan, recovery.target_store_identity); + await restoreAuthorityArchive(archivePath, target, recovery.archive_sha256); + const proof = await auditAuthorityArchive(archivePath, target, recovery.archive_sha256); + if (proof.status !== "matched" || proof.target_store_identity !== recovery.target_store_identity) throw new Error("Migration target audit failed"); + const finalSource = await observe(root, goalId); + if (sha256(finalSource.source) !== sha256(plan.source)) throw new Error("Migration source changed before publication"); + publicationAttempted = true; + await publishLocalAuthoritySelection(root, goalId, plan.target_provider, recovery.target_store_identity); + const selected = await openLocalAuthorityStoreHandle(root, goalId, {}, {existingOnly: true}); + const readback = await selected.store.loadAuthority(); + const identity = await selected.store.storeIdentity(); + // The independent full audit above is still protected by the same writer + // guard. Publication readback binds that proof to the selected head; do + // not re-scan the entire archive a second time for the identical revision. + if (selected.provider !== plan.target_provider || identity.status !== "available" || + identity.store_identity !== recovery.target_store_identity || readback.status !== "loaded" || + readback.cursor !== proof.captured_target_cursor || readback.provider_revision !== proof.captured_target_provider_revision || + sha256(readback.head) !== plan.source.projection_sha256) throw new Error("Migration publication readback failed; retry the same plan"); + await durableWriteJson(recoveryPath, {...recovery, phase: "completed"}); + return {...base, status: "migrated", authority_changed: true, plan_sha256: digest, + selected_provider: selected.provider, archive: {...archive}, audit: proof, + publication_readback: {cursor: readback.cursor, provider_revision: readback.provider_revision}}; + }); + } catch (error) { + return {...base, status: "failed", authority_changed: publicationAttempted ? null : false, + reason_code: publicationAttempted ? "migration_publication_uncertain" : "migration_rejected", + reason: error instanceof Error ? error.message : "Migration unavailable"}; + } +} diff --git a/loopx/control_plane/coordination/local_authority_provider.ts b/loopx/control_plane/coordination/local_authority_provider.ts index 1364202641..5a654b1808 100644 --- a/loopx/control_plane/coordination/local_authority_provider.ts +++ b/loopx/control_plane/coordination/local_authority_provider.ts @@ -82,7 +82,7 @@ export function localAuthorityOpenFailure(error: unknown): Record, goalId: strin }; } -async function openSelectedSqlite(root: string, goalId: string, storeIdentity: string): Promise { - const p = paths(root, goalId); +export function decodeLocalAuthoritySelection(config: unknown, goalId: string): + LocalPostgreSqlAuthoritySelection | {schema_version: typeof SCHEMA; provider: "file" | "sqlite"; goal_id: string; store_identity: string} { + if (!isAuthorityJsonObject(config) || config.schema_version !== SCHEMA || config.goal_id !== goalId || + (config.provider !== "file" && config.provider !== "sqlite" && config.provider !== "postgresql")) { + throw selectorError("Invalid local authority provider selector"); + } + if (config.provider === "postgresql") return decodePostgreSqlSelection(config, goalId); + if (!hasExactAuthorityKeys(config, ["schema_version", "provider", "goal_id", "store_identity"]) || + typeof config.store_identity !== "string" || !new RegExp(`^${config.provider}:[0-9a-f]{32}$`).test(config.store_identity)) { + throw selectorError("Invalid local authority provider selector"); + } + return {schema_version: SCHEMA, provider: config.provider, goal_id: goalId, store_identity: config.store_identity}; +} + +/** Both local providers reject a missing or replaced selected lineage. */ +async function openSelectedLocal(root: string, goalId: string, provider: "file" | "sqlite", storeIdentity: string): Promise { + const source = sourceFor(provider); + const directory = localAuthorityProviderPaths(root, goalId)[provider]; + const Store = provider === "file" ? FileAuthorityStore : SqliteAuthorityStore; try { - // Validate metadata before comparing selector lineage so drift has its own recovery signal. - const store = new SqliteAuthorityStore(p.sqlite, goalId, {existingOnly: true}); - // Do not recreate a lost database and thereby silently reset its lineage. + const store = new Store(directory, goalId, {existingOnly: true}); try { await stat(store.path); } catch (error) { - if ((error as NodeJS.ErrnoException).code === "ENOENT") { - throw new LocalAuthorityProviderOpenError("sqlite_v0", "local_authority_provider_missing", "Selected SQLite authority database is missing"); - } - throw new LocalAuthorityProviderOpenError("sqlite_v0", "local_authority_provider_open_failed", "Selected SQLite authority database could not be opened"); + const missing = (error as NodeJS.ErrnoException).code === "ENOENT"; + throw new LocalAuthorityProviderOpenError(source, + missing ? "local_authority_provider_missing" : "local_authority_provider_open_failed", + missing ? `Selected ${provider === "sqlite" ? "SQLite authority database" : "File authority document"} is missing` : `Selected ${provider} authority could not be opened`); } const identity = await store.storeIdentity(); if (identity.status !== "available") { - throw new LocalAuthorityProviderOpenError("sqlite_v0", "local_authority_provider_open_failed", identity.reason, identity.reason_code); + throw new LocalAuthorityProviderOpenError(source, "local_authority_provider_open_failed", identity.reason, identity.reason_code); } if (identity.store_identity !== storeIdentity) { - throw new LocalAuthorityProviderOpenError("sqlite_v0", "local_authority_provider_identity_mismatch", "Selected SQLite authority identity changed"); + throw new LocalAuthorityProviderOpenError(source, "local_authority_provider_identity_mismatch", `Selected ${provider} authority identity changed`); } - // Keep every subsequent operation fenced to the selected lineage. - return new SqliteAuthorityStore(p.sqlite, goalId, {existingOnly: true, expectedIdentity: storeIdentity}); + // The binding is checked again on every operation, including cached reads. + return new Store(directory, goalId, {existingOnly: true, expectedIdentity: storeIdentity}); } catch (error) { if (error instanceof LocalAuthorityProviderOpenError) throw error; - throw new LocalAuthorityProviderOpenError("sqlite_v0", "local_authority_provider_open_failed", "Selected SQLite authority could not be opened"); + throw new LocalAuthorityProviderOpenError(source, "local_authority_provider_open_failed", `Selected ${provider} authority could not be opened`); } } +/** Called under the Goal's canonical writer guard, only after the target's + * complete history/receipts were independently verified. No implicit fallback. */ +export async function publishLocalAuthoritySelection(root: string, goalId: string, + provider: "file" | "sqlite", storeIdentity: string): Promise { + if (!new RegExp(`^${provider}:[0-9a-f]{32}$`).test(storeIdentity)) throw selectorError("Invalid local store identity"); + await durableWriteJson(localAuthorityProviderPaths(root, goalId).marker, + {schema_version: SCHEMA, provider, goal_id: goalId, store_identity: storeIdentity}); +} + async function openSelectedPostgreSql( selection: LocalPostgreSqlAuthoritySelection, dependencies: LocalAuthorityProviderDependencies, @@ -197,7 +221,7 @@ export async function openLocalAuthorityStoreHandle( dependencies: LocalAuthorityProviderDependencies = {}, options: {existingOnly?: boolean} = {}, ): Promise { - const p = paths(root, goalId); + const p = localAuthorityProviderPaths(root, goalId); let raw: string; try { raw = await readFile(p.marker, "utf8"); } catch (error) { @@ -218,21 +242,13 @@ export async function openLocalAuthorityStoreHandle( let config: unknown; try { config = JSON.parse(raw); } catch { throw new LocalAuthorityProviderOpenError(null, "local_authority_selector_invalid", "Invalid local authority provider selector JSON"); } - if (!isAuthorityJsonObject(config) || config.schema_version !== SCHEMA || config.goal_id !== goalId || - (config.provider !== "sqlite" && config.provider !== "postgresql")) { - throw selectorError("Invalid local authority provider selector"); - } - if (config.provider === "sqlite") { - if (!hasExactAuthorityKeys(config, ["schema_version", "provider", "goal_id", "store_identity"]) || - typeof config.store_identity !== "string" || !/^sqlite:[0-9a-f]{32}$/.test(config.store_identity)) { - throw selectorError("Invalid SQLite authority provider selector"); - } - const store = await openSelectedSqlite(root, goalId, config.store_identity); - return {store, provider: "sqlite", sourceAuthority: sourceFor("sqlite")}; + const selection = decodeLocalAuthoritySelection(config, goalId); + if (selection.provider === "postgresql") { + const store = await openSelectedPostgreSql(selection, dependencies); + return {store, provider: "postgresql", sourceAuthority: sourceFor("postgresql")}; } - const selection = decodePostgreSqlSelection(config, goalId); - const store = await openSelectedPostgreSql(selection, dependencies); - return {store, provider: "postgresql", sourceAuthority: sourceFor("postgresql")}; + const store = await openSelectedLocal(root, goalId, selection.provider, selection.store_identity); + return {store, provider: selection.provider, sourceAuthority: sourceFor(selection.provider)}; } export async function openLocalAuthorityStore( @@ -246,10 +262,11 @@ export async function openLocalAuthorityStore( /** Administrative opt-in for an empty, unpromoted goal; no implicit migration. */ export async function selectLocalSqliteAuthority(root: string, goalId: string, execute: boolean) { - const p = paths(root, goalId); + const p = localAuthorityProviderPaths(root, goalId); return withFileMutationLock(shadowMaintenanceLockPath(root, goalId), async () => { if (existsSync(p.marker)) { - await openLocalAuthorityStore(root, goalId); + const selected = await openLocalAuthorityStoreHandle(root, goalId); + if (selected.provider !== "sqlite") throw new Error("Provider selection cannot replace an existing authority; use a reviewed migration"); return {ok: true, provider: "sqlite", changed: false, executed: execute}; } const fence = await loadLegacyCoordinationWriterFence(root, goalId); diff --git a/tests/control_plane/test_authority_archive.py b/tests/control_plane/test_authority_archive.py index 06891ba843..f104f9a184 100644 --- a/tests/control_plane/test_authority_archive.py +++ b/tests/control_plane/test_authority_archive.py @@ -157,3 +157,37 @@ def test_all_known_upgrade_roots_are_registry_owned_and_do_not_create_stores(tmp assert all(p.read_bytes() == data for p, data in before.items()) assert not (project / ".loopx/runtime").exists() assert not (common / "authority").exists() + + +@pytest.mark.parametrize("initial,target", [("file", "sqlite"), ("sqlite", "file")]) +def test_provider_migration_cli_preserves_readback_and_rollback(tmp_path, monkeypatch, initial, target): + from canonical_authority_fixture import promoted_create_fixture + + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, state = promoted_create_fixture(tmp_path, provider=initial) + before = state.read_bytes() + + def cli(*args, expected=0): + result = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), + "--format", "json", *args], cwd=REPO, capture_output=True, text=True, timeout=90) + assert result.returncode == expected, result.stdout + result.stderr + return json.loads(result.stdout) + + try: + for index, destination in enumerate((target, initial)): + plan = tmp_path / f"plan-{index}.json" + planned = cli("authority-archive", "plan-migration", "--goal-id", "goal-a", "--provider", destination, "--plan", str(plan)) + arguments = ("authority-archive", "migrate", "--goal-id", "goal-a", "--plan", str(plan), "--plan-sha256", planned["plan_sha256"]) + assert cli(*arguments)["status"] == "planned" + migrated = cli(*arguments, "--execute") + assert migrated["status"] == "migrated", migrated + assert migrated["selected_provider"] == destination + assert migrated["audit"]["compared_commits"] == "1" + assert cli(*arguments, "--execute")["status"] == "already_applied" + # Independent production CLI reads must follow the published selector. + snapshot = cli("authority-archive", "export", "--goal-id", "goal-a", "--archive", str(tmp_path / f"after-{index}.jsonl")) + assert snapshot["archive"]["source_provider"] == destination + assert state.read_bytes() == before + finally: + subprocess.run([sys.executable, "-c", "from loopx.control_plane.effect_runtime import effect_runtime_result; effect_runtime_result('runtime.shutdown',{},retry_safe=False)"], + cwd=REPO, capture_output=True, text=True, timeout=30, check=True) diff --git a/tests/control_plane_ts/authority_format_upgrade.test.ts b/tests/control_plane_ts/authority_format_upgrade.test.ts index 7848a81a8e..e06dda5dec 100644 --- a/tests/control_plane_ts/authority_format_upgrade.test.ts +++ b/tests/control_plane_ts/authority_format_upgrade.test.ts @@ -128,10 +128,18 @@ test("content inspection separates store, archive, backup and selector regardles await assert.rejects(inspectAuthorityFormat(backup), /digest mismatch/); const selector = join(runtime, "selector.json"); await writeFile(selector, JSON.stringify({schema_version: "loopx_local_authority_provider_v0", - provider: "postgresql", goal_id: goal, tenant_id: "tenant-a", store_identity: "postgresql:synthetic"})); + provider: "postgresql", goal_id: goal, tenant_id: "tenant-a", store_identity: "postgresql:" + "a".repeat(32)})); const selected = await inspectAuthorityFormat(selector); assert.equal(selected.artifact_kind, "provider_selector"); assert.equal(selected.verification, "metadata_only"); + for (const provider of ["file", "sqlite"]) { + const binding = {schema_version: "loopx_local_authority_provider_v0", provider, goal_id: goal, + store_identity: provider + ":" + "a".repeat(32)}; + await writeFile(selector, JSON.stringify(binding)); + assert.equal((await inspectAuthorityFormat(selector)).provider, provider); + await writeFile(selector, JSON.stringify({...binding, store_identity: "wrong-provider"})); + await assert.rejects(inspectAuthorityFormat(selector), /selector/); + } await writeFile(misleading, JSON.stringify({schema_version: "loopx_file_authority_store_v999"})); assert.equal((await inspectAuthorityFormat(misleading)).status, "unsupported"); }); diff --git a/tests/control_plane_ts/local_authority_migration.test.ts b/tests/control_plane_ts/local_authority_migration.test.ts new file mode 100644 index 0000000000..6477c0d94e --- /dev/null +++ b/tests/control_plane_ts/local_authority_migration.test.ts @@ -0,0 +1,256 @@ +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 {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {durableWriteJson} from "../../loopx/control_plane/effect_runtime_io.ts"; +import {manageLocalAuthorityArchive as manage} from "../../loopx/control_plane/coordination/local_authority_archive.ts"; +import {openLocalAuthorityStoreHandle as selected, localAuthorityProviderPaths} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; +import {legacyCoordinationWriterFencePath, LEGACY_COORDINATION_WRITER_FENCE_SCHEMA} from "../../loopx/control_plane/coordination/legacy_writer_fence.ts"; +import {withCanonicalWriter} from "../../loopx/control_plane/coordination/local_authority_write.ts"; +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import {inspectAuthorityFormat} from "../../loopx/control_plane/coordination/authority_format_inspection.ts"; + +const goal = "migration-goal"; +const request = (root: string, fields: JsonObject) => ({schema_version: "loopx_authority_archive_admin_request_v0", runtime_root: root, goal_id: goal, ...fields}); +async function append(root: string, i: number, leases: JsonObject[] = []) { + return withCanonicalWriter(root, goal, false, async () => { + const {store} = await selected(root, goal); + const head = await store.loadAuthority(); + const projection = authorityProjectionFixture(goal, [{todo_id: "todo_a", role: "agent", status: "open", done: false, + text: `Independent observation ${i}`, archive_state: "active"}], leases, "native", {generation: i}); + const result = await store.commitAuthority({operation_id: `op-${i}`, expected_provider_revision: head.status === "loaded" ? head.provider_revision : null, + next_projection: projection, events: [{kind: "observation", sequence: i}], receipts: [{accepted: true, sequence: i}]}); + assert.equal(result.status, "applied", JSON.stringify(result)); + return result; + }); +} +async function fixture(t: {after: (f: () => Promise) => void}) { + const root = await mkdtemp(join(tmpdir(), "local-migration-")); + t.after(() => rm(root, {recursive: true, force: true})); + await append(root, 1); await append(root, 2); + await durableWriteJson(legacyCoordinationWriterFencePath(root, goal), { + schema_version: LEGACY_COORDINATION_WRITER_FENCE_SCHEMA, state: "engaged", goal_id: goal, + fence_id: "fixture-fence", source_version: "fixture-source", source_projection_sha256: "a".repeat(64), + expected_shadow_provider_revision: "fixture-revision"}); + return root; +} +async function plan(root: string, target: "file" | "sqlite", name = target) { + const path = join(root, `${name}-plan.json`); + const result = await manage(request(root, {action: "plan-migration", provider: target, plan: path})); + assert.equal(result.status, "planned", JSON.stringify(result)); + return {action: "migrate", plan: path, plan_sha256: result.plan_sha256, execute: true}; +} + +test("real providers: File → SQLite → File preserves history, receipts and later writes", async t => { + const root = await fixture(t); + const original = await (await selected(root, goal)).store.scanCommitted(null, 10); + const first = await plan(root, "sqlite"); + const preview = await manage(request(root, {...first, execute: false})); + assert.equal(preview.status, "planned"); + assert.equal((await selected(root, goal)).provider, "file"); + assert.equal((await manage(request(root, first))).status, "migrated"); + assert.equal((await selected(root, goal)).provider, "sqlite"); + await append(root, 3); + const resumed = await manage(request(root, first)); + assert.equal(resumed.status, "already_applied", JSON.stringify(resumed)); + assert.equal((resumed.audit as JsonObject).captured_target_cursor, "3"); + const back = await plan(root, "file"); + assert.equal((await manage(request(root, back))).status, "migrated"); + const active = await selected(root, goal); + assert.equal(active.provider, "file"); + assert.equal((await inspectAuthorityFormat(localAuthorityProviderPaths(root, goal).marker)).provider, "file"); + const copy = await active.store.scanCommitted(null, 10); + assert.equal(copy.status, "page"); + if (copy.status !== "page" || original.status !== "page") throw Error("missing history"); + assert.deepEqual(copy.transactions.slice(0, 2), original.transactions); + assert.deepEqual(copy.transactions[2].receipts, [{accepted: true, sequence: 3}]); + await append(root, 4); + const retired = await manage(request(root, first)); + assert.equal(retired.status, "failed"); + assert.match(String(retired.reason), /superseded/); + assert.equal((await active.store.loadAuthority()).status, "loaded"); +}); + +test("reviewed plans reject stale source, wrong digest and another runtime without switching", async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + assert.equal((await manage(request(root, {...p, plan_sha256: "0".repeat(64)}))).status, "failed"); + await append(root, 3); + const stale = await manage(request(root, p)); + assert.match(String(stale.reason), /source changed/); + const other = await fixture(t); + const wrongRoot = await manage(request(other, p)); + assert.match(String(wrongRoot.reason), /runtime or Goal mismatch/); + assert.equal((await selected(root, goal)).provider, "file"); +}); + +test("active leases, including expired leases, require actual settlement before planning", async t => { + const root = await fixture(t); + for (const [i, expiry] of [[3, 0], [4, 9999999999]]) { + await append(root, i, [{todo_id: "todo_a", status: "active", owner: "worker", idempotency_key: "lease-key", + version: 1, lease_epoch: 1, expires_at: expiry}]); + const result = await manage(request(root, {action: "plan-migration", provider: "sqlite", plan: join(root, `lease-${i}.json`)})); + assert.equal(result.status, "failed"); + assert.match(String(result.reason), /settled task leases/); + } + await append(root, 5, [{todo_id: "todo_a", status: "released", owner: "worker", idempotency_key: "lease-key", version: 2, lease_epoch: 1}]); + await plan(root, "sqlite"); +}); + +test("missing legacy fence and divergent target never publish a provider", async t => { + const root = await fixture(t); + const fence = legacyCoordinationWriterFencePath(root, goal); + const saved = await readFile(fence); + await rm(fence); + const denied = await manage(request(root, {action: "plan-migration", provider: "sqlite", plan: join(root, "no-fence.json")})); + assert.match(String(denied.reason), /writer fence/); + await writeFile(fence, saved); + const p = await plan(root, "sqlite"); + // Explicitly retain the source binding while staging an unrelated target. + const paths = localAuthorityProviderPaths(root, goal); + const identity = await (await selected(root, goal)).store.storeIdentity(); + assert.equal(identity.status, "available"); + if (identity.status !== "available") throw Error("identity"); + await durableWriteJson(paths.marker, {schema_version: "loopx_local_authority_provider_v0", provider: "file", goal_id: goal, store_identity: identity.store_identity}); + const target = new SqliteAuthorityStore(paths.sqlite, goal); + await target.commitAuthority({operation_id: "unrelated", expected_provider_revision: null, events: [], receipts: [], next_projection: {goal_id: goal, other: true}}); + const failed = await manage(request(root, p)); + assert.equal(failed.status, "failed"); + assert.equal(failed.authority_changed, false); + assert.equal((await selected(root, goal)).provider, "file"); +}); + +test("File handles fence their store identity after opening", async t => { + const root = await fixture(t); + const paths = localAuthorityProviderPaths(root, goal); + const original = await new FileAuthorityStore(paths.file, goal).storeIdentity(); + if (original.status !== "available") throw Error("identity"); + await durableWriteJson(paths.marker, {schema_version: "loopx_local_authority_provider_v0", provider: "file", goal_id: goal, store_identity: original.store_identity}); + const handle = await selected(root, goal); + await writeFile(join(paths.file, "store-identity"), "file:" + "f".repeat(32)); + assert.notEqual((await handle.store.loadAuthority()).status, "loaded"); + await assert.rejects(selected(root, goal), /identity/); +}); + +// Kept separate from the production owner: fault hooks instrument real effects +// only in a subprocess that the parent actually kills. +import {fork} from "node:child_process"; +import {once} from "node:events"; +async function crash(root: string, p: JsonObject, boundary: "restoring" | "published", beforeKill?: () => Promise) { + const child = fork(new URL("./local_authority_migration_process.ts", import.meta.url), + [JSON.stringify(request(root, p)), boundary], {execArgv: ["--no-warnings", "--experimental-strip-types", "--experimental-sqlite"], stdio: ["ignore", "ignore", "pipe", "ipc"]}); + let errors = ""; + child.stderr?.on("data", data => { errors += data; }); + const timeout = setTimeout(() => child.kill("SIGKILL"), 15000); + try { + const [message] = await Promise.race([once(child, "message"), once(child, "exit").then(() => { throw Error(errors || "worker exited before boundary"); })]); + assert.deepEqual(message, {boundary}); + await beforeKill?.(); + const exited = once(child, "exit"); + child.kill("SIGKILL"); await exited; + } finally { clearTimeout(timeout); child.kill("SIGKILL"); } +} +for (const boundary of ["restoring", "published"] as const) { + test(`SIGKILL ${boundary}: retries retain receipts and selected writers reopen the right provider`, async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + await crash(root, p, boundary); + const current = await selected(root, goal); + assert.equal(current.provider, boundary === "published" ? "sqlite" : "file"); + if (boundary === "published") await append(root, 3); + const retry = await manage(request(root, p)); + assert.equal(retry.status, boundary === "published" ? "already_applied" : "migrated", JSON.stringify(retry)); + const result = await (await selected(root, goal)).store.readReceipt("op-1"); + assert.equal(result.status, "found"); + assert.deepEqual(result.status === "found" ? result.receipts : [], [{accepted: true, sequence: 1}]); + }); +} +test("a queued canonical writer continues on the selected provider after publication", async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + // Kill after publication but before readback/completion; the writer must recover + // the process-owned lock and use SQLite, not retain a pre-lock File handle. + let pending: Promise | undefined; + await crash(root, p, "published", async () => { + let finished = false; + pending = append(root, 3).then(value => { finished = true; return value; }); + await new Promise(resolve => setTimeout(resolve, 60)); + assert.equal(finished, false, "canonical writer must wait for the live migration lock"); + }); + await pending; + const resumed = await manage(request(root, p)); + assert.equal(resumed.status, "already_applied"); + assert.equal((resumed.audit as JsonObject).captured_target_cursor, "3"); +}); +test("replacement of a prepared target identity is rejected without republishing source", async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + await crash(root, p, "restoring"); + const paths = localAuthorityProviderPaths(root, goal); + await rm(paths.sqlite, {recursive: true}); + const replaced = new SqliteAuthorityStore(paths.sqlite, goal); + await replaced.storeIdentity(); + const retry = await manage(request(root, p)); + assert.equal(retry.status, "failed", JSON.stringify(retry)); + assert.equal((await selected(root, goal)).provider, "file"); + assert.equal((await replaced.loadAuthority()).status, "missing"); +}); +test("a source write after interrupted copy invalidates the old plan, a fresh plan can reuse the verified prefix", async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + await crash(root, p, "restoring"); + await append(root, 3); + const stale = await manage(request(root, p)); + assert.match(String(stale.reason), /source changed/); + const fresh = await plan(root, "sqlite", "fresh"); + assert.equal((await manage(request(root, fresh))).status, "migrated"); + assert.equal((await (await selected(root, goal)).store.readReceipt("op-3")).status, "found"); +}); + +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; +test("mixed production-scale graph keeps archived Todos, decisions and lease generations through both providers", async t => { + const root = await fixture(t); + const mixed = productionScaleCoordinationFixture(goal, "native").projection; + const quiescent = {...mixed, leases: (mixed.leases as JsonObject[]).map(lease => ({...lease, status: "released"}))}; + const source = (await selected(root, goal)).store; + let head = await source.loadAuthority(); + for (let i = 0; i < 3; i++) { + assert.equal(head.status, "loaded"); + if (head.status !== "loaded") throw Error("head"); + const result = await source.commitAuthority({operation_id: `mixed-${i}`, expected_provider_revision: head.provider_revision, + events: [{kind: "observation", iteration: i}], receipts: [{accepted: i}], next_projection: {...quiescent, observation: i}}); + assert.equal(result.status, "applied"); head = await source.loadAuthority(); + } + const original = await source.scanCommitted(null, 10); + if (original.status !== "page") throw Error("history"); + for (const target of ["sqlite", "file"] as const) { + const p = await plan(root, target); + const result = await manage(request(root, p)); + assert.equal(result.status, "migrated", JSON.stringify(result)); + const copy = await (await selected(root, goal)).store.scanCommitted(null, 10); + if (copy.status !== "page") throw Error("history"); + assert.deepEqual(copy.transactions.map(({provider_revision, ...row}) => row), + original.transactions.map(({provider_revision, ...row}) => row)); + } +}); + + +test("post-publication IO failure reports uncertainty, same-plan retry reads the durable result", async t => { + const root = await fixture(t); + const p = await plan(root, "sqlite"); + const child = fork(new URL("./local_authority_migration_process.ts", import.meta.url), + [JSON.stringify(request(root, p)), "publication-error"], + {execArgv: ["--no-warnings", "--experimental-strip-types", "--experimental-sqlite"], stdio: ["ignore", "ignore", "inherit", "ipc"]}); + const exited = once(child, "exit"); + const [message] = await once(child, "message"); + await exited; + assert.equal(message.unexpected.status, "failed"); + assert.equal(message.unexpected.authority_changed, null); + assert.equal(message.unexpected.reason_code, "migration_publication_uncertain"); + assert.equal((await selected(root, goal)).provider, "sqlite"); + assert.equal((await manage(request(root, p))).status, "already_applied"); +}); diff --git a/tests/control_plane_ts/local_authority_migration_process.ts b/tests/control_plane_ts/local_authority_migration_process.ts new file mode 100644 index 0000000000..e20911e452 --- /dev/null +++ b/tests/control_plane_ts/local_authority_migration_process.ts @@ -0,0 +1,38 @@ +/** Real-process fault boundaries: SIGKILL leaves durable provider/selector bytes + * and the actual writer lock behind, rather than simulating an in-memory error. */ +import fs from "node:fs/promises"; +import {syncBuiltinESMExports} from "node:module"; +import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {localAuthorityProviderPaths} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; +import {manageLocalAuthorityArchive} from "../../loopx/control_plane/coordination/local_authority_archive.ts"; + +const request = JSON.parse(process.argv[2]); +request.runtime_root = await fs.realpath(request.runtime_root); +const boundary = process.argv[3]; +async function suspend() { + process.send?.({boundary}); + setInterval(() => {}, 1000); + await new Promise(() => {}); +} +if (boundary === "restoring") { + const commit = SqliteAuthorityStore.prototype.commitAuthority; + SqliteAuthorityStore.prototype.commitAuthority = async function (value) { + const result = await commit.call(this, value); + if (result.status === "applied") await suspend(); + return result; + }; +} else if (boundary === "published" || boundary === "publication-error") { + const rename = fs.rename; + const marker = localAuthorityProviderPaths(request.runtime_root, request.goal_id).marker; + fs.rename = async (from, to) => { + await rename(from, to); + if (to === marker && JSON.parse(await fs.readFile(marker, "utf8")).provider === "sqlite") { + if (boundary === "publication-error") throw new Error("Injected failure after selector rename"); + await suspend(); + } + }; + syncBuiltinESMExports(); +} +const result = await manageLocalAuthorityArchive(request); +process.send?.({unexpected: result}); +process.disconnect?.(); From 6c910ef1f481c4e741b8597f409de8e1565d3eda Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:21:23 +0800 Subject: [PATCH 2/4] docs(authority): document cutover recovery and reconcile remaining delivery Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../2026-09-27-recovery-audit.md | 23 ++++++ .../2026-09-27-recovery-audit.zh-CN.md | 28 +++++++- docs/reference/file-authority-state-log.md | 71 +++++++++++++++++++ 3 files changed, 121 insertions(+), 1 deletion(-) 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 index 1b8f81653e..56d7db82c6 100644 --- 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 @@ -129,3 +129,26 @@ the yielding and in-flight proof lifecycle belong to File. #4931's digest window remains a separate optimization. The regression uses a private real server and the existing mixed Todo/lease/decision fixture; production locators, active Goals and raw evidence are never modified or published. + +## Local provider cutover (`76ff7c73c`) + +Recovery/audit #5140 and runtime fairness #5156 are merged. This delivery adds +reviewed File ↔ SQLite cutover for an already-promoted canonical Goal: verified +backup, exact source/fence binding, identity-bound target, history/receipt audit, +serialized selector publication, crash resume and reverse migration carrying +new writes. Persisted active leases hold migration even after expiry. It does +not stop Hosts or migrate independently-owned Turn/spend state. + +The current inventory is three existing open PRs (#5054 retirement, #4931 SQLite +proof encoding, #5144 managed Host supervision), this cutover PR, and remaining +whole-Goal integration/default-entry scopes. Whole-Goal acceptance still needs +source drain and all retained consumers/external execution boundaries; defaults +still need new-Goal/settings/install adoption and bounded Python writer removal. +Do not subtract one from a broad work package merely because its cutover subitem +shipped. An exact remaining PR count cannot be promised until that integration +inventory and D2 results determine whether additional bounded fixes are needed. +Capacity, platform coverage and natural-time soak remain evidence gates rather +than PR quotas. PostgreSQL keeps the shared logical archive/audit contract; +service authentication, tenancy and failover are not qualified by this command. + +See [the reviewed cutover journey](../../../../reference/file-authority-state-log.md#reviewed-filesqlite-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 index 6c2fbf7705..5760bea60a 100644 --- 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 @@ -1,6 +1,7 @@ # 本地默认切换:恢复审计与剩余交付范围 -- 核对基线:2026-09-27 `157ab7b11`,加本次交付。 +- 最新核对基线:`76ff7c73c`;下方历史段落保留其当时基线。 +- 当前数量以末尾“本地 provider 切换”表为准,历史规划不是剩余 PR 倒计时。 - 归属:总目标 #4574 R5/G2;shared authority D2/D3;TS T3/T4。 - 取代[九月二十四日清单](2026-09-24-default-cutover-reconciliation.zh-CN.md)的 **当前数量口径**,不覆盖历史证据。 @@ -97,3 +98,28 @@ p95 或跨平台资格。单笔巨大事务、JSON 解析、其他同步 handler 让步及进行中证明的生命周期归 File;#4931 的 digest window 仍是独立优化。 回归使用私有真实 server 和既有混合 Todo/lease/decision fixture,不改生产 locator 及活跃 Goal,不发布原始证据。 + +## 本地 provider 切换(`76ff7c73c`) + +#5140 恢复审计、#5156 共享运行时延迟修复已合入,不能重复列为未完成。 +本次交付“已晋升 canonical Goal 的 File ↔ SQLite 审核切换”:备份并核对完整历史与 +原回执,绑定 source revision/fence 和 target identity,串行发布 selector,支持进程 +中断后的续传,以及携带最新历史的反向迁移。未结算租约(包括过期 active)拦截。 +此处不自动停止 Host,不迁移 Turn/spend 的独立状态,也不等于全部旧 Goal 晋升。 + +| 当前交付范围 | PR / 状态 | 仍需证明的结果 | +| --- | --- | --- | +| 旧 Todo events 退役与 supervisor 日志隔离 | 已有 #5054,开放 | 消费者迁走后的旧分支删除 | +| SQLite retained proof 编码 | 已有 #4931,开放 | 在 #4224 冻结负载上的正式复测,不以小型迁移耗时替代 | +| 受管 Host 执行区间保护 | 已有 #5144,开放 | 续约/取消/旧 executor 接管;attached Host 边界另行明确 | +| 整 Goal 激活与回退集成 | 本次交付其中的本地 provider 切换子项 | 旧来源 drain、全部保留消费者和外部执行状态的组合验收仍未关闭 | +| 默认入口与有界 Python 退役 | 尚未实现的后续范围 | 新 Goal、设置、安装和各入口采用合格 profile;只删除 caller 已迁走的业务 writer | + +因此当前可确定的是 **3 个已有开放 PR、当前 1 个切换 PR,以及上述剩余集成/默认 +入口范围**。本次没有把宽泛的“整 Goal”行直接勾完,也没有据此将总数机械减一。 +只有补齐消费者清单和 D2 实测后,才能判断剩余集成可合成一个 PR,还是需按具体 +失败拆分;目前不能准确承诺“再 N 个就全量切换”。容量/平台/自然时间 soak 是独立 +证据门,不是编码 PR 配额。PostgreSQL 仍复用共享逻辑历史与回执审计,但本地切换 +入口明确不接受 PostgreSQL,服务认证/tenant/failover 不在此处偷换为已完成。 + +操作与恢复边界见[审核切换](../../../../reference/file-authority-state-log.md#reviewed-filesqlite-cutover)。 diff --git a/docs/reference/file-authority-state-log.md b/docs/reference/file-authority-state-log.md index 670899bda6..0689f2fe84 100644 --- a/docs/reference/file-authority-state-log.md +++ b/docs/reference/file-authority-state-log.md @@ -234,3 +234,74 @@ platform CI evidence; POSIX validation does not substitute for it. 升级失败时可能已有部分 store 完成,必须据实报告并重试,不能覆盖之后产生的写入。 备份仍是旧格式,恢复时应先在隔离目录升级。只回退二进制并不等于安全回退数据。 该方案减少重复存储和后续写入耗时,但冷校验仍验证全历史,可能更慢。 + +## Reviewed File/SQLite cutover + +An already-promoted, quiescent canonical Goal can change between the built-in +File and SQLite providers. Stop its writers and settle its task leases first. +An expired but still `active` lease is a hold: expiry does not prove that the +old Host stopped. This command does not stop Hosts, settle Turns, retire leases, +change a registry, or remove the legacy writer fence. PostgreSQL service +activation and migration from a legacy Markdown Goal are separate operations. + +```bash +# Use the same explicit runtime root for planning, preview and execution. +loopx --runtime-root /absolute/runtime --format json authority-archive plan-migration \ + --goal-id example --provider sqlite --plan /absolute/migration-plan.json +# Review the saved plan and take PLAN_SHA256 from the planning response. +loopx --runtime-root /absolute/runtime --format json authority-archive migrate \ + --goal-id example --plan /absolute/migration-plan.json --plan-sha256 PLAN_SHA256 +loopx --runtime-root /absolute/runtime --format json authority-archive migrate \ + --goal-id example --plan /absolute/migration-plan.json --plan-sha256 PLAN_SHA256 --execute +``` + +The plan binds the canonical runtime path, Goal, source store identity, revision, +cursor, projection digest, writer-fence digest and destination provider. A +changed source requires a new reviewed plan. Plans are created exclusively; +choose a new path instead of overwriting an already-reviewed artifact. + +Execution shares the maintenance guard used by canonical command writers. It +saves a verified logical backup below the runtime's +`authority-transition/local-provider//`, binds the target identity, +restores missing history, independently audits every retained transaction and +receipt, then atomically publishes the selector. Selected local stores check +their identity on every operation. A missing/replaced selected store fails +closed; it is never recreated or silently replaced with another provider. + +An implicit File source becomes an explicit identity-bound File selection before +SQLite preparation. This keeps the source readable while the target is being +built. It does not change the source's domain history. Backup and target creation +consume additional disk space; retain the backup and recovery record for retries +and investigation. No automatic cleanup deletes these recovery artifacts. + +Retry **the same plan and digest** after interruption. A partial restore audits +its prefix before appending. If publication already happened, recovery audits +the retained prefix without discarding later target writes. A publication error +can return `authority_changed: null`: the outcome is uncertain, so retry rather +than assume the source is still selected. A completed plan superseded by another +migration cannot reactivate its former target. + +To return to File, create a **new** plan with `--provider file`, then preview and +execute it. This carries the current SQLite history back to File and retains +new acknowledged writes. An old File history is reusable only if it is an exact +prefix. Divergent/extra target history rejects migration; do not delete it or +copy old bytes over live state to force acceptance. + +Provider rollback is distinct from binary downgrade. Earlier binaries that do +not recognize explicit File selectors cannot operate this runtime. Keep a +migration-capable binary; do not remove the selector to make an old binary start. +The generic historical format check alone does not prove selector compatibility. +This operation qualifies local storage continuity, not D2 capacity/soak, all +Host lifecycles, all Goal consumers, or permission to enable a provider by default. + +### 已晋升 Goal 的本地 provider 切换 + +先停止写入并结算租约,再通过 `plan-migration` 保存并审核计划,用输出的摘要调用 +`migrate` 预览和执行。未结算的租约即使过期也会阻止迁移。TS 在同一个 canonical +写锁下完成备份、历史与回执核对、selector 发布及回读;Python 只传参数和结果。 + +中断后用同一计划重试;切换后产生的新写入会被保留。回退是另做一个 `--provider +file` 的新计划,把最新历史迁回 File,不能覆盖旧备份。此处不负责停 Host、结算 +Turn、旧 Markdown Goal 晋升或 PostgreSQL 服务部署,也不代表默认值资格通过。 +旧版本若不认识显式 File selector,会拒绝读取;保留支持迁移的运行时,不删除 +selector 绕过检查。 From 698d0813b4864a67f94ba6e055ff1bffc473b705 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:24:43 +0800 Subject: [PATCH 3/4] fix(authority): preserve uncertainty after migration transport loss Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/authority_archive.py | 6 ++++- tests/control_plane/test_authority_archive.py | 25 +++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/loopx/cli_commands/authority_archive.py b/loopx/cli_commands/authority_archive.py index 48d19f9efb..fede9785e0 100644 --- a/loopx/cli_commands/authority_archive.py +++ b/loopx/cli_commands/authority_archive.py @@ -100,7 +100,11 @@ def handle_authority_archive_command( "coordination.authority_archive.manage", request, timeout=300.0, retry_safe=False ) except (OSError, RuntimeError, ValueError) as error: - result = {"status": "failed", "reason": str(error), "authority_changed": False} + uncertain = args.authority_archive_action == "migrate" and args.execute + result = {"status": "failed", "reason": str(error), + "authority_changed": None if uncertain else False} + if uncertain: + result.update(reason_code="migration_outcome_unknown", requires_same_plan_retry=True) if (args.authority_archive_action == "upgrade" and args.require_current and any(row.get("status") == "planned" for row in result.get("results", []))): result.update(status="failed", reason="Authority format upgrade required before activating this runtime.") diff --git a/tests/control_plane/test_authority_archive.py b/tests/control_plane/test_authority_archive.py index f104f9a184..95f10d23e9 100644 --- a/tests/control_plane/test_authority_archive.py +++ b/tests/control_plane/test_authority_archive.py @@ -191,3 +191,28 @@ def cli(*args, expected=0): finally: subprocess.run([sys.executable, "-c", "from loopx.control_plane.effect_runtime import effect_runtime_result; effect_runtime_result('runtime.shutdown',{},retry_safe=False)"], cwd=REPO, capture_output=True, text=True, timeout=30, check=True) + + +@pytest.mark.parametrize("execute", [False, True]) +def test_migration_transport_loss_never_asserts_that_execution_did_not_publish(tmp_path, monkeypatch, execute): + from argparse import Namespace + from loopx.cli_commands import authority_archive + + def disconnected(*args, **kwargs): + assert kwargs["retry_safe"] is False + raise RuntimeError("Connection lost after request delivery") + + monkeypatch.setattr(authority_archive, "effect_runtime_result", disconnected) + registry = tmp_path / "registry.json" + registry.write_text("{}") + args = Namespace(command="authority-archive", authority_archive_action="migrate", goal_id="example", + plan=tmp_path / "plan.json", plan_sha256="a" * 64, execute=execute) + outputs = [] + status = authority_archive.handle_authority_archive_command( + args, registry_path=registry, runtime_root_arg=str(tmp_path / "runtime"), + print_payload=lambda payload, *_: outputs.append(payload), output_format=lambda _: "json") + assert status == 1 + assert outputs[0]["authority_changed"] is (None if execute else False) + if execute: + assert outputs[0]["requires_same_plan_retry"] is True + assert outputs[0]["reason_code"] == "migration_outcome_unknown" From caf36810524a65250fb341dbc09426b422de829c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:27:06 +0800 Subject: [PATCH 4/4] test(authority): prove Todo metadata fidelity across provider migration Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/reference/file-authority-state-log.md | 8 ++++ loopx/cli_commands/authority_archive.py | 2 +- .../local_authority_migration.test.ts | 39 +++++++++++++++++++ 3 files changed, 48 insertions(+), 1 deletion(-) diff --git a/docs/reference/file-authority-state-log.md b/docs/reference/file-authority-state-log.md index 0689f2fe84..15d9b8bf86 100644 --- a/docs/reference/file-authority-state-log.md +++ b/docs/reference/file-authority-state-log.md @@ -260,6 +260,14 @@ cursor, projection digest, writer-fence digest and destination provider. A changed source requires a new reviewed plan. Plans are created exclusively; choose a new path instead of overwriting an already-reviewed artifact. +Migration copies complete committed projections, events and original receipts; +it does not rebuild Todos from display columns or reapply a metadata allowlist. +Historical metadata values, absent keys, explicit nulls, false, zero and empty +arrays remain distinct. Auditing excludes only the transaction's physical +provider revision, which legitimately changes with the backend. It cannot prove +that an earlier legacy-to-canonical capture included every external field, nor +does it copy attachment files or independently-owned Host/Turn stores. + Execution shares the maintenance guard used by canonical command writers. It saves a verified logical backup below the runtime's `authority-transition/local-provider//`, binds the target identity, diff --git a/loopx/cli_commands/authority_archive.py b/loopx/cli_commands/authority_archive.py index fede9785e0..fbcfa023a6 100644 --- a/loopx/cli_commands/authority_archive.py +++ b/loopx/cli_commands/authority_archive.py @@ -8,7 +8,7 @@ from ..control_plane.effect_runtime import effect_runtime_result from ..paths import DEFAULT_RUNTIME_ROOT, global_registry_path, resolve_runtime_root -from ..history import load_registry +from ..control_plane.projects.registry_codec import load_registry def register_authority_archive_command( diff --git a/tests/control_plane_ts/local_authority_migration.test.ts b/tests/control_plane_ts/local_authority_migration.test.ts index 6477c0d94e..89e90856da 100644 --- a/tests/control_plane_ts/local_authority_migration.test.ts +++ b/tests/control_plane_ts/local_authority_migration.test.ts @@ -254,3 +254,42 @@ test("post-publication IO failure reports uncertainty, same-plan retry reads the assert.equal((await selected(root, goal)).provider, "sqlite"); assert.equal((await manage(request(root, p))).status, "already_applied"); }); + +test("Todo metadata survives every historical row, including null, false, empty arrays and removed keys", async t => { + const root = await fixture(t); + const source = (await selected(root, goal)).store; + const metadata: JsonObject = {priority: "P1", task_class: "advancement_task", task_domain: "code", + required_capabilities: ["shell", "filesystem_write"], required_write_scopes: ["src/**", "tests/**"], + target_capabilities: ["coordination_authority"], excluded_agents: [], claimed_by: null, + global_gate: false, completion_validation_revision: 0, completion_validation_revision_history: [], + note: "保留多行说明\n- 原始 metadata", evidence: "validation://retained-metadata"}; + for (let i = 0; i < 2; i++) { + const head = await source.loadAuthority(); + if (head.status !== "loaded") throw Error("head"); + const record = {...(head.head.todos as JsonObject[])[0], ...metadata}; + if (i === 1) { record.note = null; delete record.evidence; } + const projection = authorityProjectionFixture(goal, [record], [], "native", { + retained_annotation: {optional: null, count: 0, enabled: false, empty: [], nested: {"中文": "值"}}, + }); + const committed = await source.commitAuthority({operation_id: `metadata-${i}`, expected_provider_revision: head.provider_revision, + events: [{kind: "metadata_update", note_present: true}], receipts: [{metadata_preserved: true}], next_projection: projection}); + assert.equal(committed.status, "applied"); + } + const before = await source.scanCommitted(null, 10); + if (before.status !== "page") throw Error("history"); + for (const target of ["sqlite", "file"] as const) { + assert.equal((await manage(request(root, await plan(root, target)))).status, "migrated"); + const after = await (await selected(root, goal)).store.scanCommitted(null, 10); + if (after.status !== "page") throw Error("history"); + assert.deepEqual(after.transactions.map(({provider_revision, ...row}) => row), + before.transactions.map(({provider_revision, ...row}) => row)); + const old = (after.transactions[2].projection.todos as JsonObject[])[0]; + const current = (after.transactions[3].projection.todos as JsonObject[])[0]; + for (const [key, value] of Object.entries(metadata)) assert.deepEqual(old[key], value, key); + assert.equal(current.note, null); + assert.equal(Object.hasOwn(current, "evidence"), false); + assert.deepEqual(current.excluded_agents, []); + assert.equal(current.global_gate, false); + assert.equal(current.completion_validation_revision, 0); + } +});