diff --git a/apps/presentation/dashboard/src/data/goal-storage.ts b/apps/presentation/dashboard/src/data/goal-storage.ts index 90dc342664..14b8c28b3c 100644 --- a/apps/presentation/dashboard/src/data/goal-storage.ts +++ b/apps/presentation/dashboard/src/data/goal-storage.ts @@ -4,10 +4,16 @@ import {ChatApiError, requestJson} from "./chat"; // provider graph into the browser; this schema validates its HTTP projection. const provider = z.enum(["file", "sqlite"]); export type MigrationProvider = z.infer; -export const storageCarrierSchema = z.object({ +const migrationCarrierSchema = z.object({ goal_id: z.string().min(1), preview_id: z.string().regex(/^[a-f0-9]{32}$/), plan_sha256: z.string().regex(/^[a-f0-9]{64}$/), }); +// Reuse the cold-import owner's operation identity; old migration carriers +// remain readable without a new protocol discriminator or guessed source. +export const storageCarrierSchema = z.union([migrationCarrierSchema, z.object({ + goal_id: z.string().min(1), operation_id: z.string().regex(/^[a-f0-9]{32}$/), + plan_sha256: z.string().regex(/^[a-f0-9]{64}$/), +})]); export type StorageCarrier = z.infer; export const storageSourceSchema = z.object({ goal_id: z.string(), canonical: z.boolean(), provider: provider.nullable(), @@ -31,6 +37,14 @@ export const storageResultSchema = z.object({ }).optional(), recovery: z.object({phase: z.enum(["prepared", "completed"]), target_store_identity: z.string(), archive_sha256: z.string()}).nullable().optional(), reason_code: z.string().optional(), + operation_id: z.string().regex(/^[a-f0-9]{32}$/).optional(), + target_handoff_mode: z.enum(["soft_claim", "hard_lease"]).optional(), + source_inventory: z.object({todo_count: z.number().int().nonnegative(), + archived_todo_count: z.number().int().nonnegative(), lease_count: z.number().int().nonnegative(), + source_handoff_mode: z.string()}).optional(), + legacy_writer_fenced: z.boolean().nullable().optional(), + coordination_source_backup_verified: z.boolean().optional(), + complete_goal_backup_verified: z.literal(false).optional(), }); export type StorageResult = z.infer; @@ -44,9 +58,14 @@ async function read(url: string, init?: RequestInit): Promise { export function fetchGoalStorage(goalId: string) { return read(`/api/chat/goal-storage?${new URLSearchParams({goal_id: goalId})}`); } -export function previewGoalStorage(goalId: string, target: MigrationProvider) { +export function previewGoalStorage(goalId: string, target: MigrationProvider, mode?: "soft_claim" | "hard_lease") { + if (mode) return read("/api/chat/goal-storage/import/preview", {method: "POST", + body: JSON.stringify({goal_id: goalId, provider: target, handoff_mode: mode})}); return read("/api/chat/goal-storage/preview", {method: "POST", body: JSON.stringify({goal_id: goalId, provider: target})}); } export function recoverGoalStorage(carrier: StorageCarrier, apply = false) { + const saved = storageCarrierSchema.parse(carrier); + if ("operation_id" in saved) return read(`/api/chat/goal-storage/import/${apply ? "apply" : "recover"}`, { + method: "POST", body: JSON.stringify({...saved, ...(apply ? {writers_stopped: true} : {})})}); return read(`/api/chat/goal-storage/${apply ? "apply" : "recover"}`, {method: "POST", body: JSON.stringify(storageCarrierSchema.parse(carrier))}); } diff --git a/apps/presentation/dashboard/src/features/personal-workspace/goal-storage-settings.tsx b/apps/presentation/dashboard/src/features/personal-workspace/goal-storage-settings.tsx index a6383c6111..cb9ff37d5a 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/goal-storage-settings.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/goal-storage-settings.tsx @@ -13,6 +13,7 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan const [carrier, setCarrier] = useState(null); const [result, setResult] = useState(null); const [target, setTarget] = useState("sqlite"); + const [mode, setMode] = useState<"" | "soft_claim" | "hard_lease">(""); const [confirmed, setConfirmed] = useState(false); const [busy, setBusy] = useState(false); const [error, setError] = useState(null); @@ -22,7 +23,7 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan const generation = useRef(0); useEffect(() => { const token = ++generation.current; - setCurrent(null); setCold(undefined); setResult(null); setCarrier(null); setConfirmed(false); setInvalidSaved(false); setError(null); setBusy(true); + setCurrent(null); setCold(undefined); setResult(null); setCarrier(null); setConfirmed(false); setMode(""); setInvalidSaved(false); setError(null); setBusy(true); let saved: StorageCarrier | null = null; try { const raw = localStorage.getItem(key); @@ -54,11 +55,13 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan }, [goalId, key, reload, t]); async function submit(apply: boolean) { - if (inFlight.current || busy || invalidSaved || (apply && (!carrier || !confirmed))) return; + if (inFlight.current || busy || invalidSaved || (apply && (!carrier || !confirmed)) || + (!apply && !current?.canonical && !mode)) return; inFlight.current = true; setBusy(true); setError(null); const token = generation.current; try { - const next = apply ? await recoverGoalStorage(carrier!, true) : await previewGoalStorage(goalId, target); + const next = apply ? await recoverGoalStorage(carrier!, true) : await previewGoalStorage(goalId, target, + current?.canonical ? undefined : mode || undefined); if (token !== generation.current) return; setResult(next); if (!apply && next.ok) { @@ -87,7 +90,9 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan catch { setError(t("storage.savedInvalid")); return; } setCarrier(null); setResult(null); setConfirmed(false); setInvalidSaved(false); setError(null); } - const completed = result?.recovery?.phase === "completed"; + const coldCarrier = carrier && "operation_id" in carrier; + const completed = result?.recovery?.phase === "completed" || (coldCarrier && + (result?.status === "applied" || result?.status === "recovered" || result?.status === "replayed")); return

{t("storage.title")}

@@ -103,28 +108,33 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan

