diff --git a/apps/presentation/dashboard/src/data/goal-storage.ts b/apps/presentation/dashboard/src/data/goal-storage.ts index c0aab968e5..90dc342664 100644 --- a/apps/presentation/dashboard/src/data/goal-storage.ts +++ b/apps/presentation/dashboard/src/data/goal-storage.ts @@ -23,6 +23,12 @@ export const storageResultSchema = z.object({ target_provider: provider.optional(), selected_provider: provider.optional(), reviewed_source: z.object({provider, cursor: z.string(), provider_revision: z.string(), store_identity: z.string()}).optional(), current: storageSourceSchema.nullable().optional(), + cold_source: z.object({ + active_todo_count: z.number().int().nonnegative(), archived_todo_count: z.number().int().nonnegative(), + lease_file_count: z.number().int().nonnegative(), unsettled_lease_count: z.number().int().nonnegative(), + capture_artifacts_present: z.boolean(), outbox_files_present: z.boolean(), + import_ready: z.literal(false), writer_stop_verified: z.literal(false), outbox_reconciliation_verified: z.literal(false), + }).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(), }); 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 64c87ae0a6..a6383c6111 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 @@ -9,6 +9,7 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan const {t} = useWorkspaceI18n(); const key = `loopx-storage-preview:${goalId}`; const [current, setCurrent] = useState(null); + const [cold, setCold] = useState(); const [carrier, setCarrier] = useState(null); const [result, setResult] = useState(null); const [target, setTarget] = useState("sqlite"); @@ -21,7 +22,7 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan const generation = useRef(0); useEffect(() => { const token = ++generation.current; - setCurrent(null); setResult(null); setCarrier(null); setConfirmed(false); setInvalidSaved(false); setError(null); setBusy(true); + setCurrent(null); setCold(undefined); setResult(null); setCarrier(null); setConfirmed(false); setInvalidSaved(false); setError(null); setBusy(true); let saved: StorageCarrier | null = null; try { const raw = localStorage.getItem(key); @@ -33,10 +34,11 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan try { const observed = await fetchGoalStorage(goalId); if (token !== generation.current) return; + setCold(observed.cold_source); if (observed.ok && observed.current) { setCurrent(observed.current); setTarget(observed.current.provider === "sqlite" ? "file" : "sqlite"); - } else setError(t("storage.unavailable")); + } else { setResult(observed); setError(t(saved ? "storage.unavailable" : "storage.readUnavailable")); } if (saved) { const recovered = await recoverGoalStorage(saved); if (token !== generation.current) return; @@ -45,7 +47,7 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan if (!recovered.ok) setError(t("storage.rejected")); else if (!recovered.current) setError(t("storage.unavailable")); } - } catch { if (token === generation.current) setError(t("storage.unavailable")); } + } catch { if (token === generation.current) setError(t(saved ? "storage.unavailable" : "storage.readUnavailable")); } finally { if (token === generation.current) setBusy(false); } })(); return () => { generation.current++; }; @@ -89,10 +91,17 @@ export function GoalStorageSettings({goalId, onChanged}: {goalId: string; onChan return

{t("storage.title")}

-

{t("storage.boundary")}

+ {current?.canonical || carrier ?

{t("storage.boundary")}

: null} {current ?
{t("storage.current")} - {current.provider ?? t("ownership.unpromoted")} - {current.canonical ? {t("storage.counts", {todos: current.todo_count ?? 0, leases: current.unsettled_lease_count ?? 0})} :

{t("storage.promoteFirst")}

} + {current.provider ?? t("storage.oldSource")} + {current.canonical ? {t("storage.counts", {todos: current.todo_count ?? 0, leases: current.unsettled_lease_count ?? 0})} : <> + {cold ? <> + {t("storage.coldCounts", {active: cold.active_todo_count, archived: cold.archived_todo_count, leases: cold.unsettled_lease_count})} + {cold.capture_artifacts_present ? {t("storage.coldCapture")} : null} + {cold.outbox_files_present ? {t("storage.coldOutbox")} : null} + : null} +

{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 ?? "?"})}

diff --git a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx index be34e27e7a..e730bacc45 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx @@ -17,7 +17,10 @@ const en = { "storage.current": "Current source (fresh readback)", "storage.counts": "{todos} tasks · {leases} unsettled leases", "storage.target": "Target storage", - "storage.promoteFirst": "Previous Markdown state: use the backed-up CLI promotion first. This page switches existing canonical File/SQLite stores.", + "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.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", @@ -38,7 +41,9 @@ const en = { "ownership.soft_claim": "Collaborative claims", "ownership.hard_lease": "Exclusive execution leases", "ownership.unpromoted": "Not on canonical storage", - "ownership.promoteFirst": "Promote this Goal through the backed-up CLI migration first. This policy form does not migrate 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.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.", @@ -1360,7 +1365,10 @@ const zhCN: Record = { "storage.current": "当前来源(实时读回)", "storage.counts": "{todos} 项任务 · {leases} 项未结算 lease", "storage.target": "目标存储", - "storage.promoteFirst": "旧 Markdown 状态:请先通过有备份的 CLI 晋升。这里切换已存在的 canonical File/SQLite 存储。", + "storage.coldCounts": "{active} 项当前任务 · {archived} 项归档任务 · {leases} 项未结算 lease", + "storage.coldCapture": "已发现原 capture 文件,历史原样保留。", + "storage.coldOutbox": "已发现 outbox 文件,处理结果尚未验证。", + "storage.coldBoundary": "已盘点旧 Markdown 源。这里尚不能导入:仍需验证 writer/Host 停止、lease 结算、outbox 处置与备份绑定的导入。此次未生成 capture、迁移数据或授予执行权限。", "storage.reviewed": "已审核来源:{source},游标 {cursor} → {target}。来源变化时拒绝应用。", "storage.confirm": "我已停止写入方并结算 lease,确认此预览的存储变更。", "storage.apply": "备份并切换存储", @@ -1381,7 +1389,9 @@ const zhCN: Record = { "ownership.soft_claim": "协作认领", "ownership.hard_lease": "独占执行租约", "ownership.unpromoted": "尚未使用统一状态存储", - "ownership.promoteFirst": "请先通过带备份的 CLI 迁移将 Goal 晋升到统一存储。此策略表单不迁移存储。", + "storage.oldSource": "旧 Markdown 来源", + "storage.readUnavailable": "未能读取当前存储,请重试读回。", + "ownership.promoteFirst": "此 Goal 尚无 canonical 存储。变更所有权需要先完成经过核对的存储导入;此策略表单不执行导入。", "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 959daf5fc5..c2fa6e9252 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -881,6 +881,38 @@ L3 checkpoint: standalone acquisition/takeover, atomic claim admission and maint - **Exit:** use shared-authority Section 7.2's separate decisions for a bounded change, reversible opt-in cohort and released default. Each requires affected real CLI/backend and independent baseline/negative/recovery evidence at its own scope. Formal D2 retains applicable volume and at least ten-day evidence; a cohort need not wait for that certificate. D3 retains explicit cutover authority. This plan runs no soak or provider promotion. - **Rollback:** reviewed fenced export/import and schema-aware downgrade; replacing a binary cannot restore old write authority. +Cold-source inventory is now an explicit read-only CLI prerequisite through the +same TS source/lease owners. It includes unreferenced archives and retained +leases, with full text and source-byte witnesses, before any shadow opt-in. +Original outbox disposition is also qualified through the shipped TS effects +with Python Todo/bootstrap/capture producers absent from a disposable receiver: +markerless abandoned/committed recovery, receipt replay without duplicate effects, +unchanged active leases, and refusal plus same-operation archival of ambiguous +originals. Retain the OS-lock adapter and original history readers. This covers +the existing disposition owner, not global Host stop, lease settlement, import +confirmation, a canonical cutover, or permission to delete active old writers. +The same observation discovers original capture stores/identity, management +operation files, outbox bytes and Goal-bound rollback archives through the +existing typed owners; present history is validated without replay or drain. +Compact readback never replaces the witnessed historical files or proves Host +stop. Invalid history and missing completed rollback archives refuse inspection. +Active original outbox inspection reuses the native drain proof owner to show +pending records and receipt-proven residue per partition without any effects. +An unavailable disposition retains raw witnesses; inactive/interrupted capture +still needs its original management recovery. This preview never settles the +outbox, updates a cursor or grants cleanup/import authority. +Expired or orphan `active` leases still require settlement; canonical selectors +and fences prevent treating display Markdown as an old authority source. +This observation creates no capture or import receipt and never reports import +readiness. Goal settings now consume that same observation: current/archive +task counts, unsettled historical leases and retained capture/outbox presence, +with path-free readback, unavailable-source refusal and fresh retry. This is the +inventory stage only; no import, capture activation or execution grant is exposed. +R5/D1 and T4/C1 still need stopped writer/Host proof, original outbox disposition, +the backup-bound reviewed target and its Goal storage frontend, +confirmation and same-operation import/recovery. Keep the cold-import work open; +neither this prerequisite nor archive export qualifies writer cutoff or defaults. + Cold-source retention checkpoint: current-project full backup discovers its registered custom state and source-registry routes. Real CLI backups and inert independent extraction preserve complete Markdown bytes, unreferenced archived diff --git a/docs/reference/reviewed-coordination-promotion.md b/docs/reference/reviewed-coordination-promotion.md index ded6c9ebf6..d57b00f082 100644 --- a/docs/reference/reviewed-coordination-promotion.md +++ b/docs/reference/reviewed-coordination-promotion.md @@ -12,6 +12,87 @@ fencing and receipt proof; Python only loads the file and transports the request ## Preview and execute +### Inspect a cold old source + +An unpromoted Markdown Goal can be inventoried before enabling capture: + +```bash +loopx --format json coordination-shadow inspect-source \ + --goal-id example-goal > old-source-inventory.json +``` + +This read-only command includes all supported active and archived Todo records, +not just archives needed by the current dependency graph. It preserves full +archived text and supported metadata through the existing record codec. The +source witness binds the registered Goal, state bytes, registry and every lease +file; TS revalidates it under the existing source/maintenance locks. Retained +leases are returned separately from current graph edges. Every `active` lease +requires settlement, including expired leases and leases for removed Todos: +expiry is not proof that the Host stopped. Missing archive roles, duplicate +identities, invalid historical leases and stale sources reject the inspection. +An existing canonical selector/document or writer fence rejects using Markdown +as a cold import source, including when the selected provider is unavailable. + +`source_inventory.capture` also inventories the original management operations, +active outbox, runtime store and identity, this Goal's retained rollback stores, +and legacy observation directory. File witnesses use raw byte hashes; malformed +outbox files remain visible without being drained or rewritten. Existing TS +readers validate present active history. Their compact readback is not a copy +of all receipts: preserve the original witnessed files. Missing or altered +completed rollback archives, invalid history and unsafe file layouts refuse +inspection. Interrupted management remains an unfinished original operation. + +For an active original capture, `capture.outbox_review` uses the existing native +drain verifier to distinguish pending entries from residue with exact original +receipts, independently for Todos and leases. Its partition plans include the +original entry identities and proposed cursor/reclamation readback. `planned` +is a read-only preview: `executed` and `execution_authority_granted` are false. +Only a fresh drain through the original capture owner can act on those facts; +the inventory cannot authorize deletion, replay or a new import receipt. +Malformed files, foreign lineage, unproved markers, changed receipt bytes or +unanchored cursors make this review `failed` with the owning reason code; raw +file witnesses remain available and unchanged. An inactive or interrupted +capture returns no outbox review: recover its original management operation +first. Neither case establishes outbox reconciliation or Host stop. + +Original disposition does not need the old Python capture producer. The +shipped TS drain can prove a markerless write from the locked original source: +unchanged previous bytes settle an abandoned no-op; exact new bytes prove the +commit, only when those byte versions differ and later entries are excluded. +Already committed candidate receipts replay without a second effect. +The OS-lock Host adapter remains necessary. A lease receipt neither releases +that lease nor grants work. An ambiguous source, including A→B→A with later +entries, remains unproved and byte-preserved. Revision-bound rollback can retain +the complete candidate and outbox in their original management archive; retry +reads that same operation. Archiving is not proof of settlement or import +readiness. These paths are covered with the old producers physically absent +from a disposable receiver, not a declaration that they can all be deleted yet. + +In the App, open **Goal settings → Task ownership → Goal data storage** to read +the same verified inventory. It shows task/archive and unsettled-lease counts +and retained capture/outbox presence, without exposing source text, local paths +or execution keys. **Read current storage** performs a fresh observation; a +failed read clears old counts. Import is explicitly unavailable at this stage. +There is no confirmation or migration control for an old Markdown source. + +Read `source_inventory.active_todo_count`, `archived_todo_count`, `capture`, +`retained_leases` and `leases_requiring_settlement`. `import_ready`, +`writer_stop_verified` and `outbox_reconciliation_verified` remain **false**. +There is no `--execute` switch, writer fence, bootstrap, provider initialization, +lease grant or capture receipt. This is a shared CLI/App inventory prerequisite, +not the complete import journey or a Lark operation. The Goal storage owner +still needs explicit Host/writer stop, original +outbox reconciliation, backup, reviewed target, confirmation and same-operation +import recovery. Existing shadow qualification below retains its original gates. + +Keep the JSON private: it contains full Todo text, identities and local source +paths. It is an observation, not a complete historical backup. Preserve the +original source and supported receipts through [configuration backup](configuration-backup.md). +No capability setting is enabled; stopping this inspection needs no rollback. +Delete the operator-owned JSON when it is no longer needed. + +### Qualified shadow promotion + Use an explicitly enabled, bootstrapped and qualified runtime shadow. Its qualification must cover real mutations and required event classes; an empty shadow or a saved JSON file cannot substitute for that evidence. Fresh CLI diff --git a/docs/reference/reviewed-coordination-promotion.zh-CN.md b/docs/reference/reviewed-coordination-promotion.zh-CN.md index 4b0bb365d3..e2dc7d4fe4 100644 --- a/docs/reference/reviewed-coordination-promotion.zh-CN.md +++ b/docs/reference/reviewed-coordination-promotion.zh-CN.md @@ -8,6 +8,65 @@ fence 与 receipt 证明由 TypeScript 协调边界负责;Python 只读文件 ## 操作 +### 盘点冷旧源 + +未晋升的 Markdown Goal 可以先盘点,不需要启用 shadow 或让旧 writer 制造捕获历史: + +```bash +loopx --format json coordination-shadow inspect-source \ + --goal-id example-goal > old-source-inventory.json +``` + +读取 `source_inventory` 中的 active/archive Todo 数量、完整支持格式的记录、 +`retained_leases` 和 `leases_requiring_settlement`。全部归档均进入这次盘点, +不只包含当前图需要的归档依赖;长正文和已支持 metadata 不按注意力摘要截断。 +原源字节、registry 和全部 lease 文件绑定到同一次来源见证,TS 在既有锁下复核。 +历史 lease 单独返回,不变成当前执行 grant。仍标记 `active` 的租约均需结算, +即使已过期或所属 Todo 已删除;过期不能证明 Host 已停止。 +缺少归档角色、重复身份、非法历史 lease 和源变化会拒绝;已有 canonical selector、 +文档或 writer fence 时拒绝把 Markdown 投影当作冷旧导入源,provider 不可用也不回退。 + +`source_inventory.capture` 同时盘点原 management 操作、活跃 outbox、runtime store +与身份、此 Goal 保留的 rollback store,以及旧 observation 目录。文件指纹绑定原始 +字节;损坏的 outbox 文件仍可见,不 drain 或重写。现有 TS reader 校验存在的活跃 +历史,紧凑读回不包含全部原回执,必须保留所指原文件。已完成 rollback 的归档缺失 +或变化、历史非法及不安全文件布局均拒绝;中断的 management 仍是未完成的原操作。 + +对原活跃 capture,`capture.outbox_review` 复用现有原生 drain 的证明规则,分别列出 +Todo 和 lease 分区中待处理的原记录及有精确原回执的残留。分区计划带原 entry 身份与 +拟议 cursor/回收读回;`planned` 仅为只读预览,`executed` 和 +`execution_authority_granted` 均为 false。只有原 capture owner 下的新一次 drain 才能 +据当前事实执行,盘点不授权删除、replay 或制造导入回执。损坏文件、外来 lineage、 +无证明 marker、回执字节变化及无锚 cursor 使该预览 `failed`,保留原 reason code 和 +未改动的原文件见证。inactive/中断 capture 不生成该预览,应先恢复原 management +操作;这些结果均不证明 outbox 已处置或 Host 已停。 + +原记录处置不需要重启旧 Python capture 生产者。现有 TS drain 在原来源锁下核验无 +marker 的记录:新旧字节不同且排除后续记录时,仍为旧字节则结算为未发生的 no-op, +精确新字节则证明提交;已进入 +candidate 的原回执只 replay,不产生第二次效果。OS-lock Host 适配器仍有必要;lease +回执既不释放 lease,也不授予工作。来源有歧义(包括后续记录存在时的 A→B→A)则 +拒绝并保留原字节。绑定精确 revision 的 rollback 可在原 management 归档中完整 +保留 candidate/outbox,重试读回同一操作;归档不证明结算或导入就绪。独立接收端中 +物理移除旧生产者后的测试覆盖了这些路径,尚不代表可以删除全部旧 writer。 + +App 中打开 **Goal 设置 → 任务所有权 → Goal 数据存储**,读取同一 TS owner 核验的 +盘点。页面只显示当前/归档任务及未结算 lease 数量、原 capture/outbox 文件是否存在, +不暴露源正文、本机路径或执行密钥。“读回当前存储”重新观察,失败时清除旧数量。 +此阶段明确显示导入尚不可用,不为旧 Markdown 源提供确认或迁移按钮。 + +`import_ready`、`writer_stop_verified`、`outbox_reconciliation_verified` 始终为 false。 +此命令没有 `--execute`,不启用配置、创建 shadow/provider/fence、授予 lease 或制造捕获回执。 +它交付共享 CLI/App 盘点前置;Goal storage owner 仍须完成明确停 +writer/Host、原 outbox 对账、备份、审核目标、确认和同操作导入恢复。下方 shadow +晋升仍遵守原资格,不因盘点通过而放宽。 + +JSON 含完整正文、身份与本机路径,应留在操作员私有存储。它是观察结果,不是完整 +历史备份;原始来源和受支持回执仍按[配置备份](configuration-backup.md)保留。 +没有能力设置被启用,停止盘点无需回退;不再需要时删除本人的输出文件即可。 + +### 合格 shadow 的晋升 + 先显式启用并 bootstrap runtime shadow,让它捕获真实变更并通过资格校验。 新 CLI 预览默认 `preserve`,只切换存储权威、保留当前所有权策略;显式 `--handoff-mode-migration hard_lease` 才审核策略升级及现有 claim/lease。 diff --git a/examples/personal-workspace-browser/goal-storage.mjs b/examples/personal-workspace-browser/goal-storage.mjs new file mode 100644 index 0000000000..1df2039d38 --- /dev/null +++ b/examples/personal-workspace-browser/goal-storage.mjs @@ -0,0 +1,140 @@ +// Packaged Goal settings -> real HTTP -> typed cold-source owner. +// Workspace discovery is synthetic; every storage response comes from the +// production handler and disposable source, never a browser response fixture. +import assert from "node:assert/strict"; +import {spawn} from "node:child_process"; +import {createInterface} from "node:readline"; +import {mkdtemp, mkdir, readFile, rm, writeFile} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {resolve} from "node:path"; +import {resolveTestPython} from "../../scripts/test-python.mjs"; +import {launchBrowser, loadPlaywright, waitForHttp} from "../dashboard-browser-smoke-support.mjs"; +import {openWorkspacePage} from "./scenario-context.mjs"; +import {outputDir, packaged, port, repoRoot, startServer} from "./fixture.mjs"; + +async function startAuthority() { + 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" + + " \n" + + "## Todo Archive\n- [x] Private full archived requirement\n" + + " \n"; + await writeFile(resolve(root, "state.md"), source); + await writeFile(resolve(root, "registry.json"), JSON.stringify({ + common_runtime_root: resolve(root, "runtime"), goals: [{id: "multi-agent-projection", + repo: root, state_file: "state.md", coordination: {registered_agents: ["agent-a"]}}], + })); + 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", + 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: []})); + const child = spawn(resolveTestPython(), ["-I", "-u", "-c", ` +from pathlib import Path +import sys +from loopx.chat_server import ChatHTTPServer, ChatRequestHandler +server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler) +server.registry_path = Path(sys.argv[1]) / "registry.json" +server.runtime_root = Path(sys.argv[1]) / "runtime" +server.runtime_root_override = str(server.runtime_root) +server.verbose = False +print(server.server_port, flush=True) +server.serve_forever() +`, root], {cwd: repoRoot, stdio: ["ignore", "pipe", "pipe"]}); + const lines = createInterface({input: child.stdout}); + let diagnostic = ""; + child.stderr.on("data", chunk => { diagnostic = (diagnostic + chunk).slice(-2000); }); + try { + const port = await new Promise((accept, reject) => { + const timeout = setTimeout(() => reject(new Error(`Storage authority timeout: ${diagnostic}`)), 20_000); + child.once("error", error => { clearTimeout(timeout); reject(error); }); + child.once("exit", () => { clearTimeout(timeout); reject(new Error(`Storage authority exited: ${diagnostic}`)); }); + lines.once("line", line => { clearTimeout(timeout); accept(Number(line)); }); + }); + assert.ok(Number.isSafeInteger(port) && port > 0); + return {root, source, url: `http://127.0.0.1:${port}`, async close() { + lines.close(); + if (child.exitCode === null && child.signalCode === null) { + const exited = new Promise(accept => child.once("exit", accept)); + child.kill("SIGTERM"); await exited; + } + await rm(root, {recursive: true, force: true}); + }}; + } catch (error) { + lines.close(); child.kill("SIGTERM"); await rm(root, {recursive: true, force: true}); throw error; + } +} + +const goalStorageScenario = { + id: "goal-storage", + async run({browser, url}) { + const authority = await startAuthority(); + let context; + let writes = 0; + try { + context = await openWorkspacePage(browser, url, {beforeGoto: async (_api, page) => { + await page.route(/\/api\/chat\/goal-(storage|ownership)(?:[/?]|$)/, async route => { + if (route.request().method() !== "GET") writes++; + const parsed = new URL(route.request().url()); + await route.fulfill({response: await route.fetch({url: authority.url + parsed.pathname + parsed.search})}); + }); + }}); + const {page} = context; + async function open(language = "zh") { + await page.locator(".personal-goal-link", {hasText: "Multi Agent Projection"}).click(); + await page.getByRole("button", {name: language === "zh" ? "Goal 设置" : "Goal settings", exact: true}).click(); + await page.getByRole("button", {name: language === "zh" ? "任务所有权" : "Task ownership", exact: true}).click(); + return page.getByRole("region", {name: language === "zh" ? "Goal 数据存储" : "Goal data storage"}); + } + 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.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"); + await mkdir(outbox, {recursive: true}); + const residue = resolve(outbox, "original.json"); + await writeFile(residue, "{unrecognized original bytes"); + await panel.getByRole("button", {name: "读回当前存储", exact: true}).click(); + await panel.getByText("已发现 outbox 文件,处理结果尚未验证。", {exact: true}).waitFor(); + await writeFile(resolve(authority.root, "state.md"), authority.source.replace("todo_old", "todo_current")); + await panel.getByRole("button", {name: "读回当前存储", exact: true}).click(); + await panel.getByRole("alert").waitFor(); + assert.equal(await panel.getByText(/1 项当前任务/).count(), 0, "failed read must clear stale facts"); + await writeFile(resolve(authority.root, "state.md"), authority.source); + await panel.getByRole("button", {name: "读回当前存储", exact: true}).click(); + await panel.getByText("1 项当前任务 · 1 项归档任务 · 1 项未结算 lease", {exact: true}).waitFor(); + await page.setViewportSize({width: 390, height: 844}); + await panel.getByRole("button", {name: "读回当前存储", exact: true}).focus(); + assert.equal(await page.evaluate(() => document.documentElement.scrollWidth > innerWidth + 1), false); + await page.screenshot({path: resolve(outputDir, "goal-storage-cold-mobile.png"), animations: "disabled"}); + await page.setViewportSize({width: 1512, height: 982}); + await page.evaluate(() => localStorage.setItem("loopx-pw-locale", "en")); + 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.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."}; + } finally { await context?.close(); await authority.close(); } + }, +}; + +// Standalone entry keeps this bounded acceptance out of unrelated scenario lists. +await mkdir(outputDir, {recursive: true}); +const server = await startServer(); +let browser; +try { + const url = `http://127.0.0.1:${port}/${packaged ? "chat/" : ""}?statusUrl=/status.json`; + await waitForHttp(url); + browser = await launchBrowser(loadPlaywright().chromium); + const result = await goalStorageScenario.run({browser: {newPage: options => browser.newPage({locale: "zh-CN", ...options})}, url}); + console.log(JSON.stringify({status: "PASS", ...result})); +} finally { await browser?.close(); server.kill("SIGTERM"); } diff --git a/loopx/cli_commands/coordination_shadow.py b/loopx/cli_commands/coordination_shadow.py index c7ca198aae..69fcecfd6c 100644 --- a/loopx/cli_commands/coordination_shadow.py +++ b/loopx/cli_commands/coordination_shadow.py @@ -50,6 +50,7 @@ def register_coordination_shadow_command( ) for name, help_text in ( ("inspect", "Compare the current legacy projection with the file shadow."), + ("inspect-source", "Inventory old Markdown Todos, retained leases and original capture artifacts without enabling a shadow or importing."), ( "qualify", "Validate bounded parity and transaction coverage for the active outbox lineage.", @@ -155,6 +156,18 @@ def _render(payload: dict[str, object]) -> str: configuration = payload.get("configuration") if isinstance(configuration, dict): lines.append(f"- configuration: `{configuration.get('reason_code')}`") + inventory = payload.get("source_inventory") + if isinstance(inventory, dict): + lines.extend([ + f"- source_inventory: `{inventory.get('status')}`", + f"- active_todos: `{inventory.get('active_todo_count')}`", + f"- archived_todos: `{inventory.get('archived_todo_count')}`", + f"- retained_lease_files: `{inventory.get('lease_file_count')}`", + f"- leases_requiring_settlement: `{inventory.get('leases_requiring_settlement')}`", + "- import_ready: `false` — writer/Host stop and outbox reconciliation remain unverified.", + ]) + if inventory.get("reason"): + lines.append(f"- reason: {inventory['reason']}") inspection = payload.get("inspection") if isinstance(inspection, dict): lines.extend( @@ -242,6 +255,26 @@ def handle_coordination_shadow_command( } print_payload(payload, output_format(args), _render) return 0 if payload["ok"] else 1 + if args.coordination_shadow_command == "inspect-source": + from ..control_plane.coordination.local_authority_shadow_projection import source_effect_runtime_result + + _, _, 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) + inventory = source_effect_runtime_result("coordination.source.inspect", { + "schema_version": "loopx_cold_source_inspection_request_v0", + "runtime_root": str(runtime_root.expanduser().absolute()), "goal_id": args.goal_id, + "projection": projection, "source_snapshot": snapshot, + }) + payload = {"ok": inventory.get("status") == "inspected", + "schema_version": "loopx_coordination_shadow_admin_v0", + "action": "inspect-source", "goal_id": args.goal_id, + "executed": False, "source_inventory": inventory, + "decision_read_from_shadow": False} + print_payload(payload, output_format(args), _render) + return 0 if payload["ok"] else 1 if reviewed_plan is not None and ( args.minimum_operations is not None or args.require_event_kind ): diff --git a/loopx/control_plane/coordination/cold_source_inspection.ts b/loopx/control_plane/coordination/cold_source_inspection.ts new file mode 100644 index 0000000000..156aef079e --- /dev/null +++ b/loopx/control_plane/coordination/cold_source_inspection.ts @@ -0,0 +1,165 @@ +/** Read-only complete old-source inventory. This is neither a reviewed import + * plan nor evidence of stopped Hosts, settled effects or qualified shadow. */ +import {readFile, lstat} from "node:fs/promises"; +import {dirname, join} from "node:path"; +import type {JsonObject} from "../effect_program.ts"; +import {canonicalAuthorityObject, canonicalAuthoritySha256} from "./authority_store_codec.ts"; +import {canonicalTaskLease} from "./task_lease_state.ts"; +import {decodeRuntimeShadowRequest, verifyShadowSourceSnapshot, withShadowSourceLocks} from "./runtime_shadow.ts"; +import {loadLegacyCoordinationWriterFence} from "./legacy_writer_fence.ts"; +import {withShadowMaintenanceLock, ShadowManagementError, readShadowManagementState, + readRetainedShadowArtifacts, shadowManagementDirectory, requireShadowCaptureBinding} from "./shadow_management.ts"; +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 {FileAuthorityStore} from "./file_authority_store.ts"; + +export const COLD_SOURCE_INSPECTION_REQUEST_SCHEMA = "loopx_cold_source_inspection_request_v0"; +export const COLD_SOURCE_INSPECTION_RESULT_SCHEMA = "loopx_cold_source_inspection_result_v0"; + +/** Observe the existing drain owner's decisions without publishing a cursor, + * replaying an entry or reclaiming bytes. Failure keeps the raw inventory; + * an unknown original cannot be silently classified as settled. M is held. */ +async function reviewRetainedOutbox(root: string, goal: string, view: JsonObject): Promise { + const boundary = {executed: false, execution_authority_granted: false}; + try { + const binding = await requireShadowCaptureBinding(root, goal); + const partitions: JsonObject[] = []; + for (const partition of ["todos", "leases"] as const) { + const inventory = await drainInventory(root, goal, partition); + partitions.push(planShadowDrain({schema_version: SHADOW_DRAIN_PLAN_REQUEST_SCHEMA, + runtime_root: root, goal_id: goal, partition, + capture_lineage_id: binding.capture_lineage_id, store_identity: binding.store_identity, + source_root_digest: binding.source_root_digest, cursor: inventory.cursor, entries: inventory.entries, + remaining_entries: inventory.entries.length, budget_open: true, acknowledgement: null}, view)); + } + return {status: "planned", ...boundary, partitions}; + } catch (error) { + const failure = error as {reasonCode?: string; reason_code?: string; code?: string}; + return {status: "failed", ...boundary, + reason_code: failure.reasonCode ?? failure.reason_code ?? failure.code ?? "shadow_drain_request_invalid"}; + } +} + +export async function inspectColdCoordinationSource(value: unknown): Promise { + try { + const request = decodeRuntimeShadowRequest(value, COLD_SOURCE_INSPECTION_REQUEST_SCHEMA); + const management = shadowManagementDirectory(request.runtime_root, request.goal_id); + // The registered runtime may have a supported path alias. Its artifact + // subdirectories must not redirect a read or a maintenance lock elsewhere. + for (const path of [join(request.runtime_root, "authority-shadow"), + join(request.runtime_root, "authority-shadow", "file"), + join(request.runtime_root, "authority-shadow", "outbox"), + join(request.runtime_root, "authority-transition"), dirname(management), management]) { + try { + if (!(await lstat(path)).isDirectory()) throw new ShadowManagementError("shadow_outbox_layout_invalid"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + } + return await withShadowMaintenanceLock(request.runtime_root, request.goal_id, () => + withShadowSourceLocks(request, async () => { + const fence = await loadLegacyCoordinationWriterFence(request.runtime_root, request.goal_id); + if (fence.status !== "missing") throw new ShadowManagementError( + fence.status === "loaded" ? "legacy_authority_already_promoted" : fence.reason_code); + const paths = localAuthorityProviderPaths(request.runtime_root, request.goal_id); + for (const path of [paths.marker, new FileAuthorityStore(paths.file, request.goal_id, {existingOnly: true}).path]) { + try { await lstat(path); } + catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") continue; + throw error; + } + throw new ShadowManagementError("cold_source_canonical_authority_present", + "Use the canonical provider's storage/recovery path; Markdown is no longer an import source"); + } + await verifyShadowSourceSnapshot(request); + const artifacts = await readRetainedShadowArtifacts(request.runtime_root, request.goal_id); + const managementState = await readShadowManagementState(request.runtime_root, request.goal_id); + if (managementState?.status === "active" && artifacts.runtime_store === null) { + throw new ShadowManagementError("provider_read_unavailable"); + } + async function historyReadback(storeKind: "runtime_shadow" | "legacy_observation", present: boolean): Promise { + if (!present) return null; + const result = await readLocalAuthorityShadow({ + schema_version: LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA, + runtime_root: request.runtime_root, goal_id: request.goal_id, + store_kind: storeKind, read_model: "proof", + scan_limit: storeKind === "runtime_shadow" && managementState?.status === "active" ? 10000 : 0, + }); + if (result.status !== "loaded") throw new ShadowManagementError(String(result.reason_code ?? "provider_read_unavailable")); + return result; + } + const runtimeReadback = await historyReadback("runtime_shadow", artifacts.runtime_store !== null); + const legacyReadback = await historyReadback("legacy_observation", artifacts.legacy_store !== null); + // Inactive/interrupted capture must recover its own management operation + // first. A directory alone cannot supply the missing lineage authority. + const outboxReview = managementState?.status === "active" && runtimeReadback !== null + ? await reviewRetainedOutbox(request.runtime_root, request.goal_id, runtimeReadback) : null; + // Full proof is an internal input to both partition plans, not extra + // response history. Preserve the existing compact readback boundary. + if (runtimeReadback !== null) (runtimeReadback.proof as JsonObject).transactions = []; + const retainedLeases: JsonObject[] = []; + for (const entry of request.source_snapshot.lease_inventory as JsonObject[]) { + const name = String(entry.name); + // Source verification binds every file to its Goal and filename. Keep + // historical bytes separate from the graph's current execution edges. + const lease = canonicalAuthorityObject(JSON.parse(await readFile( + join(request.runtime_root, "goals", request.goal_id, "task-leases", name), "utf8")), "retained lease"); + canonicalTaskLease(lease, request.goal_id, name.slice(0, -5)); + retainedLeases.push(lease); + } + await verifyShadowSourceSnapshot(request); + if (canonicalAuthoritySha256(artifacts) !== canonicalAuthoritySha256( + await readRetainedShadowArtifacts(request.runtime_root, request.goal_id))) { + throw new ShadowManagementError("source_changed_retry"); + } + const todos = request.projection.todos as JsonObject[]; + return { + schema_version: COLD_SOURCE_INSPECTION_RESULT_SCHEMA, status: "inspected", + goal_id: request.goal_id, executed: false, import_ready: false, + writer_stop_verified: false, outbox_reconciliation_verified: false, + active_todo_count: todos.filter(todo => todo.archive_state === "active").length, + archived_todo_count: todos.filter(todo => todo.archive_state === "archive").length, + lease_file_count: retainedLeases.length, + // Expiry revokes execution; it never proves the writer process stopped. + leases_requiring_settlement: retainedLeases.filter(lease => lease.status === "active") + .map(lease => lease.todo_id), + retained_leases: retainedLeases, + capture: {management_state: managementState, artifacts, + runtime_shadow_readback: runtimeReadback, legacy_observation_readback: legacyReadback, + outbox_review: outboxReview}, + projection: request.projection, source_snapshot: request.source_snapshot, + decision_read_from_shadow: false, + }; + })); + } catch (error) { + return {schema_version: COLD_SOURCE_INSPECTION_RESULT_SCHEMA, status: "failed", + executed: false, import_ready: false, + reason_code: error instanceof ShadowManagementError ? error.reason_code : "cold_source_inspection_invalid", + reason: error instanceof Error ? error.message : "cold source inspection failed"}; + } +} + +/** App read model of the same verified observation. No source bytes, paths, + * execution identities or original receipts cross the HTTP boundary. */ +export async function inspectColdCoordinationStorage(value: unknown): Promise { + const observed = await inspectColdCoordinationSource(value); + const base = {ok: observed.status === "inspected", status: observed.status, + authority_changed: false, execution_authority_granted: false}; + if (observed.status !== "inspected") return {...base, current: null, reason_code: observed.reason_code}; + const capture = observed.capture as JsonObject; + const artifacts = capture.artifacts as JsonObject; + const outbox = artifacts.outbox as JsonObject | null; + const outboxInventory = outbox?.inventory as JsonObject | undefined; + return {...base, current: {goal_id: observed.goal_id, canonical: false, provider: null}, + cold_source: { + active_todo_count: observed.active_todo_count, archived_todo_count: observed.archived_todo_count, + lease_file_count: observed.lease_file_count, + unsettled_lease_count: (observed.leases_requiring_settlement as string[]).length, + capture_artifacts_present: artifacts.management_state !== null || artifacts.runtime_store !== null || + artifacts.legacy_store !== null || (artifacts.rollback_archives as JsonObject[]).length > 0, + outbox_files_present: ((outboxInventory?.entries ?? []) as JsonObject[]).some(entry => entry.kind === "file"), + import_ready: false, writer_stop_verified: false, outbox_reconciliation_verified: false, + }}; +} diff --git a/loopx/control_plane/coordination/runtime_shadow.py b/loopx/control_plane/coordination/runtime_shadow.py index ded83dcebe..d6ec31b35b 100644 --- a/loopx/control_plane/coordination/runtime_shadow.py +++ b/loopx/control_plane/coordination/runtime_shadow.py @@ -265,6 +265,7 @@ def capture_todo_archive_dependencies(todos: list[dict[str, Any]], state_text: s def build_runtime_shadow_source_snapshot( *, goal: Mapping[str, Any], runtime_root: Path, state_path: Path, registry_path: Path, + include_all_archived_todos: bool = False, ) -> tuple[dict[str, object], dict[str, object]]: """Bind the supplied Goal and every derived fact to one registry observation.""" from ...agent_registry import registered_agent_ids_for_goal @@ -281,6 +282,7 @@ def build_runtime_shadow_source_snapshot( projection, snapshot = _build_runtime_shadow_source_snapshot( goal=current, runtime_root=runtime_root, state_path=state_path, registry_path=registry_path, registry=registry, + include_all_archived_todos=include_all_archived_todos, ) snapshot["registry_source"] = { **witness, "registered_agents": registered_agent_ids_for_goal(current), @@ -334,6 +336,7 @@ def _source_lease_bytes(path: Path) -> bytes: def _build_runtime_shadow_source_snapshot( *, goal: Mapping[str, Any], runtime_root: Path, state_path: Path, registry_path: Path, registry: dict[str, Any], + include_all_archived_todos: bool = False, ) -> tuple[dict[str, object], dict[str, object]]: """Project exactly the bytes carried by one ephemeral source precondition. @@ -383,10 +386,26 @@ def read_evidence(path: Path) -> bytes | None: raise ShadowManagementError("legacy_todo_event_source_retired", "legacy_todo_event_source_retired: preserve and export legacy Todo events with a compatible older release before migration") - fields = parse_active_state_todos(state_text, goal=dict(goal), state_path=state_path, item_limit=None, rollout_events=rollout_events) - todos = todo_summaries_from_fields(fields=fields, source="markdown_active_state", rollout_events=rollout_events, roles=["user", "agent"], status=None, - todo_id=None, agent_id=None, limit=None).todos - todos = capture_todo_archive_dependencies(todos, state_text) + if include_all_archived_todos: + from ..todos.active_state_todo_parser import parse_todo_source + from ..todos.todo_summary import structured_todo_item, canonical_todo_read_record + active, archived, _ = parse_todo_source(state_text) + for item in archived: + if item.get("role") not in {"agent", "user"}: + raise ShadowManagementError("cold_source_archive_role_unproved", + "Archived Todo role must be explicit; cold inspection cannot infer its owner from prose") + # A cold inventory reads persisted records, not attention summaries or + # the live graph's dependency-only archive. Reuse the full record codec + # for both sections, retaining their original heading and full text. + todos = [canonical_todo_read_record(structured_todo_item( + item, role=item["role"], source_section=item["source_section"], + archive_state=item["archive_state"], text_limit=None)) + for item in [*active["user"], *active["agent"], *archived]] + else: + fields = parse_active_state_todos(state_text, goal=dict(goal), state_path=state_path, item_limit=None, rollout_events=rollout_events) + todos = todo_summaries_from_fields(fields=fields, source="markdown_active_state", rollout_events=rollout_events, roles=["user", "agent"], status=None, + todo_id=None, agent_id=None, limit=None).todos + todos = capture_todo_archive_dependencies(todos, state_text) leases: list[dict[str, Any]] = [] inventory: list[dict[str, object]] = [] for path in _source_lease_paths(runtime_root, goal_id): diff --git a/loopx/control_plane/coordination/shadow_management.ts b/loopx/control_plane/coordination/shadow_management.ts index 76f163c932..02698590d2 100644 --- a/loopx/control_plane/coordination/shadow_management.ts +++ b/loopx/control_plane/coordination/shadow_management.ts @@ -509,6 +509,84 @@ async function inventory(path: string): Promise { await visit(path, ""); return { entries, digest: managementDigest(entries) }; } + +/** Discover original capture files without reinterpreting them as delivery or + * import receipts. The caller holds maintenance and primary source locks. */ +export async function readRetainedShadowArtifacts(root: string, goal: string): Promise { + async function retainedFile(path: string): Promise { + try { + if (!(await lstat(dirname(path))).isDirectory() || !(await lstat(path)).isFile()) { + throw new ShadowManagementError("shadow_outbox_layout_invalid"); + } + const bytes = await readFile(path); + return {path, size: bytes.length, sha256: `sha256:${createHash("sha256").update(bytes).digest("hex")}`}; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; + throw error; + } + } + async function directory(path: string): Promise { + const value = await inventory(path); + return value === null ? null : {path, inventory: value}; + } + const management = shadowManagementDirectory(root, goal); + const state = await readShadowManagementState(root, goal); + const operations = join(management, "operations"); + const operationInventory = await directory(operations); + const store = new FileAuthorityStore(join(root, "authority-shadow", "file-v0"), goal, {existingOnly: true}); + const legacy = new FileAuthorityStore(join(root, "authority-shadow", "file", goal), goal, {existingOnly: true}); + const rollbackArchives: JsonObject[] = []; + // File rollback archives share a provider directory across Goals. Derive + // only this Goal's originals through its existing management manifests. + const entries = (operationInventory?.inventory as JsonObject | undefined)?.entries as JsonObject[] | undefined; + for (const entry of entries ?? []) { + const relative = String(entry.path); + if (entry.kind !== "file" || relative.split("/").length !== 2 || !relative.endsWith("/manifest.json")) continue; + const raw = await readJson(join(operations, relative)); + if (!raw) throw new ShadowManagementError("shadow_management_manifest_invalid"); + const request = requestOf(raw.request); + if (request.goal_id !== goal || !sameSourceRoot(request.runtime_root, root) + || join(operations, relative) !== manifestPath({...request, runtime_root: root})) { + throw new ShadowManagementError("shadow_management_manifest_invalid"); + } + const terminal = await readJson(resultPath(request)); + const current = state?.operation.operation_id === request.operation_id ? state : null; + const manifest = await loadManifest(request, {operation: current?.operation ?? terminal + ?? {manifest_digest: managementDigest(raw)}}); + if (terminal !== null) { + if (terminal.request_digest !== requestDigest(request) || !isAuthorityJsonObject(terminal.result)) { + throw new ShadowManagementError("shadow_management_result_invalid"); + } + await validateReplayResult(request, manifest, terminal.result, state); + } else if (current?.result) await validateReplayResult(request, manifest, current.result, state); + const completed = terminal !== null || (state?.operation.operation_id === request.operation_id && state.result !== null); + if (manifest.kind !== "rollback") continue; + if (completed && !same(await inventory(archiveOutboxPath(request)), manifest.outbox)) { + throw new ShadowManagementError("rollback_archive_readback_mismatch"); + } + if (manifest.candidate === null) continue; + const archive = await retainedFile(store.authorityArchivePath(request.operation_id)); + if (archive === null && completed) throw new ShadowManagementError("rollback_archive_readback_mismatch"); + // A prepared rollback may not have renamed the source yet. Keep that + // absence visible; inspection never finishes the interrupted operation. + if (archive !== null) { + if (archive.sha256 !== (manifest.candidate as JsonObject).sha256) { + throw new ShadowManagementError("rollback_archive_readback_mismatch"); + } + rollbackArchives.push(archive); + } + } + return { + management_state: await retainedFile(shadowManagementStatePath(root, goal)), + management_operations: operationInventory, + outbox: await directory(join(root, "authority-shadow", "outbox", goal)), + runtime_store: await retainedFile(store.path), + runtime_store_identity: await retainedFile(store.identityPath), + rollback_archives: rollbackArchives, + legacy_observation: await directory(legacy.directory), + legacy_store: await retainedFile(legacy.path), + }; +} async function fileDigest(path: string): Promise { try { return `sha256:${createHash("sha256").update(await readFile(path)).digest("hex")}`; } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; throw error; } diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index f04aff9971..2f911e193f 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -285,6 +285,8 @@ export function createEffectRuntimeHandlers( ["todo.archive.capture_dependencies", lazyHandler(() => Promise.all([import("./todos/archive_capture.ts"), import("./coordination/source_transfer.ts")]), ([{captureArchivedTodoDependencies}, {withCoordinationSourceTransfer}]) => withCoordinationSourceTransfer("todo.archive.capture_dependencies", captureArchivedTodoDependencies))], ["agent.supervisor.plan_append", lazyHandler(() => import("./agents/supervisor_event_append.ts"), ({planSupervisorEventAppend}) => planSupervisorEventAppend)], ["coordination.source.project", lazyHandler(() => Promise.all([import("./coordination/source_projection.ts"), import("./coordination/source_transfer.ts")]), ([{projectCoordinationSource}, {withCoordinationSourceTransfer}]) => withCoordinationSourceTransfer("coordination.source.project", projectCoordinationSource))], + ["coordination.source.inspect", lazyHandler(() => Promise.all([import("./coordination/cold_source_inspection.ts"), import("./coordination/source_transfer.ts")]), ([{inspectColdCoordinationSource}, {withCoordinationSourceTransfer}]) => withCoordinationSourceTransfer("coordination.source.inspect", inspectColdCoordinationSource))], + ["coordination.source.inspect_storage", lazyHandler(() => Promise.all([import("./coordination/cold_source_inspection.ts"), import("./coordination/source_transfer.ts")]), ([{inspectColdCoordinationStorage}, {withCoordinationSourceTransfer}]) => withCoordinationSourceTransfer("coordination.source.inspect_storage", inspectColdCoordinationStorage))], ["todo.monitor_metadata.plan", lazyHandler(() => import("./todos/monitor_metadata.ts"), ({planMonitorMetadata}) => planMonitorMetadata)], ["todo.authoring_scope.plan", lazyHandler(() => import("./todos/authoring_scope.ts"), ({planTodoAuthoringScope}) => planTodoAuthoringScope)], ["todo.contract_diagnostics.evaluate", lazyHandler(() => import("./todos/authoring_scope.ts"), ({evaluateTodoContractDiagnostics}) => evaluateTodoContractDiagnostics)], diff --git a/loopx/presentation/goal_storage_api.py b/loopx/presentation/goal_storage_api.py index 5389f68e07..2d146aa732 100644 --- a/loopx/presentation/goal_storage_api.py +++ b/loopx/presentation/goal_storage_api.py @@ -44,6 +44,7 @@ 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", + "cold_source", ) if key in result} payload["goal_id"] = goal_id if preview_id is not None: @@ -56,9 +57,29 @@ def _storage_inspect(self) -> None: if set(query) != {"goal_id"} or len(query["goal_id"]) != 1: raise ValueError("goal_id is required exactly once") goal_id = query["goal_id"][0] - self._storage_send(self._storage_owner(goal_id, action="migration-readback"), goal_id) + registry, goal = self._registry_and_goal(goal_id) except (KeyError, TypeError, ValueError): self._send_error("Choose a registered Goal.", status=400, error_code="invalid_goal_storage_request") + return + try: + result = self._storage_owner(goal_id, action="migration-readback") + current = result.get("current") + if result.get("ok") and isinstance(current, dict) and current.get("canonical") is False: + from ..control_plane.coordination.local_authority_shadow_projection import source_effect_runtime_result + from ..control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot + from ..state_refresh import resolve_goal_state + + _, _, state_path = resolve_goal_state(registry=registry, goal_id=goal_id, + project_override=None, state_file_override=None) + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, + runtime_root=self.server.runtime_root, state_path=state_path, + registry_path=self.server.registry_path, include_all_archived_todos=True) + result = source_effect_runtime_result("coordination.source.inspect_storage", { + "schema_version": "loopx_cold_source_inspection_request_v0", + "runtime_root": str(self.server.runtime_root.expanduser().absolute()), + "goal_id": goal_id, "projection": projection, "source_snapshot": snapshot, + }) + self._storage_send(result, goal_id) except Exception: # noqa: BLE001 - local provider errors stay private. self._send_error("Current storage unavailable.", 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 a66ddbac18..d11f18fb2d 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": 218, + "line": 231, "column": 20, "kind": "codec_read", "api": "load_project_registry", @@ -1111,7 +1111,7 @@ }, { "site": "loopx/control_plane/coordination/runtime_shadow.py::.build_runtime_shadow_source_snapshot::codec_read:load_project_registry#1", - "line": 277, + "line": 278, "column": 20, "kind": "codec_read", "api": "load_project_registry", diff --git a/tests/control_plane/test_chat_goal_storage.py b/tests/control_plane/test_chat_goal_storage.py index c96d3c897a..ac053f6472 100644 --- a/tests/control_plane/test_chat_goal_storage.py +++ b/tests/control_plane/test_chat_goal_storage.py @@ -199,3 +199,98 @@ def test_original_receipt_is_readable_when_live_provider_is_unavailable(api): assert operate(call, plan)[0] == 409 finally: unavailable.rename(directory) + + +@pytest.fixture(params=[False, True], ids=["cold", "retained-capture"]) +def old_api(tmp_path, monkeypatch, request): + from tests.control_plane.test_cold_source_inspection import cold_workspace + from tests.control_plane.shadow_e2e_fixture import workspace + + isolate_sqlite_runtime(tmp_path, monkeypatch) + ws = workspace(tmp_path) if request.param else cold_workspace(tmp_path) + if request.param: + ws.add("Private original history must stay on disk") + outbox = ws.runtime / "authority-shadow" / "outbox" / ws.goal / "todos" + (outbox / "original-unrecognized.json").write_bytes(b"{malformed original outbox") + directory = ws.runtime / "goals" / ws.goal / "task-leases" + directory.mkdir(parents=True, exist_ok=True) + (directory / "removed.json").write_text(json.dumps({ + "schema_version": "task_lease_v0", "goal_id": ws.goal, "todo_id": "removed", + "owner": "agent-a", "status": "active", "idempotency_key": "old-work", + "version": 4, "lease_epoch": 2, "expires_at": "2000-01-01T00:00:00Z", "write_scopes": [], + })) + server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler) + server.runtime_root, server.registry_path = ws.runtime, ws.registry + server.runtime_root_override, server.verbose = str(ws.runtime), False + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + + def call(): + connection = http.client.HTTPConnection("127.0.0.1", server.server_port, timeout=90) + connection.request("GET", f"/api/chat/goal-storage?goal_id={ws.goal}") + response = connection.getresponse() + result = json.loads(response.read()) + connection.close() + encoded = json.dumps(result) + assert str(tmp_path) not in encoded and "Private original history" not in encoded + assert "old-work" not in encoded and "malformed original outbox" not in encoded + return response.status, result + + yield call, ws, request.param + server.shutdown() + thread.join(5) + server.server_close() + restart_effect_runtime() + + +def test_registered_old_source_read_failure_is_unavailable_and_recovers(old_api): + call, ws, _ = old_api + source = ws.state.read_bytes() + ws.state.write_bytes(b"\xff") + code, failed = call() + assert code == 503 and failed["error_code"] == "goal_storage_unavailable" + assert "cold_source" not in failed and "current" not in failed + ws.state.write_bytes(source) + code, result = call() + assert code == 200 and result["cold_source"]["import_ready"] is False + + +def test_old_source_http_inventory_keeps_history_and_does_not_grant_import(old_api): + from tests.control_plane.test_cold_source_inspection import capture_bytes + + call, ws, captured = old_api + before = capture_bytes(ws) + source = ws.state.read_bytes(), ws.registry.read_bytes() + for _ in range(2): + code, result = call() + assert code == 200 and result["ok"], result + assert result["current"]["canonical"] is False and result["current"]["provider"] is None + old = result["cold_source"] + assert old["active_todo_count"] == 1 + assert old["archived_todo_count"] == (0 if captured else 1) + # Even an orphan, expired active lease needs settlement; not Host stop. + assert old["unsettled_lease_count"] == 1 + assert old["capture_artifacts_present"] is captured + assert old["outbox_files_present"] is captured + assert old["writer_stop_verified"] is old["outbox_reconciliation_verified"] is old["import_ready"] is False + assert result["authority_changed"] is result["execution_authority_granted"] is False + assert "outbox_review" not in json.dumps(result) + assert before == capture_bytes(ws) + assert source == (ws.state.read_bytes(), ws.registry.read_bytes()) + assert not (ws.runtime / "authority").exists() + + +def test_old_source_http_refuses_canonical_presence_and_recovers_fresh_observation(old_api): + import hashlib + + call, ws, _ = old_api + marker = ws.runtime / "authority" / ("provider-" + hashlib.sha256(ws.goal.encode()).hexdigest() + ".json") + marker.parent.mkdir(parents=True) + marker.write_text('{"provider":"sqlite"}') + code, result = call() + assert code == 409 and result["ok"] is False and result["current"] is None, result + assert result["reason_code"] == "cold_source_canonical_authority_present" + assert "cold_source" not in result + assert marker.read_text() == '{"provider":"sqlite"}' + marker.unlink() + assert call()[0] == 200 diff --git a/tests/control_plane/test_cold_source_disposition_e2e.py b/tests/control_plane/test_cold_source_disposition_e2e.py new file mode 100644 index 0000000000..aeaeef33b6 --- /dev/null +++ b/tests/control_plane/test_cold_source_disposition_e2e.py @@ -0,0 +1,173 @@ +"""Retained originals are handled without bringing back their Python producers. + +The original CLI creates real interrupted writes. A separate copied runtime +then dispatches the shipped TS effects with those producer files absent. The +OS-lock adapter remains: language retirement must preserve useful Host IO. +""" +from __future__ import annotations + +import json +from pathlib import Path +import shutil +import subprocess +import sys + +import loopx +import pytest + +from loopx.control_plane.coordination.coordination_state_contract_generated import ( + LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA, +) +from tests.control_plane.shadow_e2e_fixture import workspace + + +pytestmark = pytest.mark.stage2c_e2e + + +@pytest.fixture(scope="module") +def receiver(tmp_path_factory): + root = tmp_path_factory.mktemp("receiver") + package = root / "loopx" + shutil.copytree(Path(loopx.__file__).parent, package, + ignore=shutil.ignore_patterns("__pycache__", "testing", "*.pyc")) + removed = ["todos.py", "bootstrap.py", + "control_plane/coordination/runtime_shadow_writer_adapter.py", + "control_plane/coordination/local_authority_shadow_outbox.py"] + for relative in removed: + (package / relative).unlink() + assert all(not (package / relative).exists() for relative in removed) + assert (package / "control_plane/coordination/shadow_lock_host.py").is_file() + module = package / "control_plane/effect_runtime_handlers.ts" + script = ( + f"import {{createEffectRuntimeHandlers, dispatchEffectRuntimeMethod}} from {json.dumps(module.as_uri())};" + "let raw=''; for await (const chunk of process.stdin) raw+=chunk;" + "const {method,request}=JSON.parse(raw);" + "const handlers=createEffectRuntimeHandlers({fingerprint:'disposable',requestShutdown:()=>{}});" + "console.log(JSON.stringify(await dispatchEffectRuntimeMethod(handlers,method,request)));" + ) + + def dispatch(method, request): + result = subprocess.run( + ["node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", script], + input=json.dumps({"method": method, "request": request}), + cwd=root, capture_output=True, text=True, timeout=45, + ) + assert result.returncode == 0, result.stdout + result.stderr + return json.loads(result.stdout) + + return dispatch + + +def drain(receiver, fixture): + return receiver("coordination.runtime_shadow.drain", { + "schema_version": "loopx_shadow_drain_v0", "runtime_root": str(fixture.runtime), + "goal_id": fixture.goal, "python_executable": sys.executable, "config_enabled": True, + "max_entries": 32, "budget_seconds": 10, "lock_timeout_seconds": 2, + }) + + +def inspect(fixture): + return fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + + +def originals(path): + return {str(p.relative_to(path)): p.read_bytes() for p in path.rglob("*") + if p.is_file() and not p.name.endswith(".lock")} + + +@pytest.mark.parametrize("window,resolution,no_op", [ + ("before_replace", "abandoned", True), + ("before_marker", "committed_proven_by_readback", False), + ("before_commit", "committed", False), + ("after_commit", "committed", False), +]) +def test_stopped_original_todo_recovers_once_without_python_producer( + tmp_path, receiver, window, resolution, no_op, +): + fixture = workspace(tmp_path) + fixture.crash(window, "todo", "add", "--role", "agent", "--text", "Interrupted original") + source = fixture.state.read_bytes() + before = inspect(fixture) + partition = fixture.runtime / "authority-shadow/outbox" / fixture.goal / "todos" + prepared = json.loads(next(partition.glob("*.prepared.json")).read_bytes()) + initial = before["capture"]["runtime_shadow_readback"]["proof"]["last_applied_sequences"]["todos"] + result = drain(receiver, fixture) + assert result["ok"] is True, result + assert result["replayed"] == (1 if window == "after_commit" else 0) + assert result["delivered"] + result["reconciled"] == (0 if window == "after_commit" else 1) + assert fixture.state.read_bytes() == source + after = inspect(fixture) + view = after["capture"]["runtime_shadow_readback"] + receipt_view = receiver("coordination.runtime_shadow.outbox_read", { + "schema_version": LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA, + "runtime_root": str(fixture.runtime), "goal_id": fixture.goal, "store_kind": "runtime_shadow", + "scan_after_cursor": None, "scan_limit": 0, "read_model": "proof", + "receipt_operation_id": prepared["entry_id"], + }) + assert receipt_view["status"] == "loaded", receipt_view + receipt = receipt_view["proof"]["receipt"]["receipts"][0] + assert receipt["entry_id"] == prepared["entry_id"] + assert receipt["resolution"] == resolution and receipt["no_op"] is no_op + assert receipt["seq"] == 1 and view["proof"]["last_sequences"]["todos"] == 1 + assert view["proof"]["last_applied_sequences"]["todos"] == (0 if no_op else 1) + assert initial == (1 if window == "after_commit" else 0) + assert [p.name for p in partition.iterdir()] == ["drain-cursor.json"] + retained = originals(fixture.runtime) + assert drain(receiver, fixture)["outcome"] == "nothing_pending" + assert originals(fixture.runtime) == retained + # Successful original disposition is not stopped-Host/import qualification. + assert after["import_ready"] is False and after["writer_stop_verified"] is False + assert not (fixture.runtime / "authority").exists() + + +@pytest.mark.parametrize("window", ["before_commit", "after_commit"]) +def test_original_lease_receipt_never_grants_or_releases_work(tmp_path, receiver, window): + fixture = workspace(tmp_path) + todo = fixture.add("Original leased task")["todo_id"] + fixture.crash(window, "task-lease", "acquire", "--todo-id", todo, "--owner", "agent-a", + "--idempotency-key", "original", "--ttl-seconds", "120") + path = fixture.runtime / "goals" / fixture.goal / "task-leases" / (todo + ".json") + lease = path.read_bytes() + assert drain(receiver, fixture)["ok"] is True + assert path.read_bytes() == lease + after = inspect(fixture) + assert after["leases_requiring_settlement"] == [todo] + assert after["retained_leases"][0]["status"] == "active" + assert after["import_ready"] is False + retained = originals(fixture.runtime) + assert drain(receiver, fixture)["outcome"] == "nothing_pending" + assert originals(fixture.runtime) == retained and path.read_bytes() == lease + + +def test_unproved_original_stays_intact_and_rolls_back_same_operation(tmp_path, receiver): + fixture = workspace(tmp_path) + fixture.crash("before_marker", "handoff-mode", "set", "--mode", "soft_claim") + fixture.cli("handoff-mode", "set", "--mode", "hard_lease") + before = originals(fixture.runtime) + result = drain(receiver, fixture) + assert result["ok"] is False and result["reason_code"] == "outbox_source_unproved" + assert originals(fixture.runtime) == before + inventory = inspect(fixture) + source = fixture.state.read_bytes() + artifacts = inventory["capture"]["artifacts"] + candidate = Path(artifacts["runtime_store"]["path"]).read_bytes() + pending = originals(Path(artifacts["outbox"]["path"])) + request = {"schema_version": "loopx_coordination_runtime_shadow_rollback_v0", + "runtime_root": str(fixture.runtime), "goal_id": fixture.goal, + "operation_id": "retain-original", "expected_bootstrap_operation_id": None, + "expected_provider_revision": inventory["capture"]["runtime_shadow_readback"]["provider_revision"], + "projection": inventory["projection"], "source_snapshot": inventory["source_snapshot"]} + rolled = receiver("coordination.runtime_shadow.rollback", request) + assert rolled["status"] == "applied", rolled + assert Path(rolled["candidate_archive_path"]).read_bytes() == candidate + assert originals(Path(rolled["outbox_archive_path"])) == pending + retained = originals(fixture.runtime) + replay = receiver("coordination.runtime_shadow.rollback", request) + assert replay["status"] == "replayed" and replay["operation_id"] == rolled["operation_id"] + assert originals(fixture.runtime) == retained and fixture.state.read_bytes() == source + refused = drain(receiver, fixture) + assert refused["ok"] is False and refused["outcome"] == "stopped" + assert refused["delivered"] == 0 and refused["replayed"] == 0 + assert originals(fixture.runtime) == retained + assert inspect(fixture)["capture"]["outbox_review"] is None + assert not (fixture.runtime / "authority").exists() diff --git a/tests/control_plane/test_cold_source_inspection.py b/tests/control_plane/test_cold_source_inspection.py new file mode 100644 index 0000000000..2620fd7ea0 --- /dev/null +++ b/tests/control_plane/test_cold_source_inspection.py @@ -0,0 +1,352 @@ +"""A cold source can be inventoried without manufacturing capture history.""" +import json +import hashlib +from pathlib import Path +import subprocess +import sys + +import pytest + +from tests.control_plane.shadow_e2e_fixture import workspace + + +def capture_bytes(fixture): + """Original history files; transient lock files are not receipt sources.""" + return {str(p.relative_to(fixture.runtime)): p.read_bytes() + for root in (fixture.runtime / "authority-shadow", fixture.runtime / "authority-transition") + for p in root.rglob("*") if p.is_file() and not p.name.endswith(".lock")} + + +def test_capture_inventory_keeps_original_history_and_outbox_bytes(tmp_path): + fixture = workspace(tmp_path) + fixture.add("Original captured operation") + outbox = fixture.runtime / "authority-shadow" / "outbox" / fixture.goal / "todos" + outbox.mkdir(exist_ok=True) + residue = outbox / "unrecognized-original-receipt.json" + residue.write_bytes(b"{malformed original bytes") + before = capture_bytes(fixture) + result = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + capture = result["capture"] + assert capture["management_state"]["status"] == "active" + assert capture["runtime_shadow_readback"]["status"] == "loaded" + assert capture["runtime_shadow_readback"]["proof"]["last_applied_sequences"]["todos"] >= 1 + witness = capture["artifacts"]["runtime_store"] + original = next(data for path, data in before.items() + if path == str(Path(witness["path"]).relative_to(fixture.runtime))) + assert witness["sha256"] == "sha256:" + hashlib.sha256(original).hexdigest() + entries = capture["artifacts"]["outbox"]["inventory"]["entries"] + raw = next(e for e in entries if e["path"] == "todos/" + residue.name) + assert raw["sha256"] == "sha256:" + hashlib.sha256(residue.read_bytes()).hexdigest() + assert result["outbox_reconciliation_verified"] is False + assert result["import_ready"] is False + review = capture["outbox_review"] + assert review["status"] == "failed" + assert review["reason_code"] == "outbox_file_invalid" + assert review["executed"] is False + assert before == capture_bytes(fixture) + assert not (fixture.runtime / "authority").exists() + + +def test_capture_inventory_keeps_rolled_back_operation_archives(tmp_path): + fixture = workspace(tmp_path) + fixture.add("Original operation before rollback") + revision = fixture.cli("coordination-shadow", "inspect")["inspection"]["provider_revision"] + rollback = fixture.cli("coordination-shadow", "rollback", "--provider-revision", revision, "--execute") + assert rollback["rollback"]["status"] == "applied" + before = capture_bytes(fixture) + capture = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"]["capture"] + assert capture["management_state"]["status"] == "inactive" + assert capture["artifacts"]["runtime_store"] is None + assert capture["runtime_shadow_readback"] is None + assert capture["outbox_review"] is None + archives = capture["artifacts"]["rollback_archives"] + assert len(archives) == 1 + assert archives[0]["path"] == rollback["rollback"]["candidate_archive_path"] + assert archives[0]["sha256"] == "sha256:" + hashlib.sha256(Path(archives[0]["path"]).read_bytes()).hexdigest() + assert before == capture_bytes(fixture) + + +@pytest.mark.parametrize("window", ["before_commit", "after_commit"]) +def test_original_lease_outbox_uses_its_own_partition_without_renewal_or_cleanup(tmp_path, window): + fixture = workspace(tmp_path) + todo_id = fixture.add("Original leased task")["todo_id"] + fixture.crash(window, "task-lease", "acquire", "--todo-id", todo_id, + "--owner", "agent-a", "--idempotency-key", "original-lease", "--ttl-seconds", "120") + before = capture_bytes(fixture) + lease_path = fixture.runtime / "goals" / fixture.goal / "task-leases" / (todo_id + ".json") + lease_bytes = lease_path.read_bytes() + result = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + review = result["capture"]["outbox_review"] + assert review["status"] == "planned" and review["executed"] is False + leases = next(p for p in review["partitions"] if p["partition"] == "leases") + assert len(leases["pending_entry_ids"]) == (1 if window == "before_commit" else 0) + assert len(leases["reclaim_entry_ids"]) == (0 if window == "before_commit" else 1) + assert leases["next_seq"] == (1 if window == "before_commit" else 2) + assert todo_id in result["leases_requiring_settlement"] + assert lease_path.read_bytes() == lease_bytes and before == capture_bytes(fixture) + + +@pytest.mark.parametrize("corrupt", [False, True]) +def test_capture_inventory_reads_original_observation_through_existing_store(tmp_path, corrupt): + fixture = workspace(tmp_path) + source = fixture.runtime / "authority-shadow" / "file-v0" + retained = fixture.runtime / "authority-shadow" / "file" / fixture.goal + retained.mkdir(parents=True) + # A synthetic historical copy of a real provider document, not a mock or + # an import receipt. The historical reader owns its schema validation. + for path in source.iterdir(): + if path.is_file() and not path.name.endswith(".lock"): + (retained / path.name).write_bytes(path.read_bytes()) + if corrupt: + next(retained.glob("authority-store-*.json")).write_text('{"invalid":true}') + before = capture_bytes(fixture) + result = fixture.cli("coordination-shadow", "inspect-source", success=not corrupt) + if corrupt: + assert result["source_inventory"]["status"] == "failed" + else: + capture = result["source_inventory"]["capture"] + assert capture["legacy_observation_readback"]["status"] == "loaded" + assert capture["artifacts"]["legacy_store"]["path"].startswith(str(retained)) + assert result["source_inventory"]["import_ready"] is False + assert before == capture_bytes(fixture) + + +@pytest.mark.parametrize("unsafe", ["corrupt_history", "symlink_outbox", "symlink_shadow_root", "invalid_original_result"]) +def test_capture_inventory_refuses_unreadable_or_unsafe_original_history(tmp_path, unsafe): + fixture = workspace(tmp_path) + if unsafe == "corrupt_history": + store = next((fixture.runtime / "authority-shadow" / "file-v0").glob("authority-store-*.json")) + store.write_text('{"invalid":true}') + elif unsafe == "symlink_outbox": + outbox = fixture.runtime / "authority-shadow" / "outbox" / fixture.goal + (outbox / "external").symlink_to(fixture.state) + elif unsafe == "symlink_shadow_root": + shadow = fixture.runtime / "authority-shadow" + original = tmp_path / "external-shadow" + shadow.rename(original) + shadow.symlink_to(original, target_is_directory=True) + else: + state = next((fixture.runtime / "authority-transition").rglob("state.json")) + original = json.loads(state.read_text()) + original["result"]["primary_writeback_preserved"] = False + state.write_text(json.dumps(original)) + before = capture_bytes(fixture) + result = fixture.cli("coordination-shadow", "inspect-source", success=False) + assert result["ok"] is False + assert result["source_inventory"]["status"] == "failed" + assert result["source_inventory"]["import_ready"] is False + assert before == capture_bytes(fixture) + + +@pytest.mark.parametrize("window", ["before_commit", "after_commit"]) +def test_capture_inventory_does_not_drain_crashed_original_outbox(tmp_path, window): + fixture = workspace(tmp_path) + fixture.crash(window, "todo", "add", "--role", "agent", "--text", "Retained original operation") + before = capture_bytes(fixture) + result = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + entries = result["capture"]["artifacts"]["outbox"]["inventory"]["entries"] + assert any(e["path"].endswith(".prepared.json") for e in entries) + assert any(e["path"].endswith(".committed.json") for e in entries) + assert result["import_ready"] is False and result["outbox_reconciliation_verified"] is False + review = result["capture"]["outbox_review"] + assert review["status"] == "planned" and review["executed"] is False + plans = review["partitions"] + todos = next(p for p in plans if p["partition"] == "todos") + assert len(todos["pending_entry_ids"]) == (1 if window == "before_commit" else 0) + assert len(todos["reclaim_entry_ids"]) == (0 if window == "before_commit" else 1) + assert len(todos["replay_entries"]) == (0 if window == "before_commit" else 1) + assert result["writer_stop_verified"] is False + assert before == capture_bytes(fixture) + + +@pytest.mark.parametrize("defect", ["receipt_bytes", "orphan_marker", "foreign_lineage"]) +def test_outbox_review_refuses_unproved_disposition_but_retains_original_bytes(tmp_path, defect): + fixture = workspace(tmp_path) + window = "after_commit" if defect == "receipt_bytes" else "before_commit" + fixture.crash(window, "todo", "add", "--role", "agent", "--text", "Original operation") + directory = fixture.runtime / "authority-shadow" / "outbox" / fixture.goal / "todos" + prepared = next(directory.glob("*.prepared.json")) + if defect == "orphan_marker": + prepared.unlink() + else: + record = json.loads(prepared.read_text()) + if defect == "receipt_bytes": + record["writer"]["operation_id"] = "different-original-operation" + else: + record["capture_lineage_id"] = "foreign-lineage" + prepared.write_text(json.dumps(record)) + before = capture_bytes(fixture) + result = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + review = result["capture"]["outbox_review"] + assert review["status"] == "failed" + assert review["reason_code"] == ("outbox_receipt_mismatch" if defect == "receipt_bytes" else "outbox_file_invalid") + assert review["executed"] is False + assert result["import_ready"] is False and result["outbox_reconciliation_verified"] is False + assert before == capture_bytes(fixture) + + +@pytest.mark.parametrize("archive", ["candidate", "outbox"]) +def test_capture_inventory_refuses_missing_terminal_rollback_archive(tmp_path, archive): + fixture = workspace(tmp_path) + revision = fixture.cli("coordination-shadow", "inspect")["inspection"]["provider_revision"] + result = fixture.cli("coordination-shadow", "rollback", "--provider-revision", revision, "--execute") + if archive == "candidate": + Path(result["rollback"]["candidate_archive_path"]).unlink() + else: + Path(result["rollback"]["outbox_archive_path"], "manifest.json").unlink() + before = capture_bytes(fixture) + refused = fixture.cli("coordination-shadow", "inspect-source", success=False) + assert refused["source_inventory"]["reason_code"] == "rollback_archive_readback_mismatch" + assert before == capture_bytes(fixture) + + +def cold_workspace(tmp_path): + fixture = workspace(tmp_path, bootstrap=False) + registry = json.loads(fixture.registry.read_text()) + registry["goals"][0]["coordination"].pop("runtime_shadow") + fixture.registry.write_text(json.dumps(registry)) + fixture.state.write_text( + "---\ngoal_id: goal-e2e\nhandoff_mode: soft_claim\n---\n" + "## Agent Todo\n- [ ] Keep active\n" + " \n" + "## Todo Archive\n- [x] Keep unrelated archived evidence\n" + " \n" + ) + return fixture + + +def test_inventory_preserves_unreferenced_archive_without_shadow_or_effects(tmp_path): + fixture = cold_workspace(tmp_path) + before = {p: p.read_bytes() for p in (fixture.registry, fixture.state)} + result = fixture.cli("coordination-shadow", "inspect-source") + inventory = result["source_inventory"] + assert inventory["status"] == "inspected" + assert inventory["active_todo_count"] == 1 + assert inventory["archived_todo_count"] == 1 + archived = next(t for t in inventory["projection"]["todos"] if t["todo_id"] == "todo_archived") + assert archived["text"] == "Keep unrelated archived evidence" + assert archived["note"] == "kept" and archived["watch_only"] == "false" + assert inventory["import_ready"] is False + assert inventory["writer_stop_verified"] is False + assert inventory["outbox_reconciliation_verified"] is False + assert inventory["capture"]["outbox_review"] is None + assert result["executed"] is False + assert not (fixture.runtime / "authority-shadow").exists() + assert not (fixture.runtime / "authority").exists() + assert before == {p: p.read_bytes() for p in before} + + +@pytest.mark.parametrize("original, todo_id, heading", [ + ("Keep unrelated archived evidence", "todo_archived", "Todo Archive"), + ("Keep active", "todo_active", "Agent Todo"), +]) +def test_source_text_is_not_an_attention_summary(tmp_path, original, todo_id, heading): + fixture = cold_workspace(tmp_path) + full_text = "Preserve " + "explicit recovery requirement " * 70 + fixture.state.write_text(fixture.state.read_text().replace(original, full_text)) + inventory = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + archive = next(t for t in inventory["projection"]["todos"] if t["todo_id"] == todo_id) + assert archive["text"] == full_text.strip() + assert archive["source_section"] == heading + + +def test_stale_source_is_rejected_by_native_owner(tmp_path): + from loopx.control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot + from loopx.control_plane.coordination.local_authority_shadow_projection import source_effect_runtime_result + + fixture = cold_workspace(tmp_path) + goal = json.loads(fixture.registry.read_text())["goals"][0] + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=fixture.runtime, + state_path=fixture.state, registry_path=fixture.registry, include_all_archived_todos=True) + fixture.state.write_text(fixture.state.read_text() + "\nNew source bytes\n") + result = source_effect_runtime_result("coordination.source.inspect", { + "schema_version": "loopx_cold_source_inspection_request_v0", "goal_id": fixture.goal, + "runtime_root": str(fixture.runtime), "projection": projection, "source_snapshot": snapshot}) + assert result["status"] == "failed" and result["reason_code"] == "source_changed_retry" + assert result["import_ready"] is False + + +def test_selected_canonical_source_is_not_reinterpreted_as_markdown(tmp_path): + import hashlib + + fixture = cold_workspace(tmp_path) + marker = fixture.runtime / "authority" / ("provider-" + hashlib.sha256(fixture.goal.encode()).hexdigest() + ".json") + marker.parent.mkdir(parents=True) + marker.write_text('{"provider":"sqlite"}') + result = fixture.cli("coordination-shadow", "inspect-source", success=False) + assert result["source_inventory"]["reason_code"] == "cold_source_canonical_authority_present" + assert marker.read_text() == '{"provider":"sqlite"}' + + +def test_existing_capture_still_excludes_unreferenced_archives(tmp_path): + from loopx.control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot + + fixture = cold_workspace(tmp_path) + goal = json.loads(fixture.registry.read_text())["goals"][0] + projection, _ = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=fixture.runtime, + state_path=fixture.state, registry_path=fixture.registry) + assert [t["todo_id"] for t in projection["todos"]] == ["todo_active"] + + +def test_expired_orphan_active_lease_still_requires_settlement(tmp_path): + fixture = cold_workspace(tmp_path) + directory = fixture.runtime / "goals" / fixture.goal / "task-leases" + directory.mkdir(parents=True) + lease = {"schema_version": "task_lease_v0", "goal_id": fixture.goal, + "todo_id": "removed", "owner": "agent-a", "status": "active", + "idempotency_key": "old-work", "version": 4, "lease_epoch": 2, + "expires_at": "2000-01-01T00:00:00Z", "write_scopes": []} + path = directory / "removed.json" + path.write_text(json.dumps(lease)) + before = path.read_bytes() + inventory = fixture.cli("coordination-shadow", "inspect-source")["source_inventory"] + assert inventory["lease_file_count"] == 1 + assert inventory["leases_requiring_settlement"] == ["removed"] + assert inventory["retained_leases"] == [lease] + assert inventory["projection"]["leases"] == [] + assert inventory["import_ready"] is False + assert path.read_bytes() == before + + +@pytest.mark.parametrize("corruption", ["duplicate_archive", "unknown_archive_role", "invalid_orphan_lease"]) +def test_ambiguous_or_unsupported_history_is_not_an_empty_source(tmp_path, corruption): + fixture = cold_workspace(tmp_path) + if corruption == "duplicate_archive": + fixture.state.write_text(fixture.state.read_text().replace("todo_id=todo_archived", "todo_id=todo_active")) + elif corruption == "unknown_archive_role": + fixture.state.write_text(fixture.state.read_text().replace("todo_id=todo_archived role=agent", "todo_id=todo_archived")) + else: + directory = fixture.runtime / "goals" / fixture.goal / "task-leases" + directory.mkdir(parents=True) + (directory / "removed.json").write_text(json.dumps({"goal_id": fixture.goal, "todo_id": "removed"})) + result = fixture.cli("coordination-shadow", "inspect-source", success=False) + assert result["ok"] is False + assert not (fixture.runtime / "authority").exists() + + +def test_source_inspection_has_no_execute_switch(tmp_path): + fixture = cold_workspace(tmp_path) + result = subprocess.run([sys.executable, "-m", "loopx.cli", + *fixture.arguments("coordination-shadow", "inspect-source", "--execute")], + capture_output=True, text=True) + assert result.returncode == 2 and "unrecognized arguments" in result.stderr + + +@pytest.mark.parametrize("kind", ["unsupported_name", "directory", "symlink"]) +def test_unsafe_lease_inventory_is_not_silently_omitted(tmp_path, kind): + fixture = cold_workspace(tmp_path) + directory = fixture.runtime / "goals" / fixture.goal / "task-leases" + directory.mkdir(parents=True) + if kind == "unsupported_name": + (directory / "历史.json").write_text("{}") + elif kind == "directory": + (directory / "todo_removed.json").mkdir() + else: + target = tmp_path / "retained.json" + target.write_text("{}") + (directory / "todo_removed.json").symlink_to(target) + result = fixture.cli("coordination-shadow", "inspect-source", success=False) + assert result["ok"] is False + assert result["error"] == "source_lease_inventory_invalid" + assert not (fixture.runtime / "authority").exists() + assert not (fixture.runtime / "authority-shadow").exists()