{t("storage.coldBoundary")}

}
: null} - {current?.canonical || carrier ? <> - {carrier ?

{t("storage.reviewed", {source: result?.reviewed_source?.provider ?? "?", target: result?.target_provider ?? "?", cursor: result?.reviewed_source?.cursor ?? "?"})}

+ {current || carrier ? <> + {carrier ?

{coldCarrier ? t("storage.coldReviewed", {target: result?.target_provider ?? "?", mode: result?.target_handoff_mode ? t(`ownership.${result.target_handoff_mode}`) : "?"}) : t("storage.reviewed", {source: result?.reviewed_source?.provider ?? "?", target: result?.target_provider ?? "?", cursor: result?.reviewed_source?.cursor ?? "?"})}

: } - {carrier && !completed ? : null} + {!carrier && !current?.canonical ? : null} + {result?.source_inventory ?

{t("storage.coldInventory", {todos: result.source_inventory.todo_count, archived: result.source_inventory.archived_todo_count, leases: result.source_inventory.lease_count})}

: null} + {carrier && !completed ? : null}
- {carrier ? <> + {carrier ? <> - : } + : }
: null}
{invalidSaved ? : null}
- {completed ?

{t("storage.completed")}

: null} - {result?.recovery?.phase === "prepared" ?

{t("storage.prepared")}

: null} + {completed ?

{t(coldCarrier ? "storage.coldCompleted" : "storage.completed")}

: null} + {result?.recovery?.phase === "prepared" || (coldCarrier && result?.status === "prepared") ?

{t("storage.prepared")}

: null} {error ?

{error}

: null} {result?.reason_code ?

{result.reason_code}

: null} {current?.canonical || carrier ?
{t("storage.details")} {current?.canonical ?

{current.store_identity} · {current.provider_revision} · {current.cursor}

: null} - {carrier ?

{carrier.preview_id} · {carrier.plan_sha256}

: null} + {carrier ?

{"operation_id" in carrier ? carrier.operation_id : carrier.preview_id} · {carrier.plan_sha256}

: null}
: null} {busy ?

{t("common.loading")}

: null} diff --git a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx index e730bacc45..fe271e78e7 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx @@ -20,7 +20,14 @@ const en = { "storage.coldCounts": "{active} active tasks · {archived} archived tasks · {leases} unsettled leases", "storage.coldCapture": "Original capture files observed; history retained.", "storage.coldOutbox": "Outbox files observed; their processing is unverified.", - "storage.coldBoundary": "Old Markdown source inspected. Import is not available here yet: writer/Host stop, lease settlement, outbox disposition and a backup-bound import still need verification. Nothing was captured, migrated or granted execution authority.", + "storage.coldBoundary": "Import this previous Markdown Goal from a verified private local backup. This does not stop Hosts, settle leases or dispose of capture/outbox work for you.", + "storage.choosePolicy": "Choose the execution policy", + "storage.coldReviewed": "Reviewed Markdown source → {target}; execution policy: {mode}. Apply rechecks the original source and backup.", + "storage.coldInventory": "Reviewed inventory: {todos} tasks, including {archived} archived · {leases} retained lease records", + "storage.coldConfirm": "I stopped all writers and Hosts, settled leases and disposed of capture/outbox work. Import this reviewed source.", + "storage.coldApply": "Import reviewed Markdown source", + "storage.coldPreview": "Back up and preview import", + "storage.coldCompleted": "The original import receipt is verified. Current storage is read independently above. Retained leases and receipts grant no new execution authority.", "storage.reviewed": "Reviewed source: {source}, cursor {cursor} → {target}. A changed source rejects apply.", "storage.confirm": "I stopped writers and settled leases. I confirm this reviewed storage change.", "storage.apply": "Back up and switch storage", @@ -43,7 +50,7 @@ const en = { "ownership.unpromoted": "Not on canonical storage", "storage.oldSource": "Old Markdown source", "storage.readUnavailable": "Current storage could not be read. Try reading it again.", - "ownership.promoteFirst": "This Goal has no canonical store. Ownership changes require a reviewed storage import; this policy form does not perform it.", + "ownership.promoteFirst": "Use Data storage below to review and import this Markdown Goal first. This policy form does not migrate storage.", "ownership.target": "New policy", "ownership.softHelp": "Coordinate assignment without requiring an exclusive execution lease. Active leases must be settled first.", "ownership.hardHelp": "Ownership changes and completion require the task’s original execution lease.", @@ -1368,7 +1375,14 @@ const zhCN: Record = { "storage.coldCounts": "{active} 项当前任务 · {archived} 项归档任务 · {leases} 项未结算 lease", "storage.coldCapture": "已发现原 capture 文件,历史原样保留。", "storage.coldOutbox": "已发现 outbox 文件,处理结果尚未验证。", - "storage.coldBoundary": "已盘点旧 Markdown 源。这里尚不能导入:仍需验证 writer/Host 停止、lease 结算、outbox 处置与备份绑定的导入。此次未生成 capture、迁移数据或授予执行权限。", + "storage.coldBoundary": "从已验证的本机私有备份导入这个旧 Markdown Goal。此操作不会替你停止 Host、结算 lease 或处置 capture/outbox。", + "storage.choosePolicy": "请选择执行策略", + "storage.coldReviewed": "已审核 Markdown 来源 → {target};执行策略:{mode}。导入时会复核原来源与备份。", + "storage.coldInventory": "已审核盘点:{todos} 项任务(含 {archived} 项归档)· {leases} 条保留租约记录", + "storage.coldConfirm": "我已停止全部写入方和 Host、结算 lease 并处置 capture/outbox。确认导入这份审核来源。", + "storage.coldApply": "导入已审核的 Markdown 来源", + "storage.coldPreview": "备份并预览导入", + "storage.coldCompleted": "原导入回执已核验。上方当前存储独立读回;保留的租约和回执不授予新的执行权限。", "storage.reviewed": "已审核来源:{source},游标 {cursor} → {target}。来源变化时拒绝应用。", "storage.confirm": "我已停止写入方并结算 lease,确认此预览的存储变更。", "storage.apply": "备份并切换存储", @@ -1391,7 +1405,7 @@ const zhCN: Record = { "ownership.unpromoted": "尚未使用统一状态存储", "storage.oldSource": "旧 Markdown 来源", "storage.readUnavailable": "未能读取当前存储,请重试读回。", - "ownership.promoteFirst": "此 Goal 尚无 canonical 存储。变更所有权需要先完成经过核对的存储导入;此策略表单不执行导入。", + "ownership.promoteFirst": "请先在下方“数据存储”中审核并导入此 Markdown Goal。此策略表单不迁移存储。", "ownership.target": "新策略", "ownership.softHelp": "协调任务归属,不要求独占执行租约。仍在运行的租约必须先结算。", "ownership.hardHelp": "变更所有权和完成任务需要持有该任务原有的执行租约。", diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 344c50c131..2c8bfec2ae 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -930,6 +930,34 @@ provider. This is a bounded App companion, not full old-source/Host upgrade or D2/release-default qualification. Continue those original acceptance frontiers and retire each last caller separately. +Cold-source import now has a bounded CLI/App coordination stage: complete source +records, an immutable source/target carrier bound to actual backup member bytes, +explicit operator shutdown attestation, revalidation before the durable writer +fence, and original-receipt recovery through the existing File/SQLite owners. +It refuses unresolved capture/outbox and unsettled leases, including expired +active and orphan records, without manufacturing shadow qualification. A killed +fenced process can resume without rereading Markdown; original-receipt replay +preserves later canonical writes. Coordination-source backup verification +**does not qualify complete Goal recovery**. Packaged Goal storage settings +reuse that transaction for private backup, inventory, explicit policy/stop +confirmation and original-operation readback. Reload is read-only, including a +fenced but uncommitted operation; applying the original carrier requires fresh +confirmation. File/SQLite HTTP qualification preserves later writes and refuses +source/backup drift. The operator-led POSIX stop path now exercises actual owned +Host processes and native source leases on File/SQLite: process exit and lease +release remain separate, expired active leases refuse import, and the old grant +cannot launch a Host after cutover. This proves the existing supervisor/lease +boundary with synthetic work, not automatic Host discovery or live model use. +Continue pending outbox disposition, App loading with the old normal writer +absent and independent full-backup recovery in R5. The +installed cold-import CLI uses the existing selected dispatcher and +Goal path resolver; real File/SQLite import and original-receipt recovery pass +with the four old normal producer modules physically absent in a disposable +package. This qualifies that command's loading boundary, not every other CLI +caller or removal of those modules. Keep T4 retirement on actual callers: the +retained Python prose-write guard still serves live callers and its obligation +must survive adapter removal. This partial stage does not retire the supported +old writer or qualify a released default. Cold-import preservation checkpoint (R5/D1, T4/C1): full-state backups now witness each saved regular member's bytes in their existing manifests, including raw Markdown history, lease/receipt files, SQLite snapshots and stored diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md index b5d6178584..b9cc3919e2 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md @@ -612,6 +612,15 @@ L3 检查点:独立领取/接管、原子 claim 准入与维护共用 typed le - **退出:** 按 shared-authority 7.2 分别决定有界改动、可回退自愿 cohort、发布默认值,各自在适用范围具备真实 CLI/backend、独立基线、负例和恢复证据。正式 D2 保留适用容量及至少十日证据,cohort 不必等该证书。D3 保留明确切换权限。本计划没有启动 soak 或晋升 provider。 - **回滚:** 按已审阅的 fenced export/import 和 schema-aware downgrade,不能靠替换二进制恢复旧写权威。 +冷旧 Goal 导入复用既有 TS source/promotion/receipt owner,CLI/App 使用绑定实际 +备份字节的预览、明确确认、停写 fence 与原操作恢复。POSIX 人工停止旅程现以真实 +owned Host 进程、未晋升源的原生租约和 File/SQLite 验证:进程退出不释放租约, +过期 active 租约仍拒绝导入;原生释放后保留历史身份,切换后旧 grant 在 Host 启动前 +拒绝。这是合成工作对既有 supervisor/lease 边界的验证,不是自动发现/停止 Host +或 live 模型验收。pending outbox 逐项处置、旧正常 writer 物理缺席的 App 加载、 +完整 Goal 历史与备份恢复继续开放;协调源备份验证不结算完整恢复。详见 +[冷源导入与支持边界](../../reference/local-authority-provider-selection.md)。 + 冷旧源保留检查点:本项目全量备份现盘点注册的自定义状态与来源 registry 路由。 真实 CLI 备份和独立静态解包保留完整 Markdown 字节、未引用的归档 Todo 与 runtime 原始历史,包含 SQLite snapshot。这只修复路由遗漏,不代表审核式冷导入或完整状态 diff --git a/docs/reference/local-authority-provider-selection.md b/docs/reference/local-authority-provider-selection.md index 5a1777fe26..01264b4e8b 100644 --- a/docs/reference/local-authority-provider-selection.md +++ b/docs/reference/local-authority-provider-selection.md @@ -328,6 +328,91 @@ archive restore into a new identity. Do not substitute either journey's acceptance for this one. Once this route qualifies, an old Goal's next normal write requires reviewed import; installing the binary alone does not migrate it. +The CLI stage uses the typed `coordination.cold_source.import` transaction +(`loopx_cold_source_import_request_v0`). First stop affected writers through +their owning Host and verify that their processes have exited. Inspect each +retained task lease and release it through `task-lease release`, using its +original owner/key and current `--expected-version`. An expired lease is still +unsettled; stopping a process does not release its lease. Dispose of pending +capture/outbox through its owning workflow before preparing this import. +Then execute `backup-state` with the source state, registry, coordination +evidence and runtime root. Lease release after a saved preview changes its +source: make a fresh backup and reviewed preview. `--writers-stopped` records +an operator attestation, not an automatic process stop. For example, after +shutdown and settlement, from the registered project: + +```bash +loopx --format json backup-state --project . --execute --no-automations --no-skills +loopx --format json coordination-shadow prepare-import --goal-id GOAL --operation-id IMPORT \ + --backup-manifest SAVED-MANIFEST.json --provider sqlite --target-handoff-mode hard_lease +# Inspect the immutable plan at plan_path. Confirm using its exact plan_sha256. +loopx --format json coordination-shadow apply-import --goal-id GOAL --operation-id IMPORT \ + --plan-sha256 SAVED-SHA256 --writers-stopped --execute +# After interruption, use the original operation and digest; do not prepare a new source. +loopx --format json coordination-shadow recover-import --goal-id GOAL --operation-id IMPORT \ + --plan-sha256 SAVED-SHA256 --execute +``` + +Prepare reads complete active and archived Todo records and preserves metadata, +original evidence and settled lease history. The trusted Python tar adapter +reads actual archive members, including older backups without a member list; +the TypeScript owner verifies coverage of the coordination source bytes and +pins both saved artifacts. A missing source path in the backup refuses import. +Prepare may select an empty target and persist its identity and immutable plan, +but does not fence the writer, import state or grant execution authority. + +Apply rechecks the original source and backup under the owning locks before +engaging `loopx_cold_source_import_writer_fence_v0`. Missing identity, changed +source, unsettled leases (including expired active or orphan records) and +unresolved capture/outbox fail closed. Recovery requires that original fence, +accepts no replacement source snapshot, and reads the original receipt without +overwriting later canonical writes. Source bytes remain retained history; the +import writes a separate new receipt. Recovery is not a provider rollback or a +whole-Goal restore, and replacing the binary cannot clear the writer fence. + +`coordination_source_backup_verified=true` qualifies the prepared coordination +source witness only. `complete_goal_backup_verified=false` preserves the full +backup and original-history recovery acceptance. The packaged Goal storage +settings use the same transaction for cold import: + +1. Select the Goal and open **Goal settings → Task ownership → Data storage**. + Choose File or SQLite and an explicit supported execution policy. +2. Stop affected writers/Hosts, verify their exit and settle/dispose of refused + work through its owning workflow. The preview refuses unsettled leases, + including expired active records; a checkbox cannot bypass this refusal. +3. **Back up and preview import** creates a private local archive using the + existing backup owner, then displays the complete active/archive inventory. + Preview does not import, stop a Host, settle leases or grant execution. + Confirm **Import reviewed Markdown source**. Apply + rechecks the bound original source and backup; a changed source needs a new + reviewed preview. +4. After a lost response or reload, **Read original preview and current storage** + observes the saved operation and current store independently. Reload never + completes an unfinished import. Confirm again to retry the original operation; + do not discard its carrier while the commit is ambiguous. + +Only Goal/operation/digest identifiers survive in browser storage. Full plans, +source bytes and backup paths stay local to the server. Completed receipt +readback neither overwrites later writes nor reselects a provider. +The operator-led POSIX stop path is exercised with actual owned Host processes +and native unpromoted-source leases on File/SQLite: import refuses while the +Host runs and after it exits with an active lease; native release permits +cutover. The imported released lease retains its identity and history, and a +restart using the old grant is rejected before the actual Host launches. +This does not qualify automatic Host discovery/stop, pending outbox disposition, +live model sessions or full-history restore. The +cold-import CLI reuses the selected command dispatcher and +the existing Goal path resolver; its File/SQLite import and original-receipt +recovery run with `todos.py`, `bootstrap.py`, `runtime_shadow_writer_adapter.py` +and `local_authority_shadow_outbox.py` physically absent in a disposable package. +This proves that command's independence, not that other commands or supported +writers can lose those files. Retained prose-write guards have live callers and +must move to their owning boundary before their adapter is removed. +The packaged App uses the real File/SQLite backend; App loading with those old +modules absent remains a separate acceptance. This stage does not qualify +supported old-writer retirement, +release default or historical support cutoff. + Deliver complete, reversible PR packages in this order: 1. **Direct source import and support boundary.** Implement and qualify the diff --git a/docs/reference/local-authority-provider-selection.zh-CN.md b/docs/reference/local-authority-provider-selection.zh-CN.md index 21a6669d06..d9f5c219a6 100644 --- a/docs/reference/local-authority-provider-selection.zh-CN.md +++ b/docs/reference/local-authority-provider-selection.zh-CN.md @@ -83,7 +83,18 @@ outbox 后,预览来源/目标绑定的不可变计划、明确确认、fence File/SQLite 都须在旧正常 writer 物理缺席时验证失败、来源变化与恢复。 这是未晋升 Goal 的导入,与既有 canonical provider 切换和新身份 archive restore -分别验收;实现尚未完成。历史支持矩阵验收前保留兼容安装包用于恢复;可在兼容 +分别验收。CLI 与打包 App 已有共享事务的有界实现:在 Goal 设置的任务所有权页, +通过数据存储选择 File/SQLite 和明确执行策略,备份并预览,再明确确认导入。 +浏览器只保存 Goal/原操作/摘要标识;来源字节、完整计划与备份路径仅存于本机。 +重载只读回原操作和当前存储,不自动完成未提交的导入;失响应时保留原预览, +重新确认后重试同一操作。来源或备份变化明确拒绝;原回执核验不覆盖后续写入。 +备份预览前须先通过原 Host 停写并核对进程退出,再以原 owner/key 和当前版本 +原生释放租约;到期不等于结算,预览后释放会改变来源,须重做备份和审核预览。 +POSIX 人工停止旅程已用真实 owned Host 进程、未晋升源的原生租约和 File/SQLite +验证:Host 退出后仍 active 的租约拒绝导入;释放后保留历史身份,旧 grant 重启 +在实际 Host 启动前拒绝。这不证明自动发现/停止 Host、live 模型会话或 pending +outbox 的逐项处置。完整历史恢复及旧正常 writer 物理缺席时的 App 加载仍待验收。 +历史支持矩阵验收前保留兼容安装包用于恢复;可在兼容 1.x 先发布新增 importer 和弃用提示,但不把旧 writer 使用设为导入前置。 交付按完整包推进:CLI/App 直接导入与支持边界 → 同调用族旧 writer/私有 dispatch diff --git a/examples/personal-workspace-browser/goal-storage.mjs b/examples/personal-workspace-browser/goal-storage.mjs index 1df2039d38..7cd914fd7a 100644 --- a/examples/personal-workspace-browser/goal-storage.mjs +++ b/examples/personal-workspace-browser/goal-storage.mjs @@ -12,7 +12,7 @@ import {launchBrowser, loadPlaywright, waitForHttp} from "../dashboard-browser-s import {openWorkspacePage} from "./scenario-context.mjs"; import {outputDir, packaged, port, repoRoot, startServer} from "./fixture.mjs"; -async function startAuthority() { +async function startAuthority({unsettled = true} = {}) { const root = await mkdtemp(resolve(tmpdir(), "loopx-old-storage-browser-")); const source = "---\ngoal_id: multi-agent-projection\nhandoff_mode: soft_claim\n---\n" + "## Agent Todo\n- [ ] Private current requirement\n" + @@ -26,7 +26,7 @@ async function startAuthority() { })); const leases = resolve(root, "runtime/goals/multi-agent-projection/task-leases"); await mkdir(leases, {recursive: true}); - await writeFile(resolve(leases, "removed.json"), JSON.stringify({schema_version: "task_lease_v0", + if (unsettled) await writeFile(resolve(leases, "removed.json"), JSON.stringify({schema_version: "task_lease_v0", goal_id: "multi-agent-projection", todo_id: "removed", owner: "agent-a", status: "active", idempotency_key: "original-private-key", version: 4, lease_epoch: 2, expires_at: "2000-01-01T00:00:00Z", write_scopes: []})); @@ -66,6 +66,69 @@ server.serve_forever() } } +async function importSettled(browser, url, provider) { + const authority = await startAuthority({unsettled: false}); + let context; + let applyCount = 0; + try { + context = await openWorkspacePage(browser, url, {beforeGoto: async (_api, page) => { + await page.route(/\/api\/chat\/goal-(storage|ownership)(?:[/?]|$)/, async route => { + const parsed = new URL(route.request().url()); + const response = await route.fetch({url: authority.url + parsed.pathname + parsed.search}); + if (parsed.pathname.endsWith("/import/apply") && ++applyCount === 1) { + assert.equal(response.status(), 200); // Commit succeeded; only its response is lost. + await route.abort("failed"); + } else await route.fulfill({response}); + }); + }}); + const {page} = context; + async function open() { + const navigation = page.getByRole("button", {name: "打开 Goal 导航", exact: true}); + if (await navigation.isVisible()) await navigation.click(); + await page.locator(".personal-goal-link", {hasText: "Multi Agent Projection"}).click(); + await page.getByRole("button", {name: "Goal 设置", exact: true}).click(); + await page.getByRole("button", {name: "任务所有权", exact: true}).click(); + return page.getByRole("region", {name: "Goal 数据存储"}); + } + let panel = await open(); + await panel.getByText("1 项当前任务 · 1 项归档任务 · 0 项未结算 lease", {exact: true}).waitFor(); + await panel.getByRole("combobox", {name: "目标存储", exact: true}).selectOption(provider); + await panel.getByRole("combobox", {name: "新策略", exact: true}).selectOption("hard_lease"); + await panel.getByRole("button", {name: "备份并预览导入", exact: true}).click(); + const apply = panel.getByRole("button", {name: "导入已审核的 Markdown 来源", exact: true}); + try { await panel.getByRole("checkbox").waitFor(); } + catch (error) { throw new Error(`${error.message}; panel=${await panel.innerText()}`); } + assert.ok(await apply.isDisabled()); + const saved = await page.evaluate(() => localStorage.getItem("loopx-storage-preview:multi-agent-projection")); + assert.ok(saved && !saved.includes(authority.root) && !saved.includes("Private")); + await page.reload({waitUntil: "networkidle"}); + panel = await open(); + await panel.getByRole("checkbox").waitFor(); + await panel.getByRole("status").filter({hasText: "已准备"}).waitFor(); + assert.ok(await panel.getByRole("button", {name: "导入已审核的 Markdown 来源", exact: true}).isDisabled()); + assert.equal(applyCount, 0, "reload observes the original prepared operation without applying"); + await page.setViewportSize({width: 390, height: 844}); + assert.equal(await page.evaluate(() => document.documentElement.scrollWidth > innerWidth + 1), false); + await panel.getByRole("checkbox").focus(); + await page.keyboard.press("Space"); + assert.equal(await panel.getByRole("checkbox").isChecked(), true); + await panel.getByRole("button", {name: "导入已审核的 Markdown 来源", exact: true}).click(); + await panel.getByRole("alert").waitFor(); + assert.equal(applyCount, 1); + assert.equal(await page.evaluate(() => localStorage.getItem("loopx-storage-preview:multi-agent-projection")), saved); + await page.reload({waitUntil: "networkidle"}); + panel = await open(); + await panel.getByRole("status").filter({hasText: "原导入回执已核验"}).waitFor(); + assert.equal(applyCount, 1, "completed readback does not apply again"); + await panel.getByText(provider, {exact: true}).waitFor(); + assert.equal(await panel.getByRole("checkbox").count(), 0); + assert.equal(await readFile(resolve(authority.root, "state.md"), "utf8"), authority.source); + await panel.getByRole("status").filter({hasText: "原导入回执已核验"}).scrollIntoViewIfNeeded(); + await page.screenshot({path: resolve(outputDir, `goal-storage-import-${provider}-mobile.png`), animations: "disabled"}); + assert.equal(context.errors.filter(e => !e.includes("ERR_FAILED")).length, 0, context.errors.join("; ")); + } finally { await context?.close(); await authority.close(); } +} + const goalStorageScenario = { id: "goal-storage", async run({browser, url}) { @@ -90,8 +153,12 @@ const goalStorageScenario = { let panel = await open(); await panel.getByText("1 项当前任务 · 1 项归档任务 · 1 项未结算 lease", {exact: true}).waitFor(); assert.equal(await panel.getByRole("checkbox").count(), 0); - assert.equal(await panel.getByRole("combobox").count(), 0); - assert.ok((await panel.innerText()).includes("尚不能导入")); + assert.equal(await panel.getByRole("combobox").count(), 2); + assert.ok((await panel.innerText()).includes("此操作不会替你停止 Host")); + await panel.getByRole("combobox", {name: "新策略", exact: true}).selectOption("soft_claim"); + await panel.getByRole("button", {name: "备份并预览导入", exact: true}).click(); + await panel.getByText("cold_import_lease_requires_settlement", {exact: true}).waitFor(); + assert.equal(await panel.getByRole("checkbox").count(), 0); assert.ok(!(await panel.innerText()).includes("Private")); await page.screenshot({path: resolve(outputDir, "goal-storage-cold-desktop.png"), animations: "disabled"}); const outbox = resolve(authority.root, "runtime/authority-shadow/outbox/multi-agent-projection/todos"); @@ -116,13 +183,14 @@ const goalStorageScenario = { await page.reload({waitUntil: "networkidle"}); panel = await open("en"); await panel.getByText("1 active tasks · 1 archived tasks · 1 unsettled leases", {exact: true}).waitFor(); - assert.ok((await panel.innerText()).includes("Import is not available here yet")); - assert.equal(writes, 0); + assert.ok((await panel.innerText()).includes("This does not stop Hosts")); + assert.equal(writes, 1, "only the explicit refused preview posts"); assert.equal(await readFile(resolve(authority.root, "state.md"), "utf8"), authority.source); assert.equal(await readFile(residue, "utf8"), "{unrecognized original bytes"); // Deliberately induced HTTP failures may be logged by the browser. - assert.equal(context.errors.filter(e => !e.includes("503")).length, 0, context.errors.join("; ")); - return {note: "Cold/archived tasks, expired orphan lease, original outbox, unavailable source and fresh recovery; no import or write; packaged Chinese/English desktop/mobile."}; + assert.equal(context.errors.filter(e => !e.includes("503") && !e.includes("409")).length, 0, context.errors.join("; ")); + for (const provider of ["file", "sqlite"]) await importSettled(browser, url, provider); + return {note: "Cold/archived tasks, expired orphan lease, original outbox, unavailable source and fresh recovery; refused expired lease; confirmed File/SQLite import, readonly reload and lost-response original receipt recovery; packaged Chinese/English desktop/mobile."}; } finally { await context?.close(); await authority.close(); } }, }; diff --git a/loopx/cli_commands/coordination_shadow.py b/loopx/cli_commands/coordination_shadow.py index 69fcecfd6c..38365e44bc 100644 --- a/loopx/cli_commands/coordination_shadow.py +++ b/loopx/cli_commands/coordination_shadow.py @@ -27,7 +27,7 @@ from ..control_plane.projects.registry_codec import load_project_registry from ..paths import resolve_runtime_root from ..registry import find_registry_goal -from ..state_refresh import resolve_goal_state +from ..control_plane.goals.state_resolution import resolve_goal_state PrintPayload = Callable[ @@ -69,10 +69,13 @@ def register_coordination_shadow_command( ), ("rollback", "Quarantine one exact pre-promotion file shadow lineage."), ("recover-promotion", "Read back or recover an already-fenced saved promotion without reading legacy Markdown."), + ("prepare-import", "Save a reviewed complete cold-source import and verified existing backup; does not stop writers or import."), + ("apply-import", "Import the original reviewed cold source after explicit writer/Host stop confirmation."), + ("recover-import", "Recover only the original fenced cold import, without reading Markdown."), ): action = actions.add_parser(name, help=help_text) action.add_argument("--goal-id", required=True) - if name != "recover-promotion": + if name not in {"recover-promotion", "apply-import", "recover-import"}: action.add_argument("--project", type=Path) action.add_argument("--state-file", type=Path) if name in {"bootstrap", "rollback", "promote", "recover-promotion"}: @@ -132,6 +135,19 @@ def register_coordination_shadow_command( required=True, help="Exact Todo identity to read from the parity-matched file head.", ) + if name == "prepare-import": + action.add_argument("--operation-id", required=True) + action.add_argument("--backup-manifest", type=Path, required=True) + action.add_argument("--provider", choices=("file", "sqlite"), required=True) + action.add_argument("--target-handoff-mode", choices=("soft_claim", "hard_lease"), required=True) + if name in {"apply-import", "recover-import"}: + action.add_argument("--operation-id", required=True) + action.add_argument("--plan-sha256", required=True) + action.add_argument("--execute", action="store_true", required=True, + help="Explicitly execute this original reviewed administrative operation.") + if name == "apply-import": + action.add_argument("--writers-stopped", action="store_true", required=True, + help="Attest that the original writers and attached Hosts are stopped; expiry and an idle lock are insufficient.") def _projection_version(projection: dict[str, object]) -> str: @@ -202,6 +218,16 @@ def _render(payload: dict[str, object]) -> str: f"- legacy_writer_fenced: `{promotion.get('legacy_writer_fenced')}`", ] ) + cold_import = payload.get("cold_import") + if isinstance(cold_import, dict): + for field in ( + "status", "reason_code", "operation_id", "plan_path", "plan_sha256", "target_provider", + "target_handoff_mode", + "legacy_writer_fenced", "execution_authority_granted", + "coordination_source_backup_verified", "complete_goal_backup_verified", + ): + if field in cold_import: + lines.append(f"- {field}: `{cold_import[field]}`") bounded = qualification if isinstance(qualification, dict) else read_candidate if isinstance(bounded, dict) and bounded.get("scope") == "bounded": lines.extend([ @@ -233,6 +259,38 @@ def handle_coordination_shadow_command( if goal is None: raise ValueError(f"goal {args.goal_id!r} is not present in the registry") runtime_root = resolve_runtime_root(registry, runtime_root_arg, registry_path=registry_path) + if args.coordination_shadow_command in {"prepare-import", "apply-import", "recover-import"}: + # Backup/source codecs bind physical Host paths. Resolve aliases at + # this IO boundary, retaining the same identity across all phases. + runtime_root = runtime_root.resolve() + from ..control_plane.coordination.local_authority_shadow_projection import source_effect_runtime_result + request = {"schema_version": "loopx_cold_source_import_request_v0", + "runtime_root": str(runtime_root.expanduser().absolute()), "goal_id": args.goal_id, + "operation_id": args.operation_id} + if args.coordination_shadow_command == "prepare-import": + from ..control_plane.coordination.cold_source_backup import read_cold_source_backup + _, _, state_path = resolve_goal_state(registry=registry, goal_id=args.goal_id, + project_override=args.project, state_file_override=args.state_file) + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, + runtime_root=runtime_root, state_path=state_path, registry_path=registry_path, + include_all_archived_todos=True) + request.update(action="prepare", projection=projection, source_snapshot=snapshot, + target_handoff_mode=args.target_handoff_mode, target_provider=args.provider, + source_backup=read_cold_source_backup(args.backup_manifest)) + else: + request.update(action="apply" if args.coordination_shadow_command == "apply-import" else "recover", + expected_plan_sha256=args.plan_sha256) + if args.coordination_shadow_command == "apply-import": + request["writers_stopped"] = args.writers_stopped + result = source_effect_runtime_result("coordination.cold_source.import", request, + retry_safe=False, timeout=300.0) + payload = {"ok": result.get("status") in {"prepared", "applied", "recovered", "replayed"}, + "schema_version": "loopx_coordination_shadow_admin_v0", + "action": args.coordination_shadow_command, "goal_id": args.goal_id, + "executed": result.get("executed") is True, "cold_import": result, + "decision_read_from_shadow": False} + print_payload(payload, output_format(args), _render) + return 0 if payload["ok"] else 1 reviewed_path = getattr(args, "reviewed_plan", None) reviewed_plan = None if reviewed_path is not None: diff --git a/loopx/cli_runtime.py b/loopx/cli_runtime.py index 96f28b2aba..9db8f2ae37 100644 --- a/loopx/cli_runtime.py +++ b/loopx/cli_runtime.py @@ -67,7 +67,7 @@ _STATUS_COMMANDS = frozenset({"check", "status", "diagnose", "review-packet"}) _SELECTED_COMMANDS = _STATUS_COMMANDS | { "todo", "quota", "task-lease", "change-window", "delegation", "turn", "doctor", "commands", - "authority-archive", "extension", "slash-commands", + "authority-archive", "coordination-shadow", "extension", "slash-commands", } @@ -259,6 +259,10 @@ def _build_selected_parser(command: str) -> LoopXArgumentParser: from .cli_commands.authority_archive import register_authority_archive_command register_authority_archive_command(subparsers, add_subcommand_format) + elif command == "coordination-shadow": + from .cli_commands.coordination_shadow import register_coordination_shadow_command + + register_coordination_shadow_command(subparsers, add_subcommand_format) elif command == "extension": from .cli_commands.extension import register_extension_commands @@ -300,6 +304,13 @@ def _dispatch_common_command( args, registry_path=registry_path, runtime_root_arg=args.runtime_root, output_format=output_format, print_payload=print_payload, ) + if args.command == "coordination-shadow": + from .cli_commands.coordination_shadow import handle_coordination_shadow_command + + return handle_coordination_shadow_command( + args, registry_path=registry_path, runtime_root_arg=args.runtime_root, + output_format=output_format, print_payload=print_payload, + ) if args.command == "extension": from .cli_commands.extension import handle_extension_command diff --git a/loopx/control_plane/coordination/cold_source_backup.py b/loopx/control_plane/coordination/cold_source_backup.py new file mode 100644 index 0000000000..efe45c7ad0 --- /dev/null +++ b/loopx/control_plane/coordination/cold_source_backup.py @@ -0,0 +1,85 @@ +"""Read actual backup archive members for the typed cold-import owner. + +Python owns tar/Host IO only. Membership and source coverage decisions remain +in TypeScript. No extraction, restore, writer stop or execution grant occurs. +""" +from __future__ import annotations + +import hashlib +import json +from pathlib import Path, PurePosixPath +import stat +import tarfile +from typing import Any, BinaryIO + + +def _sha(handle: BinaryIO) -> str: + digest = hashlib.sha256() + for block in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(block) + return digest.hexdigest() + + +def read_cold_source_backup(manifest_path: Path) -> dict[str, Any]: + """Witness the existing backup format, including pre-member-list archives. + + A checksum of a manifest alone is insufficient. Read the same open archive, + compare its embedded source map and hash the bytes actually stored in each + regular member (including tar hardlinks). Symlinks grant no byte coverage. + """ + manifest_path = manifest_path.expanduser().absolute() + if not stat.S_ISREG(manifest_path.lstat().st_mode): + raise ValueError("cold_import_backup_manifest_unsafe") + manifest_bytes = manifest_path.read_bytes() + manifest = json.loads(manifest_bytes) + if (not isinstance(manifest, dict) or manifest.get("schema_version") != "loopx_state_backup_v0" + or manifest.get("dry_run") is not False or manifest.get("execute_requested") is not True): + raise ValueError("cold_import_backup_not_executed") + execution = manifest.get("execution") + if not isinstance(execution, dict): + raise ValueError("cold_import_backup_not_executed") + archive_path = Path(execution["archive_path"]).expanduser().absolute() + if not stat.S_ISREG(archive_path.lstat().st_mode): + raise ValueError("cold_import_backup_archive_unsafe") + with archive_path.open("rb") as archive: + original_sha = _sha(archive) + if original_sha != execution.get("archive_sha256"): + raise ValueError("cold_import_backup_archive_changed") + archive.seek(0) + members: list[dict[str, Any]] = [] + seen: set[str] = set() + with tarfile.open(fileobj=archive, mode="r:gz") as tar: + embedded: dict[str, Any] | None = None + for member in tar: + name = member.name + if (name in seen or name.startswith("/") or "\\" in name + or ".." in PurePosixPath(name).parts): + raise ValueError("cold_import_backup_member_unsafe") + seen.add(name) + if name == "manifest.json": + if not member.isfile(): + raise ValueError("cold_import_backup_manifest_unsafe") + stream = tar.extractfile(member) + if stream is None: + raise ValueError("cold_import_backup_manifest_missing") + with stream: + embedded = json.load(stream) + continue + if member.isfile() or member.islnk(): + stream = tar.extractfile(member) + if stream is None: + raise ValueError("cold_import_backup_member_unreadable") + with stream: + digest = _sha(stream) + members.append({"archive_path": name, "sha256": digest}) + if not isinstance(embedded, dict) or any(embedded.get(key) != manifest.get(key) + for key in ("schema_version", "runtime_root", "project", "included", "configuration_source_registry")): + raise ValueError("cold_import_backup_source_map_changed") + archive.seek(0) + if _sha(archive) != original_sha or manifest_path.read_bytes() != manifest_bytes: + raise ValueError("cold_import_backup_changed_retry") + return {"schema_version": "loopx_state_backup_source_witness_v0", + "archive_path": str(archive_path), "archive_sha256": original_sha, + "manifest_path": str(manifest_path), "manifest_sha256": hashlib.sha256(manifest_bytes).hexdigest(), + "included": [{"source_path": item["source_path"], "archive_path": item["archive_path"]} + for item in manifest["included"]], "members": members} diff --git a/loopx/control_plane/coordination/cold_source_import.ts b/loopx/control_plane/coordination/cold_source_import.ts new file mode 100644 index 0000000000..d86d11375c --- /dev/null +++ b/loopx/control_plane/coordination/cold_source_import.ts @@ -0,0 +1,342 @@ +/** Reviewed coordination cutover for a stopped source without active capture. + * The recovery carrier preserves source bytes; it is not a whole-Goal backup. + * Host shutdown is an explicit operator attestation, never inferred from locks. + */ +import {createHash} from "node:crypto"; +import {createReadStream} from "node:fs"; +import {lstat, readFile, readdir} from "node:fs/promises"; +import {join, isAbsolute, relative, resolve} from "node:path"; +import type {JsonObject} from "../effect_program.ts"; +import {durableWriteJson} from "../effect_runtime_io.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import {BARE_SHA256_PATTERN} from "../content_digest.ts"; +import {canonicalAuthorityBytes, canonicalAuthorityObject, canonicalAuthoritySha256, + hasExactAuthorityKeys, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {FileAuthorityStore} from "./file_authority_store.ts"; +import {localAuthorityProviderPaths, openLocalAuthorityStoreHandle, + requireLocalAuthorityRuntimeRoot, localAuthorityOpenFailure, selectLocalAuthorityTarget, + type LocalAuthorityProviderDependencies} from "./local_authority_provider.ts"; +import {COLD_SOURCE_IMPORT_WRITER_FENCE_SCHEMA, engageLegacyCoordinationWriterFenceUnderLocks, + loadLegacyCoordinationWriterFence} from "./legacy_writer_fence.ts"; +import {readShadowManagementState, ShadowManagementError, shadowManagementDirectory, + withShadowMaintenanceLock} from "./shadow_management.ts"; +import {verifyShadowSourceSnapshot, withShadowSourceLocks, type ShadowRequest} from "./runtime_shadow.ts"; +import {planHandoffPolicyMigration} from "./handoff_policy_migration.ts"; +import {canonicalTaskLease} from "./task_lease_state.ts"; +import {indexCoordinationProjection} from "./coordination_projection.ts"; +import {commitPromotionAndReadBack, readPromotionReceipt} from "./promotion_receipt.ts"; + +export const COLD_SOURCE_IMPORT_REQUEST_SCHEMA = "loopx_cold_source_import_request_v0"; +const PLAN_SCHEMA = "loopx_cold_source_import_plan_v0"; +const RESULT_SCHEMA = "loopx_cold_source_import_result_v0"; + +function reject(code: string): never { throw new ShadowManagementError(code); } +function digest(bytes: Buffer): string { return `sha256:${createHash("sha256").update(bytes).digest("hex")}`; } +function same(left: unknown, right: unknown): boolean { return canonicalAuthorityBytes(left).equals(canonicalAuthorityBytes(right)); } +function carrierPath(root: string, goal: string, operation: string): string { + return join(shadowManagementDirectory(root, goal), "cold-imports", `${canonicalAuthoritySha256(operation)}.json`); +} +function sourceInventory(projection: JsonObject, goal: string): JsonObject { + const index = indexCoordinationProjection(projection, goal); + return {todo_count: index.todos.size, + archived_todo_count: [...index.todos.values()].filter(row => row.archive_state === "archive").length, + lease_count: index.leases.size, source_handoff_mode: projection.handoff_mode}; +} +async function optionalJson(path: string): Promise { + try { return canonicalAuthorityObject(JSON.parse(await readFile(path, "utf8")), "import carrier"); } + catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; throw error; } +} +async function names(path: string): Promise { + try { return (await readdir(path)).sort(); } + catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return []; throw error; } +} + +/** Check every retained lease, including orphans and expired active records. + * A settled source token remains history, never a fresh execution grant. */ +async function requireStoppedSource(request: ShadowRequest): Promise { + const managed = await readShadowManagementState(request.runtime_root, request.goal_id); + if (managed && managed.status !== "inactive") reject("cold_import_capture_requires_disposition"); + if ((await names(join(request.runtime_root, "authority-shadow", "outbox", request.goal_id))).length) { + reject("cold_import_outbox_requires_disposition"); + } + const shadow = new FileAuthorityStore(join(request.runtime_root, "authority-shadow", "file-v0"), request.goal_id, {existingOnly: true}); + if ((await shadow.loadAuthority()).status !== "missing") reject("cold_import_shadow_requires_disposition"); + const directory = join(request.runtime_root, "goals", request.goal_id, "task-leases"); + try { if (!(await lstat(directory)).isDirectory()) reject("cold_import_lease_source_unsafe"); } + catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + for (const name of await names(directory)) { + if (!name.endsWith(".json")) continue; + if (!/^[A-Za-z0-9_.-]+\.json$/.test(name)) reject("cold_import_lease_source_unsupported"); + const path = join(directory, name); + if (!(await lstat(path)).isFile()) reject("cold_import_lease_source_unsafe"); + const raw = canonicalAuthorityObject(JSON.parse(await readFile(path, "utf8")), "source lease"); + const lease = canonicalTaskLease(raw, request.goal_id, name.slice(0, -5)); + if (lease.status !== "released") reject("cold_import_lease_requires_settlement"); + } + await verifyShadowSourceSnapshot(request); +} + +async function sourceWitness(request: ShadowRequest): Promise { + const snapshot = request.source_snapshot; + const inventory = snapshot.lease_inventory as JsonObject[]; + const files = [ + {path: snapshot.state_path, bytes_sha256: snapshot.state_bytes_sha256}, + {path: (snapshot.registry_source as JsonObject).path, bytes_sha256: `sha256:${(snapshot.registry_source as JsonObject).sha256}`}, + ...(snapshot.evidence_files as JsonObject[]), + ...inventory.map(item => ({path: join(request.runtime_root, "goals", request.goal_id, "task-leases", String(item.name)), bytes_sha256: item.bytes_sha256})), + ]; + return Promise.all(files.map(async file => { + let bytes: Buffer | null = null; + try { + if (!(await lstat(String(file.path))).isFile()) reject("cold_import_source_unsafe"); + bytes = await readFile(String(file.path)); + } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + if ((bytes === null ? null : digest(bytes)) !== file.bytes_sha256) reject("source_changed_retry"); + return {...file, bytes_base64: bytes === null ? null : bytes.toString("base64")}; + })); +} + +/** The trusted Host tar adapter witnesses actual archived members. This typed + * owner checks source coverage and pins both saved artifacts through apply. + * It qualifies coordination-source bytes, not complete-state reactivation. */ +async function verifySourceBackup(value: unknown, root: string, sources: JsonObject[] = []): Promise { + const backup = requireJsonObject(value, "source backup"); + if (!hasExactAuthorityKeys(backup, ["schema_version", "archive_path", "archive_sha256", + "manifest_path", "manifest_sha256", "included", "members"]) || + backup.schema_version !== "loopx_state_backup_source_witness_v0" || + !Array.isArray(backup.included) || !Array.isArray(backup.members)) reject("cold_import_backup_invalid"); + for (const kind of ["archive", "manifest"] as const) { + const path = backup[`${kind}_path`]; + const expected = backup[`${kind}_sha256`]; + if (typeof path !== "string" || !isAbsolute(path) || typeof expected !== "string" || + !BARE_SHA256_PATTERN.test(expected)) reject("cold_import_backup_invalid"); + if (!(await lstat(path)).isFile()) reject("cold_import_backup_unsafe"); + const hash = createHash("sha256"); + for await (const chunk of createReadStream(path)) hash.update(chunk); + if (hash.digest("hex") !== expected) reject("cold_import_backup_changed"); + } + const included = backup.included as JsonObject[]; + if (!included.some(item => item.source_path === root)) reject("cold_import_backup_runtime_missing"); + const members = new Map(); + for (const item of backup.members as JsonObject[]) { + if (!hasExactAuthorityKeys(item, ["archive_path", "sha256"]) || typeof item.archive_path !== "string" || + typeof item.sha256 !== "string" || !BARE_SHA256_PATTERN.test(item.sha256) || + members.has(item.archive_path)) reject("cold_import_backup_invalid"); + members.set(item.archive_path, item.sha256); + } + for (const item of included) { + if (!hasExactAuthorityKeys(item, ["source_path", "archive_path"]) || + typeof item.source_path !== "string" || !isAbsolute(item.source_path) || + typeof item.archive_path !== "string" || item.archive_path.startsWith("/") || + item.archive_path.split("/").some(part => part === "..") || item.archive_path.includes("\\")) { + reject("cold_import_backup_invalid"); + } + } + for (const file of sources) { + if (file.bytes_sha256 === null) continue; + const matched = included.some(item => { + const suffix = relative(String(item.source_path), String(file.path)); + if (isAbsolute(suffix) || suffix === ".." || suffix.startsWith(`..${process.platform === "win32" ? "\\" : "/"}`)) return false; + const member = suffix === "" ? String(item.archive_path) : `${item.archive_path}/${suffix.replaceAll("\\", "/")}`; + return `sha256:${members.get(member)}` === file.bytes_sha256; + }); + if (!matched) reject("cold_import_backup_source_missing_or_changed"); + } +} + +function verifyCarrier(raw: JsonObject, root: string, goal: string, operation: string, expected: unknown): JsonObject { + if (!hasExactAuthorityKeys(raw, ["plan", "plan_sha256"])) reject("cold_import_carrier_invalid"); + const plan = requireJsonObject(raw.plan, "import plan"); + if (!hasExactAuthorityKeys(plan, ["schema_version", "runtime_root", "goal_id", "operation_id", + "source", "source_witness", "source_backup", "target_projection", "target_provider", "target_store_identity"]) || + plan.schema_version !== PLAN_SCHEMA || plan.runtime_root !== root || plan.goal_id !== goal || + plan.operation_id !== operation || canonicalAuthoritySha256(plan) !== raw.plan_sha256 || + raw.plan_sha256 !== expected) reject("cold_import_reviewed_plan_changed"); + if (!Array.isArray(plan.source_witness)) reject("cold_import_carrier_invalid"); + const source = requireJsonObject(plan.source, "retained import source"); + const projection = requireJsonObject(plan.target_projection, "retained import target"); + if (!hasExactAuthorityKeys(source, ["runtime_root", "goal_id", "projection", "source_snapshot"]) || + source.runtime_root !== root || source.goal_id !== goal || + (projection.handoff_mode !== "soft_claim" && projection.handoff_mode !== "hard_lease")) reject("cold_import_carrier_invalid"); + const snapshot = requireJsonObject(source.source_snapshot, "retained source snapshot"); + const registration = requireJsonObject(snapshot.registry_source, "retained registry source"); + const migration = planHandoffPolicyMigration(requireJsonObject(source.projection, "retained source projection"), + goal, projection.handoff_mode, registration.registered_agents as string[], new Date()); + if (!migration.ready || !same(migration.target_projection, projection)) reject("cold_import_carrier_invalid"); + for (const item of plan.source_witness as JsonObject[]) { + if (!hasExactAuthorityKeys(item, ["path", "bytes_sha256", "bytes_base64"]) || + (item.bytes_base64 === null ? item.bytes_sha256 !== null : + typeof item.bytes_base64 !== "string" || digest(Buffer.from(item.bytes_base64, "base64")) !== item.bytes_sha256)) { + reject("cold_import_carrier_invalid"); + } + } + return plan; +} + +/** Prepare persists an immutable reviewed source/target carrier. Apply + * revalidates it under the primary locks before fencing. Recovery uses only + * the original durable carrier and fence, never a new Markdown snapshot. */ +export async function executeColdSourceImport(value: unknown, + dependencies: LocalAuthorityProviderDependencies = {}): Promise { + let fenced: boolean | null = null; + let attempted = false; + try { + const input = requireJsonObject(value, "cold source import request"); + const action = input.action; + if (input.schema_version !== COLD_SOURCE_IMPORT_REQUEST_SCHEMA || + (action !== "prepare" && action !== "apply" && action !== "recover" && action !== "readback") || + !hasExactAuthorityKeys(input, action === "prepare" + ? ["schema_version", "action", "runtime_root", "goal_id", "operation_id", "projection", "source_snapshot", "target_handoff_mode", "target_provider", "source_backup"] + : action === "apply" + ? ["schema_version", "action", "runtime_root", "goal_id", "operation_id", "expected_plan_sha256", "writers_stopped"] + : ["schema_version", "action", "runtime_root", "goal_id", "operation_id", "expected_plan_sha256"])) { + reject("cold_import_request_invalid"); + } + const root = requireLocalAuthorityRuntimeRoot(input.runtime_root); + const goal = requireAuthorityStoreId(input.goal_id, "goal id"); + if (/[/\\\0]/.test(goal) || goal === "." || goal === "..") reject("cold_import_request_invalid"); + const operation = requireAuthorityStoreId(input.operation_id, "operation id"); + const path = carrierPath(root, goal, operation); + if (action === "prepare") { + if (input.target_provider !== "file" && input.target_provider !== "sqlite") reject("cold_import_provider_unsupported"); + if (input.target_handoff_mode !== "soft_claim" && input.target_handoff_mode !== "hard_lease") reject("cold_import_handoff_mode_invalid"); + // Empty-target selection owns its own M lock, as on canonical creation. + // Preflight does not replace the locked source/backup checks below. It + // rejects absent source/backup before publishing target identity metadata. + const source: ShadowRequest = {runtime_root: root, goal_id: goal, + projection: canonicalAuthorityObject(input.projection, "source projection"), + source_snapshot: canonicalAuthorityObject(input.source_snapshot, "source snapshot")}; + await requireStoppedSource(source); + await verifySourceBackup(input.source_backup, root, await sourceWitness(source)); + await selectLocalAuthorityTarget(root, goal, input.target_provider, true); + } + return await withShadowMaintenanceLock(root, goal, async () => { + const prior = await loadLegacyCoordinationWriterFence(root, goal); + fenced = prior.status === "loaded" ? true : prior.status === "missing" ? false : null; + if (prior.status === "failed") reject(prior.reason_code); + const opened = await openLocalAuthorityStoreHandle(root, goal, dependencies, + {existingOnly: action !== "prepare"}); + if (opened.provider !== "file" && opened.provider !== "sqlite") reject("cold_import_provider_unsupported"); + const targetIdentity = await opened.store.storeIdentity(); + if (targetIdentity.status !== "available") reject(targetIdentity.reason_code); + if (action === "prepare") { + if (prior.status !== "missing") reject("cold_import_already_fenced"); + if (input.target_handoff_mode !== "soft_claim" && input.target_handoff_mode !== "hard_lease") reject("cold_import_handoff_mode_invalid"); + const targetMode = input.target_handoff_mode; + const source: ShadowRequest = {runtime_root: root, goal_id: goal, + projection: canonicalAuthorityObject(input.projection, "source projection"), + source_snapshot: canonicalAuthorityObject(input.source_snapshot, "source snapshot")}; + return await withShadowSourceLocks(source, async () => { + if ((await opened.store.loadAuthority()).status !== "missing") reject("cold_import_target_not_empty"); + await requireStoppedSource(source); + const registered = (source.source_snapshot.registry_source as JsonObject).registered_agents as string[]; + const migration = planHandoffPolicyMigration(source.projection, goal, targetMode, registered, new Date()); + if (!migration.ready) reject(migration.reason_code ?? "cold_import_handoff_conflict"); + const witness = await sourceWitness(source); + await verifySourceBackup(input.source_backup, root, witness); + const plan = {schema_version: PLAN_SCHEMA, runtime_root: root, goal_id: goal, operation_id: operation, + source, source_witness: witness, source_backup: input.source_backup, target_projection: migration.target_projection, + target_provider: opened.provider, target_store_identity: targetIdentity.store_identity}; + const carrier = {plan, plan_sha256: canonicalAuthoritySha256(plan)}; + const existing = await optionalJson(path); + if (existing !== null && !same(existing, carrier)) reject("cold_import_operation_conflict"); + if (existing === null) await durableWriteJson(path, carrier); + verifyCarrier((await optionalJson(path))!, root, goal, operation, carrier.plan_sha256); + return {schema_version: RESULT_SCHEMA, ok: true, status: "prepared", operation_id: operation, + plan_path: path, plan_sha256: carrier.plan_sha256, target_provider: opened.provider, + source_inventory: sourceInventory(source.projection, goal), authority_changed: false, + target_handoff_mode: targetMode, legacy_writer_fenced: false, + execution_authority_granted: false, coordination_source_backup_verified: true, complete_goal_backup_verified: false}; + }); + } + const raw = await optionalJson(path); + if (raw === null) reject("cold_import_carrier_missing"); + const plan = verifyCarrier(raw, root, goal, operation, input.expected_plan_sha256); + if (opened.provider !== plan.target_provider || targetIdentity.store_identity !== plan.target_store_identity) { + reject("cold_import_target_identity_changed"); + } + const fence = {schema_version: COLD_SOURCE_IMPORT_WRITER_FENCE_SCHEMA, state: "engaged", goal_id: goal, + fence_id: `cold-import:${operation}`, import_operation_id: operation, import_plan_sha256: raw.plan_sha256}; + if (prior.status === "loaded" && !same(prior.fence, fence)) reject("cold_import_fence_conflict"); + const source = plan.source as ShadowRequest; + const projection = requireJsonObject(plan.target_projection, "import target projection"); + const identity = {operation_id: operation, projection_sha256: canonicalAuthoritySha256(projection), + receipt: {schema_version: "loopx_cold_source_import_receipt_v0", goal_id: goal, + operation_id: operation, import_plan_sha256: raw.plan_sha256, source_projection_sha256: canonicalAuthoritySha256(source.projection)}}; + const marker = {operation_id: operation, import_plan_sha256: raw.plan_sha256, + target_store_identity: targetIdentity.store_identity, receipt_sha256: canonicalAuthoritySha256(identity.receipt)}; + const completed = await optionalJson(`${path}.completed.json`); + if (completed !== null && !same(completed, marker)) reject("cold_import_completion_identity_changed"); + if (action === "readback") { + // Page reload observes the original intent; it never completes a + // partially committed operation or treats a preview as confirmation. + const original = prior.status === "loaded" ? await readPromotionReceipt(opened.store, identity) : null; + if (!original?.matched) { + if (completed !== null) reject("cold_import_completed_authority_missing"); + const head = await opened.store.loadAuthority(); + if (head.status !== "missing") reject(original?.reason_code ?? "cold_import_target_not_empty"); + } + return {schema_version: RESULT_SCHEMA, ok: true, + status: original?.matched ? "replayed" : "prepared", + operation_id: operation, plan_sha256: raw.plan_sha256, target_provider: opened.provider, + target_handoff_mode: projection.handoff_mode, source_inventory: sourceInventory(source.projection, goal), + legacy_writer_fenced: prior.status === "loaded", authority_changed: false, + execution_authority_granted: false, complete_goal_backup_verified: false}; + } + // File's selected identity already exists; allow only the first head in + // that exact identity, never recreation of a lost identity. + const store = opened.provider === "file" && dependencies.createStore === undefined + ? new FileAuthorityStore(localAuthorityProviderPaths(root, goal).file, goal, {expectedIdentity: targetIdentity.store_identity}) + : opened.store; + const complete = async () => { + const original = await readPromotionReceipt(store, identity); + if (original.matched) { + if (completed === null) await durableWriteJson(`${path}.completed.json`, marker); + return {status: "replayed", ...original}; + } + if (completed !== null) reject("cold_import_completed_authority_missing"); + if ((await store.loadAuthority()).status !== "missing") reject(original.reason_code); + await verifySourceBackup(plan.source_backup, root, plan.source_witness as JsonObject[]); + attempted = true; + const result = await commitPromotionAndReadBack(store, identity, projection, + {...identity.receipt, schema_version: "loopx_cold_source_import_event_v0"}); + if (!result.readback.matched) reject(result.readback.reason_code); + await durableWriteJson(`${path}.completed.json`, marker); + return {status: result.interrupted || result.commit?.status === "ambiguous" ? "recovered" : "applied", ...result.readback}; + }; + let result: JsonObject; + if (action === "recover") { + if (prior.status !== "loaded") reject("cold_import_recovery_requires_fence"); + result = await complete(); + } else { + if (input.writers_stopped !== true) reject("cold_import_operator_stop_confirmation_required"); + result = await withShadowSourceLocks(source, async () => { + if (prior.status === "missing") { + await requireStoppedSource(source); + if (!same(await sourceWitness(source), plan.source_witness)) reject("source_changed_retry"); + await verifySourceBackup(plan.source_backup, root, plan.source_witness as JsonObject[]); + if ((await store.loadAuthority()).status !== "missing") reject("cold_import_target_not_empty"); + const applied = await engageLegacyCoordinationWriterFenceUnderLocks(root, goal, + String(source.source_snapshot.state_path), fence); + if (applied.status !== "applied" && applied.status !== "replayed") reject("cold_import_fence_unverified"); + fenced = true; + } + return await complete(); + }); + } + return {schema_version: RESULT_SCHEMA, ok: true, ...result, operation_id: operation, + authority_changed: attempted, target_handoff_mode: projection.handoff_mode, + source_inventory: sourceInventory(source.projection, goal), + executed: attempted, + plan_sha256: raw.plan_sha256, target_provider: opened.provider, legacy_writer_fenced: true, + stop_confirmation_source: "operator_attestation", execution_authority_granted: false, + complete_goal_backup_verified: false, legacy_fallback_used: false}; + }); + } catch (error) { + return {schema_version: RESULT_SCHEMA, ok: false, status: "failed", executed: attempted, + authority_changed: attempted ? null : false, execution_authority_granted: false, + reason_code: error instanceof ShadowManagementError ? error.reason_code : "cold_import_unavailable", + reason: error instanceof Error ? error.message : "cold import unavailable", + legacy_writer_fenced: fenced, legacy_fallback_used: false, ...localAuthorityOpenFailure(error)}; + } +} diff --git a/loopx/control_plane/coordination/cold_source_inspection.ts b/loopx/control_plane/coordination/cold_source_inspection.ts index 156aef079e..62822b0191 100644 --- a/loopx/control_plane/coordination/cold_source_inspection.ts +++ b/loopx/control_plane/coordination/cold_source_inspection.ts @@ -12,7 +12,7 @@ import {withShadowMaintenanceLock, ShadowManagementError, readShadowManagementSt import {readLocalAuthorityShadow, LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA} from "./local_authority_shadow.ts"; import {drainInventory} from "./shadow_drain_files.ts"; import {planShadowDrain, SHADOW_DRAIN_PLAN_REQUEST_SCHEMA} from "./shadow_drain_plan.ts"; -import {localAuthorityProviderPaths} from "./local_authority_provider.ts"; +import {decodeLocalAuthoritySelection, localAuthorityProviderPaths, openLocalAuthorityStoreHandle} from "./local_authority_provider.ts"; import {FileAuthorityStore} from "./file_authority_store.ts"; export const COLD_SOURCE_INSPECTION_REQUEST_SCHEMA = "loopx_cold_source_inspection_request_v0"; @@ -64,7 +64,23 @@ export async function inspectColdCoordinationSource(value: unknown): Promise; + try { selection = decodeLocalAuthoritySelection(JSON.parse(selected), request.goal_id); } + catch { throw new ShadowManagementError("cold_source_canonical_authority_present"); } + if (selection.provider !== "sqlite") throw new ShadowManagementError("cold_source_canonical_authority_present"); + const opened = await openLocalAuthorityStoreHandle(request.runtime_root, request.goal_id, {}, {existingOnly: true}); + if ((await opened.store.loadAuthority()).status !== "missing") { + throw new ShadowManagementError("cold_source_canonical_authority_present"); + } + } + for (const path of [new FileAuthorityStore(paths.file, request.goal_id, {existingOnly: true}).path]) { try { await lstat(path); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") continue; diff --git a/loopx/control_plane/coordination/legacy_writer_fence.ts b/loopx/control_plane/coordination/legacy_writer_fence.ts index 2f9558a2b8..d28826f7ff 100644 --- a/loopx/control_plane/coordination/legacy_writer_fence.ts +++ b/loopx/control_plane/coordination/legacy_writer_fence.ts @@ -35,6 +35,9 @@ export { LEGACY_COORDINATION_WRITE_CHECK_RESULT_SCHEMA, }; +/** A cold import pins its own source and operation; it has no shadow revision. */ +export const COLD_SOURCE_IMPORT_WRITER_FENCE_SCHEMA = "loopx_cold_source_import_writer_fence_v0"; + // Caller adapter: remediation is rendered here, never inside // checkLegacyCoordinationWriteAllowed, which owns only the stable typed reason // and the fence binding facts. Tokens are substituted in one pass, so a data @@ -107,6 +110,15 @@ export function legacyCoordinationWriterFencePath(root: string, goalId: string): export function decodeLegacyCoordinationWriterFence(value: unknown): JsonObject { const fence = canonicalAuthorityObject(value, "legacy coordination writer fence"); + if (fence.schema_version === COLD_SOURCE_IMPORT_WRITER_FENCE_SCHEMA) { + if (fence.state !== "engaged" || !hasExactAuthorityKeys(fence, + ["schema_version", "state", "goal_id", "fence_id", "import_operation_id", "import_plan_sha256"]) || + !BARE_SHA256_PATTERN.test(String(fence.import_plan_sha256))) { + throw new Error("Invalid cold source import writer fence"); + } + for (const key of ["goal_id", "fence_id", "import_operation_id"]) requireAuthorityStoreId(fence[key], key); + return fence; + } if (fence.schema_version === NEW_GOAL_WRITER_FENCE_SCHEMA) { if (fence.state !== "engaged" || !hasExactAuthorityKeys(fence, ["schema_version", "state", "goal_id", "fence_id", "creation_operation_id", "creation_identity_sha256", "creation_completed"]) || diff --git a/loopx/control_plane/coordination/runtime_shadow.py b/loopx/control_plane/coordination/runtime_shadow.py index d6ec31b35b..9b0a1b65ae 100644 --- a/loopx/control_plane/coordination/runtime_shadow.py +++ b/loopx/control_plane/coordination/runtime_shadow.py @@ -345,7 +345,7 @@ def _build_runtime_shadow_source_snapshot( """ from ...rollout_event_log import ROLLOUT_EVENT_SCHEMA_VERSION, rollout_event_log_path from ...paths import resolve_runtime_root - from ...state_refresh import resolve_goal_state + from ..goals.state_resolution import resolve_goal_state from ..goals.legacy_event_source import state_event_log_candidates from ..todos.active_state_todo_parser import parse_active_state_todos from ..todos.goal_todo_projection import todo_summaries_from_fields diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 2f911e193f..6c655bd223 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -465,6 +465,7 @@ export function createEffectRuntimeHandlers( ["coordination.runtime_shadow.commit_entry", lazyHandler(() => import("./coordination/shadow_entry_delivery.ts"), ({deliverShadowEntry}) => deliverShadowEntry)], ["coordination.runtime_shadow.outbox_read", lazyHandler(() => import("./coordination/local_authority_shadow.ts"), ({readLocalAuthorityShadow}) => readLocalAuthorityShadow)], ["coordination.runtime_shadow.drain", lazyHandler(() => import("./coordination/shadow_drain.ts"), ({drainShadowOutbox}) => drainShadowOutbox)], + ["coordination.cold_source.import", lazyHandler(() => Promise.all([import("./coordination/cold_source_import.ts"), import("./coordination/source_transfer.ts")]), ([{executeColdSourceImport}, {withCoordinationSourceTransfer}]) => withCoordinationSourceTransfer("coordination.cold_source.import", executeColdSourceImport))], [ "effect.program_from_ordered_steps", (params) => effectProgramFromOrderedSteps( diff --git a/loopx/presentation/configuration_api.py b/loopx/presentation/configuration_api.py index 9b24233034..bfc3110351 100644 --- a/loopx/presentation/configuration_api.py +++ b/loopx/presentation/configuration_api.py @@ -40,6 +40,9 @@ def _configuration_post_routes(self) -> dict[str, Callable[[], None]]: f"{storage_api.CHAT_GOAL_STORAGE_PATH}/preview": lambda: self._storage_update(action="plan-migration"), f"{storage_api.CHAT_GOAL_STORAGE_PATH}/apply": lambda: self._storage_update(action="migrate"), f"{storage_api.CHAT_GOAL_STORAGE_PATH}/recover": lambda: self._storage_update(action="migration-readback"), + f"{storage_api.CHAT_GOAL_STORAGE_PATH}/import/preview": lambda: self._storage_import(action="prepare"), + f"{storage_api.CHAT_GOAL_STORAGE_PATH}/import/apply": lambda: self._storage_import(action="apply"), + f"{storage_api.CHAT_GOAL_STORAGE_PATH}/import/recover": lambda: self._storage_import(action="readback"), f"{backup_api.CONFIGURATION_BACKUP_PATH}/export": self._configuration_backup_export, f"{backup_api.CONFIGURATION_BACKUP_PATH}/restore": self._configuration_backup_restore, f"{ownership_api.CHAT_GOAL_OWNERSHIP_PATH}/preview": lambda: self._ownership_update(execute=False), diff --git a/loopx/presentation/goal_storage_api.py b/loopx/presentation/goal_storage_api.py index 2d146aa732..78607fe09d 100644 --- a/loopx/presentation/goal_storage_api.py +++ b/loopx/presentation/goal_storage_api.py @@ -2,11 +2,13 @@ from __future__ import annotations import re +from pathlib import Path from typing import Any, TYPE_CHECKING, cast from urllib.parse import parse_qs, urlparse from uuid import uuid4 from ..control_plane.effect_runtime import effect_runtime_result +from ..control_plane.coordination.local_authority_shadow_projection import source_effect_runtime_result CHAT_GOAL_STORAGE_PATH = "/api/chat/goal-storage" @@ -44,6 +46,8 @@ def _storage_send(self, result: dict[str, Any], goal_id: str, preview_id: str | "ok", "status", "authority_changed", "execution_authority_granted", "plan_sha256", "reviewed_source", "target_provider", "selected_provider", "current", "recovery", "reason_code", + "operation_id", "source_inventory", "target_handoff_mode", + "legacy_writer_fenced", "coordination_source_backup_verified", "complete_goal_backup_verified", "cold_source", ) if key in result} payload["goal_id"] = goal_id @@ -114,3 +118,59 @@ def _storage_update(self, *, action: str) -> None: self._send_error("Invalid Goal storage request.", status=400, error_code="invalid_goal_storage_request") except Exception: # noqa: BLE001 - preserve the original carrier on ambiguity. self._send_error("Result unavailable. Read or retry the original preview.", status=503, error_code="goal_storage_unavailable") + + def _storage_import(self, *, action: str) -> None: + """Adapt registered source/backup IO; the shared TS owner owns cutover. + + No caller filenames, overrides or client-side source summaries enter + the owner. Readback is read-only, including after an interrupted apply. + """ + try: + body = self._read_json() + allowed = ({"goal_id", "provider", "handoff_mode"} if action == "prepare" else + {"goal_id", "operation_id", "plan_sha256", "writers_stopped"} if action == "apply" else + {"goal_id", "operation_id", "plan_sha256"}) + if set(body) != allowed or not isinstance(body.get("goal_id"), str): + raise ValueError("invalid import request") + goal_id = body["goal_id"] + registry, goal = self._registry_and_goal(goal_id) + operation = uuid4().hex if action == "prepare" else _token(body, "operation_id", 32) + root = self.server.runtime_root.resolve() + request: dict[str, Any] = {"schema_version": "loopx_cold_source_import_request_v0", + "action": action, "runtime_root": str(root), "goal_id": goal_id, "operation_id": operation} + if action == "prepare": + from ..control_plane.goals.state_resolution import resolve_goal_state + from ..control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot + from ..control_plane.coordination.cold_source_backup import read_cold_source_backup + from ..state_backup import build_state_backup_plan, execute_state_backup_plan + + _, project, state = resolve_goal_state(registry=registry, goal_id=goal_id, + project_override=None, state_file_override=None) + if project is None: + raise ValueError("registered project is required") + registry_path = Path(self.server.registry_path) + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=root, + state_path=state, registry_path=registry_path, include_all_archived_todos=True) + backup = execute_state_backup_plan(build_state_backup_plan(project=project, runtime_root=root, + output_dir=root / "backups" / "cold-import", backup_id=operation, + include_automations=False, include_skills=False, include_registry_projects=False, + registry_path=registry_path)) + request.update(projection=projection, source_snapshot=snapshot, + target_provider=body["provider"], target_handoff_mode=body["handoff_mode"], + source_backup=read_cold_source_backup(Path(backup["manifest_path"]))) + else: + request["expected_plan_sha256"] = _token(body, "plan_sha256", 64) + if action == "apply": + request["writers_stopped"] = body["writers_stopped"] + result = source_effect_runtime_result("coordination.cold_source.import", request, + timeout=300.0, retry_safe=False) + # Read the current store independently. Original intent/receipt + # success never certifies a later provider or hides a read failure. + observed = self._storage_owner(goal_id, action="migration-readback") + result["current"] = observed.get("current") if observed.get("ok") else None + self._storage_send(result, goal_id) + except (KeyError, TypeError, ValueError): + self._send_error("Invalid Goal import request.", status=400, error_code="invalid_goal_storage_request") + except Exception: # noqa: BLE001 - preserve the original carrier and private IO errors. + self._send_error("Import result unavailable. Keep and read the original preview.", + status=503, error_code="goal_storage_unavailable") diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index d11f18fb2d..0363e3e9f3 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -663,7 +663,7 @@ }, { "site": "loopx/cli_commands/coordination_shadow.py::.handle_coordination_shadow_command::codec_read:load_project_registry#1", - "line": 231, + "line": 257, "column": 20, "kind": "codec_read", "api": "load_project_registry", diff --git a/loopx/state_backup.py b/loopx/state_backup.py index 868f8681fc..46196a9134 100644 --- a/loopx/state_backup.py +++ b/loopx/state_backup.py @@ -154,6 +154,9 @@ def add(key: str, source_path: Path, archive_path: str) -> None: warnings.extend(target_warnings) add("runtime_root", runtime_root, "runtime-root") + # Import must bind the original registry bytes, even when its caller-owned + # route lives outside .loopx; configuration projection is not that source. + add("configuration_source_registry", configuration_source_registry, "configuration/registry.source.json") add("project_loopx", project / ".loopx", "project/.loopx") add("project_codex_goals", project / ".codex" / "goals", "project/.codex/goals") add("project_claude_goals", project / ".claude" / "goals", "project/.claude/goals") diff --git a/tests/control_plane/test_backup_registered_source_retention.py b/tests/control_plane/test_backup_registered_source_retention.py index 6dee9eeeed..ab16e8f98a 100644 --- a/tests/control_plane/test_backup_registered_source_retention.py +++ b/tests/control_plane/test_backup_registered_source_retention.py @@ -5,11 +5,12 @@ import subprocess import sys import tarfile +from pathlib import Path import pytest from loopx.control_plane.effect_runtime import restart_effect_runtime -from loopx.state_backup import build_state_backup_plan +from loopx.state_backup import build_state_backup_plan, execute_state_backup_plan from tests.control_plane.canonical_authority_fixture import isolate_sqlite_runtime @@ -119,3 +120,21 @@ def test_global_inventory_still_follows_registered_other_project(tmp_path): plan = build_state_backup_plan(project=project, runtime_root=runtime, include_skills=False, include_automations=False) assert any(row["key"] == "registry_active_state:other" for row in plan["included"]) + + +def test_backup_retains_original_registry_outside_conventional_directories(tmp_path): + project = tmp_path / "project" + project.mkdir() + runtime = tmp_path / "runtime" + runtime.mkdir() + registry = tmp_path / "original-registry.json" + original = b'{"goals": [], "common_runtime_root": "retained route"}\n' + registry.write_bytes(original) + result = execute_state_backup_plan(build_state_backup_plan(project=project, runtime_root=runtime, + output_dir=tmp_path / "saved", backup_id="original-registry", include_automations=False, + include_skills=False, include_registry_projects=False, registry_path=registry)) + assert result["ok"] + manifest = json.loads(Path(result["manifest_path"]).read_text()) + source = next(item for item in manifest["included"] if item["source_path"] == str(registry)) + with tarfile.open(result["archive_path"], "r:gz") as archive: + assert archive.extractfile(source["archive_path"]).read() == original diff --git a/tests/control_plane/test_chat_cold_source_import.py b/tests/control_plane/test_chat_cold_source_import.py new file mode 100644 index 0000000000..8ffe2bf2c6 --- /dev/null +++ b/tests/control_plane/test_chat_cold_source_import.py @@ -0,0 +1,132 @@ +"""Packaged-compatible HTTP transport with real cold source, backup and stores.""" +from __future__ import annotations + +import http.client +import json +import threading + +import pytest + +from loopx.chat_server import ChatHTTPServer, ChatRequestHandler +from loopx.control_plane.effect_runtime import effect_runtime_result, restart_effect_runtime +from test_cold_source_import_cli import workspace + + +@pytest.fixture +def cold_api(tmp_path, monkeypatch): + _, state, _, body, runtime, _, _ = workspace(tmp_path, monkeypatch) + registry = tmp_path / "project/.loopx/registry.json" + server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler) + server.runtime_root = runtime + server.registry_path = registry + server.runtime_root_override = str(runtime) + server.verbose = False + worker = threading.Thread(target=server.serve_forever, daemon=True) + worker.start() + + def call(action="", payload=None, origin=None): + client = http.client.HTTPConnection("127.0.0.1", server.server_port, timeout=90) + path = "/api/chat/goal-storage" + (f"/import/{action}" if action else "?goal_id=cold") + headers = {"Content-Type": "application/json"} + if origin: + headers["Origin"] = origin + client.request("POST" if payload is not None else "GET", path, + body=json.dumps(payload) if payload is not None else None, headers=headers) + response = client.getresponse() + value = json.loads(response.read()) + code = response.status + client.close() + public = json.dumps(value) + assert str(tmp_path) not in public and body not in public + assert "plan_path" not in value and "source_snapshot" not in value + return code, value + + def read(): + return effect_runtime_result("coordination.local_authority.todo_list", { + "schema_version": "loopx_local_coordination_todo_list_request_v0", + "runtime_root": str(runtime), "goal_id": "cold", "role": None, "status": None, + "todo_id": None, "agent_id": None, "limit": None}) + + yield call, read, state, body, runtime + server.shutdown() + worker.join(5) + server.server_close() + restart_effect_runtime() + + +def prepare(call, provider="sqlite"): + code, value = call("preview", {"goal_id": "cold", "provider": provider, "handoff_mode": "hard_lease"}) + assert code == 200 and value["ok"], value + return value + + +def carrier(plan): + return {key: plan[key] for key in ("goal_id", "operation_id", "plan_sha256")} + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_cold_http_preview_reload_confirm_original_recovery(cold_api, provider): + call, read, state, body, runtime = cold_api + original = state.read_bytes() + assert call()[1]["current"]["canonical"] is False + plan = prepare(call, provider) + assert plan["source_inventory"] == {"todo_count": 2, "archived_todo_count": 1, + "lease_count": 0, "source_handoff_mode": "legacy"} + assert plan["coordination_source_backup_verified"] is True + assert plan["complete_goal_backup_verified"] is False + assert plan["authority_changed"] is False and not plan["legacy_writer_fenced"] + saved = carrier(plan) + restart_effect_runtime() + code, observed = call("recover", saved) + assert code == 200 and observed["status"] == "prepared", observed + assert observed["operation_id"] == plan["operation_id"] + assert not observed["legacy_writer_fenced"] and not observed["current"]["canonical"] + assert state.read_bytes() == original + assert not list(runtime.rglob("writer-fence.json")) + code, refused = call("apply", {**saved, "writers_stopped": False}) + assert code == 409 and refused["reason_code"] == "cold_import_operator_stop_confirmation_required" + code, applied = call("apply", {**saved, "writers_stopped": True}) + assert code == 200 and applied["status"] == "applied", applied + assert applied["execution_authority_granted"] is False + assert applied["current"]["provider"] == provider + rows = read()["todos"] + assert next(row for row in rows if row["todo_id"] == "todo_current")["text"] == body + assert next(row for row in rows if row["todo_id"] == "todo_archived")["evidence"] == "original" + from loopx.todos import add_goal_todo + registry = state.parents[3] / ".loopx/registry.json" + add_goal_todo(registry_path=registry, goal_id="cold", role="agent", text="Keep later write", + claimed_by="agent-a", note="Original metadata") + later = read() + state.unlink() + restart_effect_runtime() + code, recovered = call("recover", saved) + assert code == 200 and recovered["status"] == "replayed", recovered + assert recovered["current"]["todo_count"] == 3 + assert read() == later + assert call("apply", {**saved, "writers_stopped": True})[0] == 200 + assert read() == later + + +def test_cold_http_changed_source_backup_and_untrusted_inputs(cold_api): + call, read, state, _, runtime = cold_api + plan = prepare(call) + saved = carrier(plan) + assert call("apply", {**saved, "writers_stopped": True, "plan": "injected"})[0] == 400 + assert call("recover", {**saved, "operation_id": "../wrong"})[0] == 400 + assert call("recover", {**saved, "goal_id": "unknown"})[0] == 400 + assert call("preview", {"goal_id": "cold", "provider": "sqlite", "handoff_mode": "hard_lease"}, + "https://untrusted.example")[0] == 403 + assert call("preview", {"goal_id": "cold", "provider": "postgresql", "handoff_mode": "hard_lease"})[0] == 409 + assert call("preview", {"goal_id": "cold", "provider": "sqlite", "handoff_mode": "legacy"})[0] == 409 + state.write_text(state.read_text() + "\nChanged after preview\n") + code, refused = call("apply", {**saved, "writers_stopped": True}) + assert code == 409 and refused["reason_code"] == "source_changed_retry", refused + assert not list(runtime.rglob("writer-fence.json")) + fresh = prepare(call) + archives = list((runtime / "backups/cold-import").glob(f"*{fresh['operation_id']}*.tar.gz")) + assert len(archives) == 1 + archives[0].write_bytes(b"Damaged after review") + code, refused = call("apply", {**carrier(fresh), "writers_stopped": True}) + assert code == 409 and refused["reason_code"] == "cold_import_backup_changed", refused + assert call()[1]["current"]["canonical"] is False + assert not list(runtime.rglob("writer-fence.json")) diff --git a/tests/control_plane/test_cold_source_import_cli.py b/tests/control_plane/test_cold_source_import_cli.py new file mode 100644 index 0000000000..a661677801 --- /dev/null +++ b/tests/control_plane/test_cold_source_import_cli.py @@ -0,0 +1,212 @@ +"""Real cold import CLI with native backup, source adapter and File/SQLite. + +Synthetic disposable Goals; the cold-import command runs with the old normal +producers physically absent. This does not qualify attached-Host stop, removal +of their remaining callers, or full Goal history recovery. +""" +from __future__ import annotations + +import json +import os +from pathlib import Path +import shutil +import subprocess +import sys + +import pytest + +from canonical_authority_fixture import isolate_sqlite_runtime +from loopx.control_plane.coordination.cold_source_backup import read_cold_source_backup +from loopx.state_backup import build_state_backup_plan, execute_state_backup_plan + +REPO = Path(__file__).resolve().parents[2] + + +def workspace(tmp_path, monkeypatch): + isolate_sqlite_runtime(tmp_path, monkeypatch) + project, runtime = tmp_path / "project", tmp_path / "runtime" + registry = project / ".loopx/registry.json" + state = project / ".local/goals/cold/ACTIVE_GOAL_STATE.md" + state.parent.mkdir(parents=True) + registry.parent.mkdir(parents=True) + body = " ".join(["Complete cold source requirement"] * 25) + state.write_text("---\ngoal_id: cold\nhandoff_mode: legacy\n---\n\n## Agent Todo\n\n" + f"- [ ] {body}\n \n\n" + "## Completed Work Archive\n\n- [x] Independently retained archived requirement\n" + " \n") + (state.parent / "GOAL.md").write_text("Original human objective, never a new execution grant.\n") + goal = {"id": "cold", "repo": str(project), "state_file": str(state.relative_to(project)), + "coordination": {"registered_agents": ["agent-a"]}} + registry.write_text(json.dumps({"schema_version": 1, "common_runtime_root": str(runtime), "goals": [goal]})) + runtime.mkdir() + backup = execute_state_backup_plan(build_state_backup_plan(project=project, runtime_root=runtime, + output_dir=tmp_path / "backups", backup_id="cold-original", include_automations=False, + include_skills=False, include_registry_projects=False, registry_path=registry)) + assert backup["ok"] + receiver = tmp_path / "receiver" + package = Path(os.environ.get("LOOPX_COLD_IMPORT_PACKAGE", str(REPO / "loopx"))) + shutil.copytree(package, receiver / "loopx", ignore=shutil.ignore_patterns("__pycache__", "*.pyc")) + env = {**os.environ, "PYTHONPATH": str(receiver)} + + def cli(action, *args, success=True, output_format="json", runtime_override=None): + child = subprocess.run([sys.executable, "-c", + "import loopx,runpy,sys; print(loopx.__file__,file=sys.stderr); runpy.run_module('loopx.cli',run_name='__main__')", + "--registry", str(registry), "--runtime-root", str(runtime_override or runtime), "--format", output_format, + "coordination-shadow", action, "--goal-id", "cold", "--operation-id", "cold-original", *map(str, args)], + cwd=tmp_path, env=env, capture_output=True, text=True, timeout=60) + assert str(receiver / "loopx/__init__.py") in child.stderr + assert "Traceback" not in child.stderr, child.stderr + assert (child.returncode == 0) is success, child.stdout + child.stderr + return json.loads(child.stdout) if output_format == "json" else child.stdout + + return cli, state, backup, body, runtime, receiver, env + + +def test_coordination_command_help_and_rejections_match_full_cli(tmp_path): + # Independently compare the existing compatibility parser/handler with the + # selected dispatcher, including diagnostics that deliberately fall back. + arguments = [["coordination-shadow", "--help"], + ["coordination-shadow", "prepare-import", "--help"], + ["coordination-shadow", "apply-import", "--help"], + ["coordination-shadow", "recover-import", "--help"], + ["coordination-shadow", "promote", "--help"], + ["coordination-shadow", "unknown-action"], + ["coordination-shadow", "prepare-import"], + ["coordination-shadow", "prepare-import", "--goal-id", "cold", "--operation-id", "original", + "--backup-manifest", "backup.json", "--provider", "sqlite", "--target-handoff-mode", "legacy"], + ["coordination-shadow", "recover-import", "--goal-id", "cold", "--operation-id", "original", + "--plan-sha256", "abc", "--exec"]] + observations = [] + package = Path(os.environ.get("LOOPX_COLD_IMPORT_PACKAGE", str(REPO / "loopx"))) + for module in ("loopx.entrypoint", "loopx.cli"): + script = f""" +import contextlib,io,json,sys,loopx +from {module} import main +sys.argv[0] = 'loopx' +observations = [] +for argv in {arguments!r}: + out, err = io.StringIO(), io.StringIO() + with contextlib.redirect_stdout(out), contextlib.redirect_stderr(err): + try: + code = main(argv) + except SystemExit as exc: + code = exc.code + observations.append([code, out.getvalue(), err.getvalue()]) +print(json.dumps({{'package': loopx.__file__, 'observations': observations}})) +""" + child = subprocess.run([sys.executable, "-c", script], cwd=tmp_path, + env={**os.environ, "PYTHONPATH": str(package.parent)}, capture_output=True, text=True, check=True, timeout=60) + result = json.loads(child.stdout) + assert Path(result["package"]).parent == package + observations.append(result["observations"]) + assert observations[0] == observations[1] + assert all(row[0] == 0 for row in observations[0][:5]) + assert all(row[0] == 2 for row in observations[0][5:]) + + +def test_unconfigured_shadow_inspection_matches_full_cli_without_import(tmp_path, monkeypatch): + _, state, _, _, runtime, receiver, env = workspace(tmp_path, monkeypatch) + registry = tmp_path / "project/.loopx/registry.json" + original = state.read_bytes() + observations = [] + for module in ("loopx.entrypoint", "loopx.cli"): + child = subprocess.run([sys.executable, "-c", + f"import loopx,sys; print(loopx.__file__,file=sys.stderr); from {module} import main; raise SystemExit(main())", + "--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "coordination-shadow", "inspect", "--goal-id", "cold"], + cwd=tmp_path, env=env, capture_output=True, text=True, timeout=60) + assert str(receiver / "loopx/__init__.py") in child.stderr + assert "Traceback" not in child.stderr + assert child.returncode == 1 + observations.append(child.stdout) + assert observations[0] == observations[1] + result = json.loads(observations[0]) + assert result["error_code"] == "coordination_shadow_not_enabled" + assert result["executed"] is False + assert state.read_bytes() == original + assert not (runtime / "authority-transition").exists() + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("old_producers_present", [True, False]) +def test_cold_cli_import_and_source_free_original_recovery( + tmp_path, monkeypatch, provider, old_producers_present, +): + cli, state, backup, body, runtime, receiver, env = workspace(tmp_path, monkeypatch) + if not old_producers_present: + for path in ("todos.py", "bootstrap.py", "control_plane/coordination/runtime_shadow_writer_adapter.py", + "control_plane/coordination/local_authority_shadow_outbox.py"): + (receiver / "loopx" / path).unlink() + original = state.read_bytes() + prepared = cli("prepare-import", "--backup-manifest", backup["manifest_path"], + "--provider", provider, "--target-handoff-mode", "hard_lease")["cold_import"] + assert prepared["status"] == "prepared", prepared + assert prepared["coordination_source_backup_verified"] is True + assert prepared["complete_goal_backup_verified"] is False + digest = prepared["plan_sha256"] + rendered = cli("prepare-import", "--backup-manifest", backup["manifest_path"], + "--provider", provider, "--target-handoff-mode", "hard_lease", output_format="markdown") + assert digest in rendered and prepared["plan_path"] in rendered + assert Path(prepared["plan_path"]).is_file() + applied = cli("apply-import", "--plan-sha256", digest, "--writers-stopped", "--execute")["cold_import"] + assert applied["status"] == "applied", applied + assert applied["target_provider"] == provider + assert applied["execution_authority_granted"] is False + assert state.read_bytes() == original + # Read the real native owner independently through the receiver package. + child = subprocess.run([sys.executable, "-c", + "import json,sys; from loopx.control_plane.effect_runtime import effect_runtime_result; " + "print(json.dumps(effect_runtime_result('coordination.local_authority.todo_list',json.load(sys.stdin))))"], + input=json.dumps({"schema_version": "loopx_local_coordination_todo_list_request_v0", + "runtime_root": str(runtime), "goal_id": "cold", "role": None, "status": None, + "todo_id": None, "agent_id": None, "limit": None}), + cwd=tmp_path, env=env, capture_output=True, text=True, check=True, timeout=60) + result = json.loads(child.stdout) + assert result["source_authority"] == f"{provider}_v0" + assert {row["todo_id"] for row in result["todos"]} == {"todo_current", "todo_archived"} + assert next(row for row in result["todos"] if row["todo_id"] == "todo_current")["text"] == body + assert next(row for row in result["todos"] if row["todo_id"] == "todo_archived")["evidence"] == "original" + state.unlink() + recovered = cli("recover-import", "--plan-sha256", digest, "--execute")["cold_import"] + assert recovered["status"] == "replayed", recovered + assert recovered["cursor"] == "1" + + +@pytest.mark.parametrize("artifact", ["archive_path", "manifest_path"]) +def test_changed_reviewed_backup_cannot_fence_or_import(tmp_path, monkeypatch, artifact): + cli, _, backup, _, runtime, _, _ = workspace(tmp_path, monkeypatch) + prepared = cli("prepare-import", "--backup-manifest", backup["manifest_path"], + "--provider", "sqlite", "--target-handoff-mode", "soft_claim")["cold_import"] + Path(backup[artifact]).write_bytes(b"Changed saved bytes after review") + refused = cli("apply-import", "--plan-sha256", prepared["plan_sha256"], + "--writers-stopped", "--execute", success=False)["cold_import"] + assert refused["reason_code"] == "cold_import_backup_changed" + assert refused["legacy_writer_fenced"] is False + assert not list((runtime / "authority-transition").rglob("writer-fence.json")) + + +def test_archive_source_map_cannot_be_forged_in_external_manifest(tmp_path, monkeypatch): + _, _, backup, _, _, _, _ = workspace(tmp_path, monkeypatch) + path = Path(backup["manifest_path"]) + manifest = json.loads(path.read_text()) + manifest["included"][0]["source_path"] = str(tmp_path / "wrong-source") + path.write_text(json.dumps(manifest)) + with pytest.raises(ValueError, match="cold_import_backup_source_map_changed"): + read_cold_source_backup(path) + + +def test_runtime_path_alias_uses_the_same_backed_up_source(tmp_path, monkeypatch): + cli, state, backup, _, runtime, _, _ = workspace(tmp_path, monkeypatch) + alias = tmp_path / "runtime-alias" + alias.symlink_to(runtime, target_is_directory=True) + prepared = cli("prepare-import", "--backup-manifest", backup["manifest_path"], + "--provider", "sqlite", "--target-handoff-mode", "soft_claim", + runtime_override=alias)["cold_import"] + assert prepared["status"] == "prepared", prepared + applied = cli("apply-import", "--plan-sha256", prepared["plan_sha256"], + "--writers-stopped", "--execute", runtime_override=alias)["cold_import"] + assert applied["status"] == "applied", applied + state.unlink() + recovered = cli("recover-import", "--plan-sha256", prepared["plan_sha256"], + "--execute", runtime_override=alias)["cold_import"] + assert recovered["status"] == "replayed", recovered diff --git a/tests/control_plane/test_cold_source_import_host.py b/tests/control_plane/test_cold_source_import_host.py new file mode 100644 index 0000000000..47dec1a1fa --- /dev/null +++ b/tests/control_plane/test_cold_source_import_host.py @@ -0,0 +1,165 @@ +"""Cold-import stop ordering with native leases and real owned Host processes. + +Disposable unpromoted sources only. This qualifies the operator-led POSIX stop +journey, not automatic Host discovery, outbox disposition or whole-Goal restore. +""" +from __future__ import annotations + +from datetime import datetime +import json +import os +import subprocess +import sys +import time + +import pytest + +from loopx.state_backup import build_state_backup_plan, execute_state_backup_plan +from test_cold_source_import_cli import workspace + + +HOST = """import os,sys,time +from pathlib import Path +Path(sys.argv[1]).write_text(str(os.getpid())) +while True: time.sleep(.05) +""" + +# Use the receiver's production transport and typed supervisor, including its +# native CLI proof. TERM unwinds the control pipe before the owner returns. +SUPERVISOR = """import json,signal,sys +from pathlib import Path +from loopx.control_plane.turn_driver.host_process_transport import run_host_process +def stop(*_): raise KeyboardInterrupt +signal.signal(signal.SIGTERM,stop) +try: + result=run_host_process(json.loads(sys.argv[1]),project=Path(sys.argv[2]), + input_text='',timeout_seconds=None,delegated_lease=json.loads(sys.argv[3])) + print(json.dumps(result)) +except KeyboardInterrupt: + print(json.dumps({'operator_stopped':True})) +""" + + +def source_execution(tmp_path, monkeypatch, ttl=120): + cli, state, _, _, runtime, receiver, env = workspace(tmp_path, monkeypatch) + # A supported old Markdown source can have explicit hard leases without + # already possessing canonical or shadow authority. + state.write_text(state.read_text().replace("handoff_mode: legacy", "handoff_mode: hard_lease")) + registry = tmp_path / "project/.loopx/registry.json" + argv = [sys.executable, "-m", "loopx.cli", "--registry", str(registry), + "--runtime-root", str(runtime), "--format", "json", "task-lease"] + selected = ["--goal-id", "cold", "--todo-id", "todo_current"] + + def command(action, *args): + child = subprocess.run([*argv, action, *selected, *map(str, args)], + cwd=tmp_path, env=env, capture_output=True, text=True, timeout=60) + assert child.returncode == 0, child.stdout + child.stderr + return json.loads(child.stdout) + + identity = ["--owner", "agent-a", "--idempotency-key", "old-execution"] + acquired = command("acquire", *identity, "--ttl-seconds", ttl) + assert acquired["ok"] and acquired["lease"]["status"] == "active", acquired + context = {"lease": acquired["lease"], "ttl_seconds": ttl, + "read_argv": [*argv, "inspect", *selected], + "renew_argv": [*argv, "renew", *selected, *identity]} + + def release(): + current = command("inspect") + result = command("release", *identity, "--expected-version", current["lease"]["version"]) + assert result["ok"] and result["lease"]["status"] == "released", result + return result["lease"] + + def prepare(provider, success=True): + backup = execute_state_backup_plan(build_state_backup_plan( + project=registry.parents[1], runtime_root=runtime, output_dir=tmp_path / "settled-backups", + backup_id=f"source-{time.monotonic_ns()}", include_automations=False, + include_skills=False, include_registry_projects=False, registry_path=registry)) + assert backup["ok"] + return cli("prepare-import", "--backup-manifest", backup["manifest_path"], + "--provider", provider, "--target-handoff-mode", "hard_lease", success=success)["cold_import"] + + def drain(record): + child = subprocess.run([sys.executable, "-c", + "import sys;from pathlib import Path;from loopx.control_plane.turn_driver.host_process_transport import execution_host_drain;print(execution_host_drain(Path(sys.argv[1])))", + str(record)], cwd=tmp_path, env=env, capture_output=True, text=True, timeout=60, check=True) + return child.stdout.strip() + + def launch(context, marker, record): + return subprocess.Popen([sys.executable, "-c", SUPERVISOR, + json.dumps([sys.executable, "-c", HOST, str(marker)]), str(tmp_path), json.dumps(context)], + cwd=tmp_path, env={**env, "LOOPX_HOST_PROCESS_RECORD": str(record)}, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) + + return cli, command, release, prepare, drain, launch, context, runtime, receiver + + +@pytest.mark.skipif(os.name == "nt", reason="POSIX process-group stop proof") +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_cold_import_stops_host_then_releases_native_lease_before_cutover(tmp_path, monkeypatch, provider): + cli, command, release, prepare, drain, launch, context, runtime, _ = source_execution(tmp_path, monkeypatch) + marker, record = tmp_path / "host.pid", tmp_path / "owned-host.json" + owner = launch(context, marker, record) + try: + deadline = time.monotonic() + 60 + while not marker.exists() and owner.poll() is None and time.monotonic() < deadline: + time.sleep(.05) + assert marker.exists(), owner.communicate(timeout=10) + pid = int(marker.read_text()) + os.kill(pid, 0) + assert drain(record) == "draining" + refused = prepare(provider, success=False) + assert refused["reason_code"] == "cold_import_lease_requires_settlement", refused + assert not list(runtime.rglob("writer-fence.json")) + # Neither prepare nor its rejection stops somebody else's process. + os.kill(pid, 0) + owner.terminate() + out, err = owner.communicate(timeout=30) + assert owner.returncode == 0 and json.loads(out) == {"operator_stopped": True}, (out, err) + assert drain(record) == "drained" + with pytest.raises(ProcessLookupError): + os.kill(pid, 0) + still_active = command("inspect")["lease"] + assert still_active["status"] == "active" + assert still_active["idempotency_key"] == "old-execution" + assert prepare(provider, success=False)["reason_code"] == "cold_import_lease_requires_settlement" + settled = release() + prepared = prepare(provider) + assert prepared["status"] == "prepared" and prepared["source_inventory"]["lease_count"] == 1 + result = cli("apply-import", "--plan-sha256", prepared["plan_sha256"], "--writers-stopped", "--execute")["cold_import"] + assert result["status"] == "applied" and result["execution_authority_granted"] is False + # The original lease is preserved as released history under the new + # provider. Its old grant cannot launch another Host after cutover. + current = command("inspect") + assert current["lease"] == settled + assert current["source_authority"] == f"{provider}_v0" + restarted_marker, restarted_record = tmp_path / "restarted.pid", tmp_path / "restarted-host.json" + restarted = launch(context, restarted_marker, restarted_record) + out, err = restarted.communicate(timeout=60) + assert restarted.returncode == 0, (out, err) + observation = json.loads(out) + assert observation["outcome"] == "cancelled" + assert observation["lease_failure"] == {"reason": "lease_inactive", "boundary": "initial_proof"} + assert not restarted_marker.exists() + # The proof CLI did run and drain; the nested actual Host did not. + assert json.loads(restarted_record.read_text())["phase"] == "not_launched" + assert drain(restarted_record) == "drained" + recovered = cli("recover-import", "--plan-sha256", prepared["plan_sha256"], "--execute")["cold_import"] + assert recovered["status"] == "replayed" and recovered["cursor"] == "1" + assert command("inspect")["lease"] == settled + finally: + if owner.poll() is None: + owner.terminate() + owner.communicate(timeout=30) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_elapsed_source_lease_requires_native_release_before_import(tmp_path, monkeypatch, provider): + _, command, release, prepare, _, _, context, runtime, _ = source_execution(tmp_path, monkeypatch, ttl=1) + expires = datetime.fromisoformat(context["lease"]["expires_at"].replace("Z", "+00:00")).timestamp() + time.sleep(max(0, expires - time.time()) + .05) + current = command("inspect")["lease"] + assert current["status"] == "active" and current["expires_at"] == context["lease"]["expires_at"] + assert prepare(provider, success=False)["reason_code"] == "cold_import_lease_requires_settlement" + assert not list(runtime.rglob("writer-fence.json")) + release() + assert prepare(provider)["status"] == "prepared" diff --git a/tests/control_plane/test_cold_source_import_runtime.py b/tests/control_plane/test_cold_source_import_runtime.py new file mode 100644 index 0000000000..68e3f920e0 --- /dev/null +++ b/tests/control_plane/test_cold_source_import_runtime.py @@ -0,0 +1,100 @@ +"""Real typed runtime/backend with the old normal producers physically absent. + +This qualifies the coordination transaction, not the installed CLI/App journey +or complete backup coverage. The source adapter remains separately owned. +""" +from __future__ import annotations + +import json +import os +from pathlib import Path +import shutil +import subprocess +import sys + +import pytest + +from canonical_authority_fixture import isolate_sqlite_runtime +from loopx.control_plane.coordination.cold_source_backup import read_cold_source_backup +from loopx.state_backup import build_state_backup_plan, execute_state_backup_plan +from loopx.control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot + +REPO = Path(__file__).resolve().parents[2] +METHOD = "coordination.cold_source.import" +SCHEMA = "loopx_cold_source_import_request_v0" + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_cold_import_real_runtime_without_original_producers(tmp_path, monkeypatch, provider): + isolate_sqlite_runtime(tmp_path, monkeypatch) + runtime, registry, state = tmp_path / "runtime", tmp_path / ".loopx/registry.json", tmp_path / ".local/goals/cold-goal/state.md" + registry.parent.mkdir(parents=True) + state.parent.mkdir(parents=True) + runtime.mkdir() + body = " ".join(["Complete cold source body"] * 30) + state.write_text( + "---\ngoal_id: cold-goal\nhandoff_mode: legacy\n---\n\n## Agent Todo\n\n" + f"- [ ] {body}\n \n\n" + "## Completed Work Archive\n\n- [x] Complete archived source body\n" + " \n", + encoding="utf-8", + ) + goal = {"id": "cold-goal", "repo": str(tmp_path), "state_file": str(state.relative_to(tmp_path)), + "coordination": {"registered_agents": ["agent-a"]}} + registry.write_text(json.dumps({"schema_version": 1, "common_runtime_root": str(runtime), "goals": [goal]})) + # Full persisted records are assembled by the shipped source adapter. + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=runtime, + state_path=state, registry_path=registry, include_all_archived_todos=True) + assert {row["todo_id"] for row in projection["todos"]} == {"todo_current", "todo_archived"} + backup = execute_state_backup_plan(build_state_backup_plan(project=tmp_path, runtime_root=runtime, + output_dir=tmp_path / "backup", backup_id="original", include_automations=False, + include_skills=False, include_registry_projects=False, registry_path=registry)) + source_backup = read_cold_source_backup(Path(backup["manifest_path"])) + request = {"schema_version": SCHEMA, "runtime_root": str(runtime), "goal_id": "cold-goal", "operation_id": "cold:original"} + receiver = tmp_path / "receiver" + package = Path(os.environ.get("LOOPX_COLD_IMPORT_PACKAGE", str(REPO / "loopx"))) + shutil.copytree(package, receiver / "loopx", ignore=shutil.ignore_patterns("__pycache__", "*.pyc")) + for path in ("todos.py", "bootstrap.py", "control_plane/coordination/runtime_shadow_writer_adapter.py", + "control_plane/coordination/local_authority_shadow_outbox.py"): + (receiver / "loopx" / path).unlink() + env = {**os.environ, "PYTHONPATH": str(receiver)} + + def call(method, fields): + child = subprocess.run([sys.executable, "-c", + "import json,sys,loopx; from loopx.control_plane.effect_runtime import effect_runtime_result; " + "v=json.load(sys.stdin); print(json.dumps({'package':loopx.__file__,'result':effect_runtime_result(v['method'],v['request'],retry_safe=False)},ensure_ascii=False))"], + input=json.dumps({"method": method, "request": fields}), cwd=tmp_path, env=env, + capture_output=True, text=True, timeout=60, check=True) + output = json.loads(child.stdout) + assert Path(output["package"]).parent == receiver / "loopx" + return output["result"] + + if provider == "sqlite": + module = receiver / "loopx/control_plane/coordination/local_authority_provider.ts" + subprocess.run(["node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + f"import {{selectLocalAuthorityTarget}} from {json.dumps(module.as_uri())}; " + f"await selectLocalAuthorityTarget({json.dumps(str(runtime))},'cold-goal','sqlite',true);"], + cwd=tmp_path, env=env, capture_output=True, text=True, check=True, timeout=30) + original_bytes = state.read_bytes() + prepared = call(METHOD, {**request, "action": "prepare", "projection": projection, + "source_snapshot": snapshot, "target_handoff_mode": "soft_claim", "target_provider": provider, "source_backup": source_backup}) + assert prepared["status"] == "prepared", prepared + applied = call(METHOD, {**request, "action": "apply", "expected_plan_sha256": prepared["plan_sha256"], "writers_stopped": True}) + assert applied["status"] == "applied", applied + assert applied["execution_authority_granted"] is False + assert applied["complete_goal_backup_verified"] is False + assert state.read_bytes() == original_bytes + # Fresh Python process; the carrier and original receipt outlive Markdown. + state.unlink() + replayed = call(METHOD, {**request, "action": "recover", "expected_plan_sha256": prepared["plan_sha256"]}) + assert replayed["status"] == "replayed", replayed + observed = call("coordination.local_authority.todo_list", { + "schema_version": "loopx_local_coordination_todo_list_request_v0", "runtime_root": str(runtime), + "goal_id": "cold-goal", "role": None, "status": None, "todo_id": None, "agent_id": None, "limit": None, + }) + assert observed["status"] == "loaded", observed + assert observed["cursor"] == "1" + assert observed["source_authority"] == f"{provider}_v0" + assert {row["todo_id"] for row in observed["todos"]} == {"todo_current", "todo_archived"} + assert next(row for row in observed["todos"] if row["todo_id"] == "todo_current")["text"] == body + assert next(row for row in observed["todos"] if row["todo_id"] == "todo_archived")["evidence"] == "original" diff --git a/tests/control_plane_ts/cold_source_import.test.ts b/tests/control_plane_ts/cold_source_import.test.ts new file mode 100644 index 0000000000..120b63a317 --- /dev/null +++ b/tests/control_plane_ts/cold_source_import.test.ts @@ -0,0 +1,328 @@ +import assert from "node:assert/strict"; +import {spawn, spawnSync} from "node:child_process"; +import {createHash} from "node:crypto"; +import {once} from "node:events"; +import {setTimeout as delay} from "node:timers/promises"; +import {mkdir, mkdtemp, readFile, rm, symlink, writeFile} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join, relative} from "node:path"; +import test, {type TestContext} from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {COLD_SOURCE_IMPORT_REQUEST_SCHEMA, executeColdSourceImport} from "../../loopx/control_plane/coordination/cold_source_import.ts"; +import {inspectColdCoordinationStorage, COLD_SOURCE_INSPECTION_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/cold_source_inspection.ts"; +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {localAuthorityProviderPaths, openLocalAuthorityStoreHandle, selectLocalAuthorityTarget} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; +import {checkLegacyCoordinationWriteAllowed, legacyCoordinationWriterFencePath, + LEGACY_COORDINATION_WRITE_CHECK_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/legacy_writer_fence.ts"; +import {shadowManagementDirectory} from "../../loopx/control_plane/coordination/shadow_management.ts"; +import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {projectCoordinationSource, SOURCE_PROJECTION_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/source_projection.ts"; +import {sourceRequest, todo, type ShadowFixture} from "./shadow_file_fixture.ts"; +import {acquireFileMutationLock, releaseFileMutationLock} from "../../loopx/control_plane/effect_runtime_io.ts"; + +/** Contract fixture for the trusted tar IO adapter. Real archive parsing and + * installed CLI coverage live in the Python production-entrypoint tests. */ +async function backupWitness(f: ShadowFixture, source: JsonObject): Promise { + const archive = join(f.root, "saved.tar.gz"), manifest = join(f.root, "saved.manifest.json"); + await writeFile(archive, "Saved archive byte identity"); await writeFile(manifest, "Saved manifest byte identity"); + const hash = (bytes: Buffer) => createHash("sha256").update(bytes).digest("hex"); + const snapshot = source.source_snapshot as JsonObject; + const files = [f.statePath, String((snapshot.registry_source as JsonObject).path), + ...(snapshot.lease_inventory as JsonObject[]).map(item => join(f.root, "goals", "goal-a", "task-leases", String(item.name)))]; + const members = await Promise.all(files.map(async path => ({archive_path: `runtime-root/${relative(f.root, path).replaceAll("\\", "/")}`, + sha256: hash(await readFile(path))}))); + return {schema_version: "loopx_state_backup_source_witness_v0", archive_path: archive, archive_sha256: hash(await readFile(archive)), + manifest_path: manifest, manifest_sha256: hash(await readFile(manifest)), + included: [{source_path: f.root, archive_path: "runtime-root"}], members}; +} + +async function fixture(t: TestContext, provider: "file" | "sqlite" = "file") { + const root = await mkdtemp(join(tmpdir(), "loopx-cold-import-")); + t.after(() => rm(root, {recursive: true, force: true})); + const statePath = join(root, "state.md"); + await writeFile(statePath, "---\ngoal_id: goal-a\nhandoff_mode: legacy\n---\n\n## Agent Todo\n\n"); + const archived = {...todo("archived", "done"), archive_state: "archive", source_section: "Agent Todo Archive", + text: "Retained complete archive body", evidence: "Original independent evidence"}; + const active = {...todo("active"), claimed_by: "agent-a", note: "Complete source metadata", + required_write_scopes: ["src/**"], successor_todo_ids: ["archived"]}; + const projection = projectCoordinationSource({schema_version: SOURCE_PROJECTION_REQUEST_SCHEMA, kind: "snapshot", + goal_id: "goal-a", handoff_mode: "legacy", todos: [active, archived], leases: [], + read_model_schema: "loopx_todo_canonical_read_record_v0"}).projection as JsonObject; + const f: ShadowFixture = {root, statePath, baseline: projection, + store: new FileAuthorityStore(join(root, "authority-shadow", "file-v0"), "goal-a")}; + await selectLocalAuthorityTarget(root, "goal-a", provider, true); + const source = await sourceRequest(f, projection); + const prepare = {schema_version: COLD_SOURCE_IMPORT_REQUEST_SCHEMA, action: "prepare", ...source, + operation_id: "cold-import:original", target_handoff_mode: "soft_claim", target_provider: provider, + source_backup: await backupWitness(f, source)}; + const operation = (action: "apply" | "recover", digest: unknown, stopped: unknown = true) => ({ + schema_version: COLD_SOURCE_IMPORT_REQUEST_SCHEMA, action, runtime_root: root, goal_id: "goal-a", + operation_id: prepare.operation_id, expected_plan_sha256: digest, + ...(action === "apply" ? {writers_stopped: stopped} : {})}); + const carrier = join(shadowManagementDirectory(root, "goal-a"), "cold-imports", + `${canonicalAuthoritySha256(prepare.operation_id)}.json`); + return {f, source, prepare, operation, carrier, active, archived}; +} + +test("empty selected SQLite preview remains inspectable, without accepting a lost or committed target", async t => { + const f = await fixture(t, "sqlite"); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal(prepared.status, "prepared"); + const inspect = {...f.source, schema_version: COLD_SOURCE_INSPECTION_REQUEST_SCHEMA}; + const observed = await inspectColdCoordinationStorage(inspect); + assert.equal(observed.ok, true, JSON.stringify(observed)); + assert.equal((observed.current as JsonObject).canonical, false); + assert.equal(observed.execution_authority_granted, false); + const opened = await openLocalAuthorityStoreHandle(f.f.root, "goal-a"); + assert.equal((await opened.store.loadAuthority()).status, "missing"); + const sqlitePath = localAuthorityProviderPaths(f.f.root, "goal-a").sqlite; + await rm(sqlitePath, {recursive: true}); + const lost = await inspectColdCoordinationStorage(inspect); + assert.equal(lost.ok, false); + assert.equal(lost.current, null); + // Missing identity must not be recreated by a read. + await assert.rejects(() => openLocalAuthorityStoreHandle(f.f.root, "goal-a")); +}); + +test("unfenced committed SQLite and malformed selectors still refuse old-source inspection", async t => { + const f = await fixture(t, "sqlite"); + const inspect = {...f.source, schema_version: COLD_SOURCE_INSPECTION_REQUEST_SCHEMA}; + const opened = await openLocalAuthorityStoreHandle(f.f.root, "goal-a"); + const committed = await opened.store.commitAuthority({expected_provider_revision: null, + operation_id: "unfenced-existing-head", next_projection: f.f.baseline, events: [], receipts: []}); + assert.equal(committed.status, "applied"); + assert.equal((await inspectColdCoordinationStorage(inspect)).reason_code, "cold_source_canonical_authority_present"); + await writeFile(localAuthorityProviderPaths(f.f.root, "goal-a").marker, "{unreadable original selector"); + assert.equal((await inspectColdCoordinationStorage(inspect)).ok, false); +}); + +for (const provider of ["file", "sqlite"] as const) { + test(`${provider}: cold cutover retains records and original receipt after later writes`, async t => { + const f = await fixture(t, provider); + const before = await readFile(f.f.statePath); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal(prepared.status, "prepared", JSON.stringify(prepared)); + assert.equal(prepared.complete_goal_backup_verified, false); + const originalRead = await executeColdSourceImport({...f.operation("recover", prepared.plan_sha256), action: "readback"}); + assert.equal(originalRead.status, "prepared"); + assert.equal(originalRead.authority_changed, false); + assert.equal(originalRead.legacy_writer_fenced, false); + assert.equal((await f.f.store.loadAuthority()).status, "missing"); + assert.equal((await executeColdSourceImport(f.operation("recover", prepared.plan_sha256))).reason_code, + "cold_import_recovery_requires_fence"); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256, false))).reason_code, + "cold_import_operator_stop_confirmation_required"); + const result = await executeColdSourceImport(f.operation("apply", prepared.plan_sha256)); + assert.equal(result.status, "applied", JSON.stringify(result)); + assert.equal(result.stop_confirmation_source, "operator_attestation"); + assert.equal(result.execution_authority_granted, false); + const opened = await openLocalAuthorityStoreHandle(f.f.root, "goal-a"); + const initial = await opened.store.loadAuthority(); + assert.equal(initial.status, "loaded"); + if (initial.status !== "loaded") throw new Error("import head absent"); + assert.equal(initial.cursor, "1"); + assert.equal(initial.head.handoff_mode, "soft_claim"); + assert.deepEqual(initial.head.todos, [f.active, f.archived]); + assert.deepEqual(initial.head.leases, []); + assert.deepEqual(await readFile(f.f.statePath), before); + const blocked = await checkLegacyCoordinationWriteAllowed({schema_version: LEGACY_COORDINATION_WRITE_CHECK_REQUEST_SCHEMA, + runtime_root: f.f.root, goal_id: "goal-a"}); + assert.equal(blocked.status, "blocked"); + const later = {...initial.head, canonical_followup: "Preserve this new write"}; + assert.equal((await opened.store.commitAuthority({operation_id: "canonical:later", + expected_provider_revision: initial.provider_revision, events: [], receipts: [], next_projection: later})).status, "applied"); + await rm(f.f.statePath); + // Restart in another process, with no source snapshot or source file. + const module = new URL("../../loopx/control_plane/coordination/cold_source_import.ts", import.meta.url).href; + const recovered = spawnSync(process.execPath, ["--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + `import {executeColdSourceImport} from ${JSON.stringify(module)}; process.stdout.write(JSON.stringify(await executeColdSourceImport(${JSON.stringify(f.operation("recover", prepared.plan_sha256))})));`], {encoding: "utf8"}); + assert.equal(recovered.status, 0, recovered.stderr); + assert.equal(JSON.parse(recovered.stdout).status, "replayed", recovered.stdout); + const current = await opened.store.loadAuthority(); + assert.equal(current.status, "loaded"); + if (current.status !== "loaded") throw new Error("later head absent"); + assert.deepEqual(current.head, later); + assert.equal(current.cursor, "2"); + assert.equal((await opened.store.readReceipt("cold-import:original")).status, "found"); + const readback = await executeColdSourceImport({...f.operation("recover", prepared.plan_sha256), action: "readback"}); + assert.equal(readback.status, "replayed"); + assert.equal(readback.authority_changed, false); + assert.deepEqual((await opened.store.loadAuthority()), current); + }); +} + +test("changed source or reviewed digest never engages the writer fence", async t => { + const f = await fixture(t); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal(prepared.status, "prepared"); + assert.equal((await executeColdSourceImport(f.operation("apply", "0".repeat(64)))).reason_code, "cold_import_reviewed_plan_changed"); + await writeFile(f.f.statePath, "Unreviewed source change\n"); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).reason_code, "source_changed_retry"); + await assert.rejects(readFile(legacyCoordinationWriterFencePath(f.f.root, "goal-a")), {code: "ENOENT"}); +}); + +test("a reviewed backup is required and cannot change before the first import", async t => { + const f = await fixture(t); + const backup = f.prepare.source_backup; + const omitted = {...f.prepare}; delete (omitted as JsonObject).source_backup; + assert.equal((await executeColdSourceImport(omitted)).reason_code, "cold_import_request_invalid"); + const missingSource = {...backup, members: []}; + assert.equal((await executeColdSourceImport({...f.prepare, source_backup: missingSource})).reason_code, + "cold_import_backup_source_missing_or_changed"); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal(prepared.status, "prepared"); + await writeFile(String(backup.archive_path), "Changed archive after review"); + const refused = await executeColdSourceImport(f.operation("apply", prepared.plan_sha256)); + assert.equal(refused.reason_code, "cold_import_backup_changed"); + assert.equal(refused.legacy_writer_fenced, false); + await assert.rejects(readFile(legacyCoordinationWriterFencePath(f.f.root, "goal-a")), {code: "ENOENT"}); +}); + +test("a committed original receipt remains readable after backup loss", async t => { + const f = await fixture(t); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).status, "applied"); + await rm(String(f.prepare.source_backup.archive_path)); + const recovered = await executeColdSourceImport(f.operation("recover", prepared.plan_sha256)); + assert.equal(recovered.status, "replayed"); + assert.equal(recovered.executed, false); + assert.equal(recovered.cursor, "1"); + assert.equal(recovered.complete_goal_backup_verified, false); +}); + +for (const variant of ["expired_active", "orphan_active", "symlink", "unsupported_filename"] as const) { + test(`cold import rejects ${variant} lease facts rather than deleting them`, async t => { + const f = await fixture(t); + const directory = join(f.f.root, "goals", "goal-a", "task-leases"); + await mkdir(directory, {recursive: true}); + const id = variant === "expired_active" ? "active" : "orphan"; + const record = {schema_version: "task_lease_v0", goal_id: "goal-a", todo_id: id, owner: "agent-a", + idempotency_key: "old-execution", status: "active", version: 3, lease_epoch: 2, + expires_at: "2000-01-01T00:00:00Z"}; + if (variant === "symlink") { + const outside = join(f.f.root, "original.json"); + await writeFile(outside, JSON.stringify(record)); + await symlink(outside, join(directory, "orphan.json")); + } else await writeFile(join(directory, variant === "unsupported_filename" ? "unsafe name.json" : `${id}.json`), JSON.stringify(record)); + const result = await executeColdSourceImport(f.prepare); + assert.equal(result.reason_code, variant === "symlink" ? "cold_import_lease_source_unsafe" + : variant === "unsupported_filename" ? "cold_import_lease_source_unsupported" : "cold_import_lease_requires_settlement", JSON.stringify(result)); + await assert.rejects(readFile(f.carrier), {code: "ENOENT"}); + await assert.rejects(readFile(legacyCoordinationWriterFencePath(f.f.root, "goal-a")), {code: "ENOENT"}); + }); +} + +test("unresolved original outbox is refused without its producer or shadow history", async t => { + const f = await fixture(t); + const directory = join(f.f.root, "authority-shadow", "outbox", "goal-a", "todos"); + await mkdir(directory, {recursive: true}); + const original = join(directory, "0000000001-original.prepared.json"); + await writeFile(original, "Original retained prepared bytes\n"); + assert.equal((await executeColdSourceImport(f.prepare)).reason_code, "cold_import_outbox_requires_disposition"); + assert.equal(await readFile(original, "utf8"), "Original retained prepared bytes\n"); +}); + +test("corrupt carrier and completed authority loss are fail-closed", async t => { + const f = await fixture(t); + const prepared = await executeColdSourceImport(f.prepare); + const original = await readFile(f.carrier, "utf8"); + const corrupted = JSON.parse(original); + corrupted.plan.target_projection.todos = []; + await writeFile(f.carrier, JSON.stringify(corrupted)); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).reason_code, "cold_import_reviewed_plan_changed"); + // Recomputing a self-checksum does not qualify a changed target projection. + corrupted.plan_sha256 = canonicalAuthoritySha256(corrupted.plan); + await writeFile(f.carrier, JSON.stringify(corrupted)); + assert.equal((await executeColdSourceImport(f.operation("apply", corrupted.plan_sha256))).reason_code, "cold_import_carrier_invalid"); + await writeFile(f.carrier, original); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).status, "applied"); + const opened = await openLocalAuthorityStoreHandle(f.f.root, "goal-a"); + assert.ok(opened.store instanceof FileAuthorityStore); + await rm((opened.store as FileAuthorityStore).path); + assert.equal((await executeColdSourceImport({...f.operation("recover", prepared.plan_sha256), action: "readback"})).reason_code, + "cold_import_completed_authority_missing"); + const result = await executeColdSourceImport(f.operation("recover", prepared.plan_sha256)); + assert.equal(result.reason_code, "cold_import_completed_authority_missing", JSON.stringify(result)); + await assert.rejects(readFile((opened.store as FileAuthorityStore).path), {code: "ENOENT"}); +}); + +test("killed fenced process recovers the original operation without reading Markdown", async t => { + const f = await fixture(t); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal(prepared.status, "prepared"); + const opened = await openLocalAuthorityStoreHandle(f.f.root, "goal-a"); + const store = opened.store as FileAuthorityStore; + const lock = await acquireFileMutationLock(store.path); + const module = new URL("../../loopx/control_plane/coordination/cold_source_import.ts", import.meta.url).href; + const child = spawn(process.execPath, ["--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + `import {executeColdSourceImport} from ${JSON.stringify(module)}; process.stdout.write(JSON.stringify(await executeColdSourceImport(${JSON.stringify(f.operation("apply", prepared.plan_sha256))})));`], {stdio: "pipe"}); + let output = ""; + child.stdout.on("data", bytes => { output += bytes; }); + const ended = once(child, "exit"); + try { + const deadline = Date.now() + 4000; + let fenced = false; + while (Date.now() < deadline) { + try { await readFile(legacyCoordinationWriterFencePath(f.f.root, "goal-a")); fenced = true; break; } + catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + await delay(20); + } + assert.equal(fenced, true, "child must reach durable fence before the injected crash"); + assert.equal((await store.loadAuthority()).status, "missing"); + child.kill("SIGKILL"); + assert.deepEqual(await ended, [null, "SIGKILL"]); + assert.equal(output, ""); + } finally { + child.kill("SIGKILL"); + await releaseFileMutationLock(store.path, lock.token); + } + await rm(f.f.statePath); + const observed = await executeColdSourceImport({...f.operation("recover", prepared.plan_sha256), action: "readback"}); + assert.equal(observed.status, "prepared"); + assert.equal(observed.legacy_writer_fenced, true); + assert.equal(observed.authority_changed, false); + assert.equal((await store.loadAuthority()).status, "missing"); + await assert.rejects(readFile(`${f.carrier}.completed.json`), {code: "ENOENT"}); + const recovered = await executeColdSourceImport(f.operation("recover", prepared.plan_sha256)); + assert.equal(recovered.status, "applied", JSON.stringify(recovered)); + assert.equal((await executeColdSourceImport(f.operation("recover", prepared.plan_sha256))).status, "replayed"); + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status !== "loaded") throw new Error("cold recovery head missing"); + assert.equal(head.cursor, "1"); + assert.deepEqual(head.head.todos, [f.active, f.archived]); +}); + +test("readback rejects a changed completion marker without repairing it", async t => { + const f = await fixture(t); + const prepared = await executeColdSourceImport(f.prepare); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).status, "applied"); + const path = `${f.carrier}.completed.json`; + await writeFile(path, JSON.stringify({operation_id: "another-operation"})); + const damaged = await readFile(path); + assert.equal((await executeColdSourceImport({...f.operation("recover", prepared.plan_sha256), action: "readback"})).reason_code, + "cold_import_completion_identity_changed"); + assert.deepEqual(await readFile(path), damaged); +}); + +test("released orphan lease bytes stay in the source witness without entering the execution graph", async t => { + const f = await fixture(t); + const directory = join(f.f.root, "goals", "goal-a", "task-leases"); + await mkdir(directory, {recursive: true}); + const old = JSON.stringify({schema_version: "task_lease_v0", goal_id: "goal-a", todo_id: "orphan", owner: "agent-a", + idempotency_key: "settled-original", status: "released", version: 9, lease_epoch: 8}); + await writeFile(join(directory, "orphan.json"), old); + const source = await sourceRequest(f.f, f.f.baseline); + const prepared = await executeColdSourceImport({...f.prepare, ...source, source_backup: await backupWitness(f.f, source)}); + assert.equal(prepared.status, "prepared", JSON.stringify(prepared)); + const carrier = JSON.parse(await readFile(f.carrier, "utf8")); + const witness = carrier.plan.source_witness.find((item: JsonObject) => item.path === join(directory, "orphan.json")); + assert.equal(Buffer.from(witness.bytes_base64, "base64").toString("utf8"), old); + assert.equal((await executeColdSourceImport(f.operation("apply", prepared.plan_sha256))).status, "applied"); + const store = (await openLocalAuthorityStoreHandle(f.f.root, "goal-a")).store; + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status !== "loaded") throw new Error("head missing"); + assert.deepEqual(head.head.leases, []); + assert.equal(await readFile(join(directory, "orphan.json"), "utf8"), old); +}); diff --git a/tests/control_plane_ts/content_digest_single_owner.test.ts b/tests/control_plane_ts/content_digest_single_owner.test.ts index eb70ca8cd7..dbd7d0c60e 100644 --- a/tests/control_plane_ts/content_digest_single_owner.test.ts +++ b/tests/control_plane_ts/content_digest_single_owner.test.ts @@ -100,6 +100,7 @@ const CANONICAL_CONSUMERS = [ "control_plane/collaboration/semantic_request.ts", "control_plane/coordination/authority_archive_read.ts", "control_plane/coordination/authority_source.ts", + "control_plane/coordination/cold_source_import.ts", "control_plane/coordination/legacy_writer_fence.ts", "control_plane/coordination/local_authority_migration.ts", "control_plane/coordination/local_authority_shadow.ts",