From 5165626e803d1ee2993187dcde34e7a5a5c2352a Mon Sep 17 00:00:00 2001 From: Steven <96154058+steven-kid@users.noreply.github.com> Date: Tue, 8 Sep 2026 21:05:08 +0800 Subject: [PATCH 1/3] fix(coordination): retain archive receipts until projection delivery Signed-off-by: Steven <96154058+steven-kid@users.noreply.github.com> --- .../todo-archive-delivery-recovery.md | 78 +++++ .../coordination/local_archive_attempt.ts | 170 ++++++++++ .../coordination/local_authority_runtime.ts | 58 +++- .../coordination/todo_terminal_lifecycle.ts | 33 +- .../control_plane/effect_runtime_handlers.ts | 2 + .../todos/provider_terminal_lifecycle.py | 34 +- .../test_archive_retry_delivery.py | 319 ++++++++++++++++++ .../test_local_coordination_authority.py | 1 + .../authority_store_conformance.ts | 40 +++ .../local_archive_attempt.test.ts | 241 +++++++++++++ .../shadow_native_writer_boundary.test.ts | 20 +- tsconfig.control-plane.json | 1 + 12 files changed, 985 insertions(+), 12 deletions(-) create mode 100644 docs/architecture/todo-archive-delivery-recovery.md create mode 100644 loopx/control_plane/coordination/local_archive_attempt.ts create mode 100644 tests/control_plane/test_archive_retry_delivery.py create mode 100644 tests/control_plane_ts/local_archive_attempt.test.ts diff --git a/docs/architecture/todo-archive-delivery-recovery.md b/docs/architecture/todo-archive-delivery-recovery.md new file mode 100644 index 0000000000..f657074427 --- /dev/null +++ b/docs/architecture/todo-archive-delivery-recovery.md @@ -0,0 +1,78 @@ +# Todo archive delivery recovery + +This contract hardens the promoted archive transaction introduced by +[#4053](https://github.com/huangruiteng/loopx/pull/4053), under the +[TypeScript control-plane migration RFC](./rfcs/typescript-control-plane-migration-v0.md). +Archive selection, standing-decision retention, CAS, and historical receipts +remain owned by the native transaction. + +## Pending delivery and retry identity + +The local TypeScript adapter durably records one pending archive attempt per +goal and role before executing the transaction. The record binds the operation +ID, retention limit, observed provider revision, and authority store identity. +It is correlation state, not a second Todo authority or a copy of the receipt. +The existing per-goal maintenance lock serializes attempt creation and retirement. + +While that attempt is pending, a repeated archive request for the same role and +limit reuses its identity. Receipt replay precedes current-head revision checks +and archive selection, so a committed result survives a process exit before +Markdown projection, a projection write failure, and unrelated head advancement. +The replay reports the original receipt, moved IDs, count, and revision without +creating another canonical commit. A different retention limit is rejected +until the pending delivery is settled. + +Fresh Python requests bind the revision they actually observed. If that head +changes before selection, the native transaction rejects the stale request. +Older v0 callers that omit the optional revision retain their original receipt +hash. A fresh attempt preserves an explicitly supplied revision; a pending +retry keeps the original attempt's binding rather than today's observed head. + +After the Python projection adapter delivers or verifies the current canonical +view, it acknowledges the exact operation through +`coordination.local_authority.todo_archive_ack`. The native owner checks the +receipt and store identity before retiring the attempt. A stale acknowledgement +cannot retire a newer attempt. Transport failure preserves the already committed +result; uncertain transaction outcomes retain the pending identity for retry. +Proven pre-commit rejection releases it so a corrected request can proceed. + +Preview bypasses pending attempts and historical receipt replay and performs no +writes. A fresh empty archive creates no canonical receipt and does not advance +the provider revision. Local correlation storage is bounded to two slots per +goal, independent of archive history size. + +## Delivery boundary and migration cost + +The recovery boundary ends at successful projection acknowledgement. A later +identical CLI invocation is a new archive request and can process a new batch. +This does not promise recovery of a CLI response lost after acknowledgement or +discover attempts made before correlation records existed. Callers requiring +that stronger guarantee need an explicit caller-retained request identity. + +A changed promoted archive uses four Python/TypeScript request-responses: +authority read, archive transaction, projection readback, and acknowledgement. +The previous path used three; empty and preview paths add no acknowledgement. +The extra durable correlation writes and acknowledgement are a correctness cost, +not a claimed migration speedup. Archive domain rules stay in TypeScript. +Delete the Python acknowledgement crossing with the terminal lifecycle facade +when the top-level Todo CLI and projection consumer run in TypeScript; retain +the provider receipt and independent crash/retry conformance coverage. + +## 中文契约 + +本变更修复 promoted archive 在 canonical commit 后、Markdown 投影交付前失败时的 +重试身份丢失。TypeScript 在执行前持久记录每个 goal/role 的未交付 operation ID、 +保留数量、观察到的 provider revision 和 store identity。相同 role/limit 的重试先 +恢复原 receipt,再考虑当前 head;不会重复提交,也不会把原归档数量错误地报告为零。 +未完成交付时修改保留数量会被拒绝。 + +新 Python 请求绑定实际读到的 revision;无历史 receipt 且 head 已变化时拒绝执行。 +旧 v0 请求省略 revision 时保持原有请求 hash。预览不读取未交付记录、不重放历史结果、 +不写状态;空归档不推进 canonical revision。确认仅在投影交付成功后进行,且必须匹配 +原 operation、receipt 和 store。过期确认不能清除新操作;结果不确定时保留重试身份, +确定未提交的拒绝允许后续修正请求继续。 + +恢复保证截至投影确认成功;之后相同 CLI 调用代表新一批归档,不承诺恢复确认之后才 +丢失的 stdout,也不追溯发现旧版本未记录的调用。发生归档的路径从三次跨运行时调用 +增加为四次,额外确认和持久化是明确的正确性成本。顶层 Todo CLI 与投影消费者迁入 +TypeScript 后删除 Python 确认桥接,保留 provider receipt 和独立崩溃恢复测试。 diff --git a/loopx/control_plane/coordination/local_archive_attempt.ts b/loopx/control_plane/coordination/local_archive_attempt.ts new file mode 100644 index 0000000000..b5a09655af --- /dev/null +++ b/loopx/control_plane/coordination/local_archive_attempt.ts @@ -0,0 +1,170 @@ +import { readFile } from "node:fs/promises"; +import { join } from "node:path"; + +import type { JsonObject } from "../effect_program.ts"; +import { durableWriteJson } from "../effect_runtime_io.ts"; +import type { AuthorityStore } from "./authority_store.ts"; +import { + canonicalAuthorityObject, + canonicalAuthoritySha256, + hasExactAuthorityKeys, + requireAuthorityStoreId, +} from "./authority_store_codec.ts"; +import { shadowManagementDirectory } from "./shadow_management.ts"; +import { + COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA, + executeCoordinationTodoArchiveCompleted, + type CoordinationTodoArchiveInput, +} from "./todo_terminal_lifecycle.ts"; + +const ATTEMPT_SCHEMA = "loopx_local_todo_archive_attempt_v0"; +export const LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA = + "loopx_local_coordination_todo_archive_ack_result_v0"; + +type Role = CoordinationTodoArchiveInput["role"]; +interface ArchiveAttempt { + operation_id: string; + expected_provider_revision: string | null; + max_active_done: number; + store_identity: string; +} + +function attemptPath(root: string, goalId: string, role: Role): string { + return join(shadowManagementDirectory(root, goalId), `archive-${role}.json`); +} + +async function readAttempt(root: string, goalId: string, role: Role): Promise { + let value: unknown; + try { + value = JSON.parse(await readFile(attemptPath(root, goalId, role), "utf8")); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; + throw error; + } + const slot = canonicalAuthorityObject(value, "archive attempt slot"); + if (!hasExactAuthorityKeys(slot, ["schema_version", "goal_id", "role", "attempt"]) || + slot.schema_version !== ATTEMPT_SCHEMA || slot.goal_id !== goalId || slot.role !== role) { + throw new Error("archive attempt slot identity mismatch"); + } + if (slot.attempt === null) return null; + const attempt = canonicalAuthorityObject(slot.attempt, "archive attempt"); + if (!hasExactAuthorityKeys(attempt, ["operation_id", "expected_provider_revision", + "max_active_done", "store_identity"]) || + !Number.isSafeInteger(attempt.max_active_done) || Number(attempt.max_active_done) < 0) { + throw new Error("invalid archive attempt"); + } + return { + operation_id: requireAuthorityStoreId(attempt.operation_id, "archive operation id"), + expected_provider_revision: attempt.expected_provider_revision === null ? null : + requireAuthorityStoreId(attempt.expected_provider_revision, "archive expected provider revision"), + max_active_done: Number(attempt.max_active_done), + store_identity: requireAuthorityStoreId(attempt.store_identity, "archive store identity"), + }; +} + +async function writeAttempt( + root: string, goalId: string, role: Role, attempt: ArchiveAttempt | null, +): Promise { + // A durable null tombstone avoids a delete/parent-fsync recovery window. + // There are at most two bounded slots per goal, regardless of history size. + await durableWriteJson(attemptPath(root, goalId, role), { + schema_version: ATTEMPT_SCHEMA, goal_id: goalId, role, + attempt: attempt === null ? null : {...attempt}, + }); +} + +/** Caller holds the goal's shadow-maintenance mutex throughout this operation. */ +export async function executeLocalArchiveAttempt( + store: AuthorityStore, root: string, input: CoordinationTodoArchiveInput, +): Promise { + if (input.dry_run) return executeCoordinationTodoArchiveCompleted(store, input); + let attempt = await readAttempt(root, input.goal_id, input.role); + if (attempt === null) { + // A retained v0 receipt may predate observed-head binding. Preserve the + // actual caller fields; adding today's revision would change its identity. + const receipt = await store.readReceipt(input.operation_id); + if (receipt.status === "missing") { + const head = await store.loadAuthority(); + if (head.status !== "loaded") return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...head, changed: false, + }; + } else if (receipt.status !== "found") return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...receipt, changed: false, + }; + } + const identity = await store.storeIdentity(); + if (identity.status !== "available") return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...identity, changed: false, + }; + if (attempt !== null && attempt.store_identity !== identity.store_identity) { + throw new Error("archive attempt belongs to a different authority store"); + } + if (attempt !== null && attempt.max_active_done !== input.max_active_done) return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + status: "conflict", changed: false, conflict_kind: "archive_attempt_pending", + operation_id: attempt.operation_id, + reason: "complete the pending archive projection before changing the retention limit", + }; + if (attempt === null) { + attempt = { + operation_id: input.operation_id, + expected_provider_revision: input.expected_provider_revision ?? null, + max_active_done: input.max_active_done, store_identity: identity.store_identity, + }; + // Persist correlation before the provider can commit. The provider alone + // owns selection, request identity, CAS, and historical result receipts. + await writeAttempt(root, input.goal_id, input.role, attempt); + } + const result = await executeCoordinationTodoArchiveCompleted(store, { + ...input, operation_id: attempt.operation_id, + expected_provider_revision: attempt.expected_provider_revision ?? undefined, + }); + const identityRejected = result.status === "failed" && + result.reason_code === "coordination_operation_identity_mismatch"; + const beforeCommitRejected = result.status === "conflict" || ( + result.status === "failed" && ["invalid_coordination_todo_archive", + "invalid_coordination_projection"].includes(String(result.reason_code)) + ); + if (result.status === "no_change" || identityRejected || (beforeCommitRejected && + (await store.readReceipt(attempt.operation_id)).status === "missing")) { + // No committed mutation exists to project. In particular, retain the + // established zero-revision-change semantics of no-change archive calls. + await writeAttempt(root, input.goal_id, input.role, null); + } + return {...result, operation_id: attempt.operation_id}; +} + +/** A successful projection acknowledges exactly one attempt, never a newer one. */ +export async function acknowledgeLocalArchiveAttempt( + store: AuthorityStore, root: string, goalId: string, role: Role, operationId: string, +): Promise { + const result = {schema_version: LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA, + goal_id: goalId, role, operation_id: operationId}; + const attempt = await readAttempt(root, goalId, role); + if (attempt === null) return {...result, status: "no_change", changed: false}; + if (attempt.operation_id !== operationId) return {...result, status: "stale", changed: false}; + const identity = await store.storeIdentity(); + if (identity.status !== "available") return {...result, ...identity, changed: false}; + if (attempt.store_identity !== identity.store_identity) { + throw new Error("archive acknowledgement belongs to a different authority store"); + } + const receipt = await store.readReceipt(operationId); + const original = receipt.status === "found" ? receipt.receipts[0] : undefined; + if (receipt.status !== "found" || receipt.receipts.length !== 1 || + original?.schema_version !== COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA || + original.operation_id !== operationId || original.goal_id !== goalId || + original.request_sha256 !== canonicalAuthoritySha256({ + goal_id: goalId, role, max_active_done: attempt.max_active_done, + dry_run: false, + ...(attempt.expected_provider_revision === null ? {} : { + expected_provider_revision: attempt.expected_provider_revision, + }), + })) return { + ...result, status: "failed", changed: false, + reason_code: "archive_ack_receipt_unavailable", + reason: "archive projection acknowledgement requires its committed receipt", + }; + await writeAttempt(root, goalId, role, null); + return {...result, status: "acknowledged", changed: true}; +} diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index aa70dee4bf..4b1b192624 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -55,9 +55,13 @@ import { import { COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, - executeCoordinationTodoArchiveCompleted, executeCoordinationTodoTerminalLifecycle, } from "./todo_terminal_lifecycle.ts"; +import { + acknowledgeLocalArchiveAttempt, + executeLocalArchiveAttempt, + LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA, +} from "./local_archive_attempt.ts"; import { editCoordinationTodo, TODO_COMPATIBILITY_EDIT_RESULT_SCHEMA } from "./todo_compatibility_edit.ts"; import { normalizeIdempotencyKey, @@ -73,6 +77,8 @@ export const LOCAL_COORDINATION_TODO_TERMINAL_LIFECYCLE_REQUEST_SCHEMA = "loopx_local_coordination_todo_terminal_lifecycle_request_v0"; export const LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA = "loopx_local_coordination_todo_archive_request_v0"; +export const LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA = + "loopx_local_coordination_todo_archive_ack_request_v0"; export { LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA, LOCAL_COORDINATION_MUTATION_RESULT_SCHEMA, @@ -155,6 +161,11 @@ function requiredNonNegativeSafeInteger(value: unknown, label: string): number { return value as number; } +function archiveRole(value: unknown): "agent" | "user" { + if (value !== "agent" && value !== "user") throw new TypeError("unsupported archive role"); + return value; +} + function optionalNonNegativeSafeInteger(value: unknown, label: string): number | null { return value === null || value === undefined ? null @@ -882,16 +893,23 @@ export async function archiveLocalCoordinationTodos( input.max_active_done, "max_active_done", ); + const role = archiveRole(input.role); + const operationId = requireAuthorityStoreId(input.operation_id, "operation id"); + const expectedRevision = input.expected_provider_revision === undefined ? undefined : + requireAuthorityStoreId(input.expected_provider_revision, "expected provider revision"); + if (typeof input.dry_run !== "boolean") throw new TypeError("dry_run must be a boolean"); + const now = claimObservedAt(input.observed_at); return await withCanonicalWriter(root, goalId, input.dry_run === true, async () => { const store = dependencies.createStore?.(authorityDirectory(root), goalId) ?? new FileAuthorityStore(authorityDirectory(root), goalId); - return {...await executeCoordinationTodoArchiveCompleted(store, { + return {...await executeLocalArchiveAttempt(store, root, { goal_id: goalId, - role: requireAuthorityStoreId(input.role, "role") as "agent" | "user", + role, max_active_done: maxActiveDone, - operation_id: requireAuthorityStoreId(input.operation_id, "operation id"), + operation_id: operationId, + expected_provider_revision: expectedRevision, dry_run: input.dry_run as boolean, - now: claimObservedAt(input.observed_at), + now, }), ...providerEvidence}; }); } catch (error) { @@ -904,6 +922,36 @@ export async function archiveLocalCoordinationTodos( } } +/** Retire one local retry identity only after its compatibility projection succeeds. */ +export async function acknowledgeLocalCoordinationTodoArchive( + value: unknown, + dependencies: LocalAuthorityRuntimeDependencies = {}, +): Promise { + try { + const input = requireJsonObject(value, "local Todo archive acknowledgement"); + if (input.schema_version !== LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA) { + throw new TypeError("local Todo archive acknowledgement schema mismatch"); + } + const root = runtimeRoot(input.runtime_root); + const goalId = requireAuthorityStoreId(input.goal_id, "goal id"); + const role = archiveRole(input.role); + const operationId = requireAuthorityStoreId(input.operation_id, "operation id"); + return await withCanonicalWriter(root, goalId, false, async () => { + const store = dependencies.createStore?.(authorityDirectory(root), goalId) ?? + new FileAuthorityStore(authorityDirectory(root), goalId); + return acknowledgeLocalArchiveAttempt(store, root, goalId, role, operationId); + }); + } catch (error) { + return { + schema_version: LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA, + status: "failed", changed: false, + reason_code: error instanceof ShadowManagementError ? error.reason_code : + "invalid_local_coordination_todo_archive_ack_request", + reason: error instanceof Error ? error.message : "invalid archive acknowledgement", + }; + } +} + /** Embedded file adapter; no Markdown input or projection write is accepted. */ export async function editLocalCoordinationTodo( value: unknown, diff --git a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts index 5993bdcacf..43b75a49a2 100644 --- a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts +++ b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts @@ -109,6 +109,7 @@ export interface CoordinationTodoArchiveInput { readonly role: TodoRole; readonly max_active_done: number; readonly operation_id: string; + readonly expected_provider_revision?: string; readonly dry_run: boolean; readonly now: Date; } @@ -1107,6 +1108,11 @@ function normalizeArchiveInput(raw: CoordinationTodoArchiveInput): CoordinationT goal_id: requireAuthorityStoreId(raw.goal_id, "goal id"), role: requireLiteral(raw.role, TODO_ROLES, "role"), operation_id: requireAuthorityStoreId(raw.operation_id, "operation id"), + ...(raw.expected_provider_revision === undefined ? {} : { + expected_provider_revision: requireAuthorityStoreId( + raw.expected_provider_revision, "expected provider revision", + ), + }), dry_run: requireBoolean(raw.dry_run, "dry_run"), now: requireDate(raw.now, "now"), }; @@ -1136,6 +1142,7 @@ function replayArchive( return { ...result, schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + operation_id: input.operation_id, status, changed: status !== "replayed" && result.changed === true, provider_revision: receipt.provider_revision, @@ -1165,15 +1172,32 @@ export async function executeCoordinationTodoArchiveCompleted( role: input.role, max_active_done: input.max_active_done, dry_run: input.dry_run, + ...(input.expected_provider_revision === undefined ? {} : { + expected_provider_revision: input.expected_provider_revision, + }), }); - const replay = replayArchive( - await store.readReceipt(input.operation_id), input, requestSha, "replayed", - ); - if (replay !== null) return replay; + // Preview observes the current snapshot without consuming or replaying a + // durable operation identity. Historical receipts precede current-head CAS. + if (!input.dry_run) { + const replay = replayArchive( + await store.readReceipt(input.operation_id), input, requestSha, "replayed", + ); + if (replay !== null) return replay; + } const head = await store.loadAuthority(); if (head.status !== "loaded") { return {schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...head, changed: false}; } + if (input.expected_provider_revision !== undefined && + head.provider_revision !== input.expected_provider_revision) { + return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + status: "conflict", changed: false, + conflict_kind: "provider_revision_mismatch", + current_provider_revision: head.provider_revision, + current_cursor: head.cursor, + }; + } let projection: ReturnType; try { projection = indexCoordinationProjection(head.head, input.goal_id); @@ -1193,6 +1217,7 @@ export async function executeCoordinationTodoArchiveCompleted( const updatedAt = input.now.toISOString().replace(/\.\d{3}Z$/u, "Z"); const result: JsonObject = { role: selection.role, + operation_id: input.operation_id, changed: moved.length > 0, active_done_before: selection.active_done_before, active_done_after: selection.active_done_after, diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 30e8a86711..16a69a4741 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -115,6 +115,7 @@ import { rollbackCoordinationRuntimeShadow, } from "./coordination/runtime_shadow.ts"; import { + acknowledgeLocalCoordinationTodoArchive, archiveLocalCoordinationTodos, claimLocalCoordinationTodo, createLocalCoordinationTodo, @@ -430,6 +431,7 @@ export function createEffectRuntimeHandlers( ["coordination.local_authority.todo_update", updateLocalCoordinationTodo], ["coordination.local_authority.todo_terminal", terminalLifecycleLocalCoordinationTodo], ["coordination.local_authority.todo_archive", archiveLocalCoordinationTodos], + ["coordination.local_authority.todo_archive_ack", acknowledgeLocalCoordinationTodoArchive], ["coordination.local_authority.todo_compatibility_edit", editLocalCoordinationTodo], ["coordination.local_authority.mutate", mutateLocalCoordinationAuthority], ["coordination.local_authority.todo_read", readLocalCoordinationTodo], diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index 6126ad7f56..2a989308f8 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -40,6 +40,7 @@ _TERMINAL_REQUEST_SCHEMA = "loopx_local_coordination_todo_terminal_lifecycle_request_v0" _ARCHIVE_REQUEST_SCHEMA = "loopx_local_coordination_todo_archive_request_v0" +_ARCHIVE_ACK_REQUEST_SCHEMA = "loopx_local_coordination_todo_archive_ack_request_v0" _ACCEPTED = {"applied", "recovered", "replayed", "no_change", "planned"} @@ -532,6 +533,7 @@ def archive_canonical_todos_if_promoted( max_active_done=max_active_done, provider_revision=provider_revision, ), + "expected_provider_revision": provider_revision, "dry_run": dry_run, "observed_at": now_local(), }, @@ -545,7 +547,7 @@ def archive_canonical_todos_if_promoted( ), payload=payload, ) - return _projection_payload( + response = _projection_payload( settle_canonical_todo_projection( {"ok": True, "dry_run": dry_run, "goal_id": goal_id, **dict(result)}, registry_path=registry_path, @@ -555,6 +557,36 @@ def archive_canonical_todos_if_promoted( state_file=state_file, ) ) + if ( + not dry_run + and response.get("moved_count", 0) > 0 + and response.get("projection_delivery") in {"delivered", "current"} + ): + # The native owner retains the attempt until its external projection + # provider succeeds. An ACK failure must preserve the committed result + # and leave the same attempt available for the next retry. + try: + acknowledgement = effect_runtime_result( + "coordination.local_authority.todo_archive_ack", + { + "schema_version": _ARCHIVE_ACK_REQUEST_SCHEMA, + "runtime_root": str(runtime_root.expanduser().resolve(strict=False)), + "goal_id": goal_id, + "role": role, + "operation_id": response.get("operation_id"), + }, + ) + response["archive_delivery_ack"] = ( + dict(acknowledgement) + if isinstance(acknowledgement, Mapping) + else {"status": "pending", "reason_code": "invalid_archive_ack_result"} + ) + except Exception as error: # noqa: BLE001 - the canonical commit already landed + response["archive_delivery_ack"] = { + "status": "pending", "reason_code": "archive_ack_unavailable", + "error_class": error.__class__.__name__, "retryable": True, + } + return response __all__ = [ diff --git a/tests/control_plane/test_archive_retry_delivery.py b/tests/control_plane/test_archive_retry_delivery.py new file mode 100644 index 0000000000..c847216ee1 --- /dev/null +++ b/tests/control_plane/test_archive_retry_delivery.py @@ -0,0 +1,319 @@ +"""Archive recovery through the public CLI and a real canonical file store.""" + +from __future__ import annotations + +import json +from pathlib import Path +import subprocess +import sys + +import pytest +from canonical_authority_fixture import initialize_canonical_authority + +from loopx.control_plane.coordination.local_authority import ( + LocalCoordinationAuthorityUnavailable, + read_canonical_todos_if_promoted, +) +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) +from loopx.control_plane.todos import provider_projection, provider_terminal_lifecycle +from loopx.control_plane.todos.active_state_editing import TODO_SECTION_HEADINGS +from loopx.todos import archive_completed_todos, complete_goal_todo, update_goal_todo + + +REPOSITORY = Path(__file__).resolve().parents[2] +CRASH_BEFORE_PROJECTION = """ +import json, os, sys +from pathlib import Path +from loopx.control_plane.todos import provider_terminal_lifecycle +from loopx.cli import main +evidence = Path(sys.argv.pop(1)) +def lose_response(payload, **kwargs): + evidence.write_text(json.dumps(payload), encoding='utf-8') + os._exit(73) +provider_terminal_lifecycle.settle_canonical_todo_projection = lose_response +raise SystemExit(main(sys.argv[1:])) +""" + + +def _fixture(tmp_path: Path) -> tuple[Path, Path, Path]: + runtime = tmp_path / "runtime" + project = tmp_path / "project" + project.mkdir() + state = project / "ACTIVE_GOAL_STATE.md" + state.write_text( + "# Goal\n\n## User Todo / Owner Review Reading Queue\n\n" + "## Agent Todo\n\n## Completed Work Archive\n", + encoding="utf-8", + ) + registry = tmp_path / "registry.json" + registry.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(runtime), + "goals": [ + { + "id": "archive-goal", + "repo": str(project), + "state_file": state.name, + "coordination": {"registered_agents": ["agent-a"]}, + } + ], + } + ), + encoding="utf-8", + ) + projection = build_todo_runtime_shadow_projection( + goal_id="archive-goal", + todos=[ + { + "schema_version": "todo_item_v0", + "done": True, + "text": "A completed delivery awaiting archival", + "todo_id": "todo_completed", + "role": "agent", + "status": "done", + "archive_state": "active", + "source_section": TODO_SECTION_HEADINGS["agent"], + "task_class": "advancement_task", + "claimed_by": "agent-a", + }, + { + "schema_version": "todo_item_v0", + "done": False, + "text": "The next delivery", + "todo_id": "todo_next", + "role": "agent", + "status": "open", + "archive_state": "active", + "source_section": TODO_SECTION_HEADINGS["agent"], + "task_class": "advancement_task", + "claimed_by": "agent-a", + }, + ], + handoff_mode="soft_claim", + ) + initialize_canonical_authority( + runtime, "archive-goal", projection, state_path=state + ) + return registry, runtime, state + + +def _archive_cli( + registry: Path, + *, + crash_evidence: Path | None = None, + execute: bool = True, + maximum: int = 0, +) -> subprocess.CompletedProcess[str]: + launcher = ( + [sys.executable, "-c", CRASH_BEFORE_PROJECTION, str(crash_evidence)] + if crash_evidence + else [sys.executable, "-m", "loopx.cli"] + ) + return subprocess.run( + [ + *launcher, + "--format", + "json", + "--registry", + str(registry), + "todo", + "archive-completed", + "--goal-id", + "archive-goal", + "--role", + "agent", + "--max-active-done", + str(maximum), + *(["--execute"] if execute else []), + ], + cwd=REPOSITORY, + capture_output=True, + text=True, + check=False, + timeout=45, + ) + + +def test_archive_cli_recovers_committed_result_after_process_exit( + tmp_path: Path, +) -> None: + registry, runtime, state = _fixture(tmp_path) + original_markdown = state.read_bytes() + evidence = tmp_path / "committed-result.json" + lost = _archive_cli(registry, crash_evidence=evidence) + assert lost.returncode == 73, lost.stderr + committed = json.loads(evidence.read_text(encoding="utf-8")) + assert committed["status"] == "applied" + assert committed["moved_count"] == 1 + assert committed["moved_todo_ids"] == ["todo_completed"] + assert state.read_bytes() == original_markdown + + retry = _archive_cli(registry) + assert retry.returncode == 0, retry.stderr or retry.stdout + recovered = json.loads(retry.stdout) + assert recovered["status"] == "replayed" + assert recovered["moved_count"] == 1 + assert recovered["moved_todo_ids"] == committed["moved_todo_ids"] + assert recovered["original_receipt"] == committed["original_receipt"] + assert recovered["provider_revision"] == committed["provider_revision"] + assert recovered["changed"] is False + assert recovered["projection_delivery"] in {"delivered", "current"} + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime, goal_id="archive-goal" + ) + assert canonical is not None + assert canonical["provider_revision"] == committed["provider_revision"] + assert canonical["todos"][0]["archive_state"] == "archive" + + # A delivered operation must not become the permanent identity for all + # future archive calls with the same role/retention arguments. + fresh = _archive_cli(registry) + assert fresh.returncode == 0, fresh.stderr + no_change = json.loads(fresh.stdout) + assert no_change["status"] == "no_change" + assert no_change["moved_count"] == 0 + assert no_change["provider_revision"] == committed["provider_revision"] + + completed = complete_goal_todo( + registry_path=registry, + goal_id="archive-goal", + todo_id="todo_next", + agent_id="agent-a", + next_agent_todo="Continue after the next delivery", + next_claimed_by="agent-a", + ) + assert completed["completed"] is True + next_batch = _archive_cli(registry) + assert next_batch.returncode == 0, next_batch.stderr or next_batch.stdout + next_result = json.loads(next_batch.stdout) + assert next_result["status"] == "applied" + assert next_result["moved_todo_ids"] == ["todo_next"] + assert ( + next_result["original_receipt"]["operation_id"] + != committed["original_receipt"]["operation_id"] + ) + + +def test_pending_projection_and_preview_preserve_archive_retry( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry, runtime, state = _fixture(tmp_path) + original_markdown = state.read_text(encoding="utf-8") + + def reject_projection(*_args, **_kwargs): + raise OSError("injected projection write failure") + + with monkeypatch.context() as patch: + patch.setattr(provider_projection, "_atomic_write_text", reject_projection) + committed = archive_completed_todos( + registry_path=registry, + goal_id="archive-goal", + max_active_done=0, + dry_run=False, + ) + assert committed["status"] == "applied" + assert committed["projection_delivery"] == "pending" + + preview = _archive_cli(registry, execute=False) + assert preview.returncode == 0, preview.stderr + assert json.loads(preview.stdout)["moved_count"] == 0 + assert state.read_text(encoding="utf-8") == original_markdown + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime, goal_id="archive-goal" + ) + assert canonical is not None + assert canonical["provider_revision"] == committed["provider_revision"] + + different = _archive_cli(registry, maximum=1) + assert different.returncode != 0 + retry = _archive_cli(registry) + assert retry.returncode == 0, retry.stderr or retry.stdout + replay = json.loads(retry.stdout) + assert replay["status"] == "replayed" + assert replay["original_receipt"] == committed["original_receipt"] + assert replay["projection_delivery"] in {"delivered", "current"} + + +def test_archive_rejects_snapshot_drift_before_selection( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry, runtime, _ = _fixture(tmp_path) + read = provider_terminal_lifecycle.read_canonical_todos_if_promoted + + def read_then_edit(**kwargs): + observed = read(**kwargs) + update_goal_todo( + registry_path=registry, + goal_id="archive-goal", + todo_id="todo_next", + agent_id="agent-a", + note="An independently committed correction", + ) + return observed + + monkeypatch.setattr( + provider_terminal_lifecycle, "read_canonical_todos_if_promoted", read_then_edit + ) + with pytest.raises(LocalCoordinationAuthorityUnavailable) as rejected: + archive_completed_todos( + registry_path=registry, + goal_id="archive-goal", + max_active_done=0, + dry_run=False, + ) + assert rejected.value.payload["status"] == "conflict" + canonical = read(runtime_root=runtime, goal_id="archive-goal") + assert canonical is not None + assert all(todo["archive_state"] == "active" for todo in canonical["todos"]) + monkeypatch.setattr( + provider_terminal_lifecycle, "read_canonical_todos_if_promoted", read + ) + retry = _archive_cli(registry) + assert retry.returncode == 0, retry.stderr or retry.stdout + assert json.loads(retry.stdout)["moved_todo_ids"] == ["todo_completed"] + + +def test_archive_ack_transport_failure_preserves_committed_result( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry, runtime, _ = _fixture(tmp_path) + execute = provider_terminal_lifecycle.effect_runtime_result + + def unavailable_ack(method, params): + if method == "coordination.local_authority.todo_archive_ack": + raise OSError("injected acknowledgement transport failure") + return execute(method, params) + + with monkeypatch.context() as patch: + patch.setattr( + provider_terminal_lifecycle, "effect_runtime_result", unavailable_ack + ) + committed = archive_completed_todos( + registry_path=registry, + goal_id="archive-goal", + max_active_done=0, + dry_run=False, + ) + assert committed["status"] == "applied" + assert committed["moved_todo_ids"] == ["todo_completed"] + assert committed["projection_delivery"] in {"delivered", "current"} + assert committed["archive_delivery_ack"]["status"] == "pending" + assert committed["archive_delivery_ack"]["retryable"] is True + retry = _archive_cli(registry) + assert retry.returncode == 0, retry.stderr or retry.stdout + recovered = json.loads(retry.stdout) + assert recovered["status"] == "replayed" + assert recovered["original_receipt"] == committed["original_receipt"] + assert recovered["archive_delivery_ack"]["status"] == "acknowledged" + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime, goal_id="archive-goal" + ) + assert canonical is not None + assert canonical["provider_revision"] == committed["provider_revision"] diff --git a/tests/control_plane/test_local_coordination_authority.py b/tests/control_plane/test_local_coordination_authority.py index de9bd1d4d8..4925baf4b7 100644 --- a/tests/control_plane/test_local_coordination_authority.py +++ b/tests/control_plane/test_local_coordination_authority.py @@ -1345,6 +1345,7 @@ def count_authority_runtime_call( "coordination.local_authority.todo_list", "coordination.local_authority.todo_archive", "coordination.local_authority.todo_list", + "coordination.local_authority.todo_archive_ack", ] canonical_after_archive = read_canonical_todos_if_promoted( runtime_root=runtime_root, diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 22e7156896..61c3a4eb47 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -547,6 +547,46 @@ export function registerAuthorityStoreConformance( assert.equal(afterNoChange.cursor, beforeNoChange.cursor); }); + test(`${providerName} conformance: archive binds the observed head and replays before current eligibility`, async (t) => { + const {store} = await factory(t); + const goalId = "goal-archive-retry"; + const seed = await store.commitAuthority({ + operation_id: "archive-retry-seed", expected_provider_revision: null, + next_projection: todoTerminalProjection(goalId), events: [], receipts: [], + }); + assert.equal(seed.status, "applied"); + if (seed.status !== "applied") return; + const request = {goal_id: goalId, role: "agent" as const, max_active_done: 0, + operation_id: "archive-retry", dry_run: false, + expected_provider_revision: seed.provider_revision, + now: new Date("2026-09-08T01:00:00Z")}; + const before = await store.loadAuthority(); + const stale = await executeCoordinationTodoArchiveCompleted(store, { + ...request, expected_provider_revision: "stale-observed-revision", + }); + assert.equal(stale.status, "conflict", JSON.stringify(stale)); + assert.equal(stale.conflict_kind, "provider_revision_mismatch"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal((await store.readReceipt(request.operation_id)).status, "missing"); + + const applied = await executeCoordinationTodoArchiveCompleted(store, request); + assert.equal(applied.status, "applied", JSON.stringify(applied)); + assert.deepEqual(applied.moved_todo_ids, ["todo-done-a", "todo-done-b"]); + const after = await store.loadAuthority(); + const replay = await executeCoordinationTodoArchiveCompleted(store, { + ...request, now: new Date("2026-09-08T02:00:00Z"), + }); + assert.equal(replay.status, "replayed", JSON.stringify(replay)); + assert.equal(replay.moved_count, 2); + assert.deepEqual(replay.original_receipt, applied.original_receipt); + assert.deepEqual(await store.loadAuthority(), after); + const changedIntent = await executeCoordinationTodoArchiveCompleted(store, { + ...request, expected_provider_revision: String(applied.provider_revision), + }); + assert.equal(changedIntent.reason_code, "coordination_operation_identity_mismatch"); + assert.deepEqual(await store.loadAuthority(), after); + }); + test(`${providerName} conformance: supersede preserves the legacy terminal continuation`, async (t) => { const {store} = await factory(t); const goalId = "goal-supersede"; diff --git a/tests/control_plane_ts/local_archive_attempt.test.ts b/tests/control_plane_ts/local_archive_attempt.test.ts new file mode 100644 index 0000000000..aa39dd6aa0 --- /dev/null +++ b/tests/control_plane_ts/local_archive_attempt.test.ts @@ -0,0 +1,241 @@ +import assert from "node:assert/strict"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import { canonicalAuthoritySha256 } from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import { FileAuthorityStore } from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import { + TODO_CANONICAL_READ_RECORD_FIELDS, + TODO_CANONICAL_READ_RECORD_SCHEMA, + prepareCoordinationProjectionCommit, +} from "../../loopx/control_plane/coordination/coordination_projection.ts"; +import { + acknowledgeLocalCoordinationTodoArchive, + archiveLocalCoordinationTodos, + LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA, + LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, +} from "../../loopx/control_plane/coordination/local_authority_runtime.ts"; +import { shadowManagementDirectory } from "../../loopx/control_plane/coordination/shadow_management.ts"; +import { executeCoordinationTodoArchiveCompleted } from "../../loopx/control_plane/coordination/todo_terminal_lifecycle.ts"; + +function todo(id: string): JsonObject { + return {schema_version: "todo_item_v0", todo_id: id, role: "agent", + status: "done", done: true, text: "Completed archive fixture", + archive_state: "active", source_section: "Agent Todo"}; +} + +async function fixture(t: test.TestContext) { + const root = await mkdtemp(join(tmpdir(), "loopx-archive-attempt-")); + t.after(() => rm(root, {recursive: true, force: true})); + const store = new FileAuthorityStore(join(root, "authority", "file-v0"), "goal-a"); + const todos = [todo("done-a"), todo("done-b")]; + const seed = await store.commitAuthority({ + operation_id: "seed", expected_provider_revision: null, events: [], receipts: [], + next_projection: {goal_id: "goal-a", handoff_mode: "soft_claim", todos, leases: [], + todo_read_model: {schema_version: TODO_CANONICAL_READ_RECORD_SCHEMA, + todo_count: todos.length, records_sha256: canonicalAuthoritySha256(todos), + contract_fields: [...TODO_CANONICAL_READ_RECORD_FIELDS]}}, + }); + assert.equal(seed.status, "applied"); + if (seed.status !== "applied") throw new Error("fixture initialization failed"); + const request = {schema_version: LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, + runtime_root: root, goal_id: "goal-a", role: "agent", max_active_done: 0, + operation_id: "archive-a", expected_provider_revision: seed.provider_revision, + dry_run: false, observed_at: "2026-09-08T01:00:00Z"}; + const ack = (operationId: string) => acknowledgeLocalCoordinationTodoArchive({ + schema_version: LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA, + runtime_root: root, goal_id: "goal-a", role: "agent", operation_id: operationId, + }); + const slotPath = join(shadowManagementDirectory(root, "goal-a"), "archive-agent.json"); + return {root, store, request, ack, slotPath}; +} + +async function addCompleted(store: FileAuthorityStore, id: string) { + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status !== "loaded") throw new Error("fixture authority missing"); + const added = await store.commitAuthority(prepareCoordinationProjectionCommit({ + goal_id: "goal-a", operation_id: `add-${id}`, + expected_provider_revision: head.provider_revision, projection: head.head, + mutations: [{kind: "todo_upsert", todo: todo(id)}], + })); + assert.equal(added.status, "applied"); + if (added.status !== "applied") throw new Error("fixture update failed"); + return added.provider_revision; +} + +test("archive retries recover the accepted batch until exact projection acknowledgement", async (t) => { + const {store, request, ack, slotPath} = await fixture(t); + const applied = await archiveLocalCoordinationTodos(request); + assert.equal(applied.status, "applied", JSON.stringify(applied)); + assert.deepEqual(applied.moved_todo_ids, ["done-a", "done-b"]); + const retained = await readFile(slotPath, "utf8"); + // The caller loses the result before projection. Later unrelated work must + // not change either the retried selection or the original receipt. + const changedRevision = await addCompleted(store, "done-c"); + const retryRequest = {...request, operation_id: "archive-new-suggestion", + expected_provider_revision: changedRevision}; + const replay = await archiveLocalCoordinationTodos(retryRequest); + assert.equal(replay.status, "replayed", JSON.stringify(replay)); + assert.equal(replay.operation_id, "archive-a"); + assert.equal(replay.changed, false); + assert.equal(replay.moved_count, 2); + assert.deepEqual(replay.original_receipt, applied.original_receipt); + assert.equal(await readFile(slotPath, "utf8"), retained); + assert.equal((await archiveLocalCoordinationTodos({...retryRequest, + max_active_done: 1})).conflict_kind, "archive_attempt_pending"); + const beforeAck = await store.loadAuthority(); + assert.equal((await ack("archive-old")).status, "stale"); + assert.equal(await readFile(slotPath, "utf8"), retained); + assert.equal((await ack("archive-a")).status, "acknowledged"); + assert.deepEqual(await store.loadAuthority(), beforeAck, "ACK cannot change domain state"); + + // A successfully delivered attempt ends. This must kill an implementation + // that permanently hashes only goal/role/limit and replays the first batch. + const next = await archiveLocalCoordinationTodos(retryRequest); + assert.equal(next.status, "applied", JSON.stringify(next)); + assert.deepEqual(next.moved_todo_ids, ["done-c"]); + const nextSlot = await readFile(slotPath, "utf8"); + assert.equal((await ack("archive-a")).status, "stale"); + assert.equal(await readFile(slotPath, "utf8"), nextSlot); + assert.equal((await ack("archive-new-suggestion")).status, "acknowledged"); + assert.equal((await ack("archive-new-suggestion")).status, "no_change"); +}); + +test("archive snapshot drift is rejected before selection and releases the uncommitted attempt", async (t) => { + const {store, request, slotPath} = await fixture(t); + const revision = await addCompleted(store, "done-c"); + const before = await store.loadAuthority(); + const rejected = await archiveLocalCoordinationTodos(request); + assert.equal(rejected.status, "conflict", JSON.stringify(rejected)); + assert.equal(rejected.conflict_kind, "provider_revision_mismatch"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal((await store.readReceipt(request.operation_id)).status, "missing"); + assert.equal(JSON.parse(await readFile(slotPath, "utf8")).attempt, null); + const fresh = await archiveLocalCoordinationTodos({...request, + expected_provider_revision: revision, operation_id: "archive-fresh"}); + assert.equal(fresh.status, "applied"); + assert.equal(fresh.moved_count, 3); +}); + +test("archive previews neither create nor acknowledge attempts; no-change retains zero-write authority", async (t) => { + const {store, request, slotPath, ack} = await fixture(t); + const before = await store.loadAuthority(); + assert.equal((await archiveLocalCoordinationTodos({...request, dry_run: true})).status, "planned"); + await assert.rejects(readFile(slotPath), {code: "ENOENT"}); + assert.deepEqual(await store.loadAuthority(), before); + const applied = await archiveLocalCoordinationTodos(request); + const pendingBytes = await readFile(slotPath, "utf8"); + const preview = await archiveLocalCoordinationTodos({...request, dry_run: true, + expected_provider_revision: applied.provider_revision}); + assert.equal(preview.status, "no_change"); + assert.equal(await readFile(slotPath, "utf8"), pendingBytes); + assert.equal((await ack(request.operation_id)).status, "acknowledged"); + const archived = await store.loadAuthority(); + const noChange = await archiveLocalCoordinationTodos({...request, + operation_id: "archive-empty", expected_provider_revision: applied.provider_revision}); + assert.equal(noChange.status, "no_change"); + assert.deepEqual(await store.loadAuthority(), archived); + assert.equal((await store.readReceipt("archive-empty")).status, "missing"); + assert.equal(JSON.parse(await readFile(slotPath, "utf8")).attempt, null); +}); + +test("concurrent same-intent archive calls share one durable receipt", async (t) => { + const {request, store} = await fixture(t); + const results = await Promise.all([request, {...request, operation_id: "contender"}] + .map((input) => archiveLocalCoordinationTodos(input))); + assert.equal(results.filter((result) => result.status === "applied").length, 1); + assert.equal(results.filter((result) => result.status === "replayed").length, 1); + assert.deepEqual(results[0].original_receipt, results[1].original_receipt); + const scan = await store.scanCommitted(null, 10); + assert.equal(scan.status, "page"); + if (scan.status === "page") assert.equal(scan.transactions.length, 2); +}); + +test("archive rejects malformed retry metadata and preview types before effectful dispatch", async (t) => { + const {request, store, slotPath} = await fixture(t); + for (const dryRun of ["true", 1, null]) { + let opened = false; + const result = await archiveLocalCoordinationTodos({...request, dry_run: dryRun}, { + createStore() {opened = true; return store;}, + }); + assert.equal(result.status, "failed"); + assert.equal(opened, false); + } + const applied = await archiveLocalCoordinationTodos(request); + assert.equal(applied.status, "applied"); + const original = await readFile(slotPath, "utf8"); + const slot = JSON.parse(original); + slot.attempt.store_identity = "foreign-store"; + await writeFile(slotPath, JSON.stringify(slot)); + const before = await store.loadAuthority(); + const failure = await archiveLocalCoordinationTodos(request); + assert.equal(failure.status, "failed"); + assert.match(String(failure.reason), /different authority store/); + assert.deepEqual(await store.loadAuthority(), before); +}); + +test("pre-binding v0 archive receipts retain their identity across local upgrade and acknowledgement", async (t) => { + const {request, store, ack} = await fixture(t); + const old = await executeCoordinationTodoArchiveCompleted(store, { + goal_id: request.goal_id, role: "agent", max_active_done: 0, + operation_id: "archive-v0", dry_run: false, now: new Date(request.observed_at), + }); + assert.equal(old.status, "applied"); + const replay = await archiveLocalCoordinationTodos({...request, + operation_id: "archive-v0", expected_provider_revision: undefined}); + assert.equal(replay.status, "replayed", JSON.stringify(replay)); + assert.deepEqual(replay.original_receipt, old.original_receipt); + assert.equal((await ack("archive-v0")).status, "acknowledged"); + const revision = await addCompleted(store, "done-c"); + // Explicitly changing the old request's head binding remains an identity + // conflict, but that rejected attempt cannot block a subsequent new batch. + const mismatch = await archiveLocalCoordinationTodos({...request, + operation_id: "archive-v0", expected_provider_revision: revision}); + assert.equal(mismatch.reason_code, "coordination_operation_identity_mismatch"); + const fresh = await archiveLocalCoordinationTodos({...request, + operation_id: "archive-after-upgrade", expected_provider_revision: revision}); + assert.equal(fresh.status, "applied", JSON.stringify(fresh)); + assert.deepEqual(fresh.moved_todo_ids, ["done-c"]); +}); + +test("archive retains an applied attempt when its receipt is temporarily invisible", async (t) => { + const {request, store, slotPath, ack} = await fixture(t); + const readReceipt = store.readReceipt.bind(store); + let hideCommittedReceipt = true; + store.readReceipt = async (operationId) => { + const receipt = await readReceipt(operationId); + return hideCommittedReceipt && operationId === request.operation_id && receipt.status === "found" + ? {status: "missing"} + : receipt; + }; + const dependencies = {createStore: () => store}; + + const interrupted = await archiveLocalCoordinationTodos(request, dependencies); + assert.equal(interrupted.status, "failed", JSON.stringify(interrupted)); + assert.equal(interrupted.reason_code, "coordination_commit_readback_mismatch"); + const committed = await store.loadAuthority(); + assert.equal(committed.status, "loaded"); + if (committed.status !== "loaded") throw new Error("archive commit was not durable"); + const durableReceipt = await readReceipt(request.operation_id); + assert.equal(durableReceipt.status, "found"); + if (durableReceipt.status !== "found") throw new Error("archive receipt was not durable"); + const retained = await readFile(slotPath, "utf8"); + assert.equal(JSON.parse(retained).attempt.operation_id, request.operation_id); + + hideCommittedReceipt = false; + const replay = await archiveLocalCoordinationTodos({...request, + operation_id: "new-suggestion-after-readback-failure", + expected_provider_revision: committed.provider_revision, + }, dependencies); + assert.equal(replay.status, "replayed", JSON.stringify(replay)); + assert.equal(replay.operation_id, request.operation_id); + assert.equal(replay.moved_count, 2); + assert.deepEqual(replay.original_receipt, durableReceipt.receipts[0]); + assert.deepEqual(await store.loadAuthority(), committed, "retry cannot commit a second archive"); + assert.equal(await readFile(slotPath, "utf8"), retained); + assert.equal((await ack(request.operation_id)).status, "acknowledged"); +}); diff --git a/tests/control_plane_ts/shadow_native_writer_boundary.test.ts b/tests/control_plane_ts/shadow_native_writer_boundary.test.ts index 41770b89ce..a1357c9cda 100644 --- a/tests/control_plane_ts/shadow_native_writer_boundary.test.ts +++ b/tests/control_plane_ts/shadow_native_writer_boundary.test.ts @@ -11,10 +11,12 @@ import { } from "../../loopx/control_plane/coordination/shadow_management.ts"; import { archiveLocalCoordinationTodos, + acknowledgeLocalCoordinationTodoArchive, createLocalCoordinationTodo, claimLocalCoordinationTodo, mutateLocalCoordinationAuthority, editLocalCoordinationTodo, terminalLifecycleLocalCoordinationTodo, LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, + LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_CREATE_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_CLAIM_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_TERMINAL_LIFECYCLE_REQUEST_SCHEMA, LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA, @@ -31,7 +33,14 @@ for (const [name, invoke, schema, requestFields] of [ linked_successor_todo_ids: [], lease_expected_version: null, }], ["archive", archiveLocalCoordinationTodos, - LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, {max_active_done: 0}], + LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, { + max_active_done: 0, role: "agent", operation_id: "archive-maintenance", + observed_at: "2026-01-01T00:00:00Z", + }], + ["archive acknowledgement", acknowledgeLocalCoordinationTodoArchive, + LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA, { + role: "agent", operation_id: "archive-maintenance", + }], ] as const) { test(`promoted ${name} checks maintenance before opening a provider`, async (t) => { const root = await mkdtemp(join(tmpdir(), "loopx-native-maintenance-")); @@ -54,7 +63,14 @@ for (const [name, invoke, schema, requestFields] of [ linked_successor_todo_ids: [], lease_expected_version: null, }], ["archive", archiveLocalCoordinationTodos, - LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, {max_active_done: 0}], + LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA, { + max_active_done: 0, role: "agent", operation_id: "archive-maintenance", + observed_at: "2026-01-01T00:00:00Z", + }], + ["archive acknowledgement", acknowledgeLocalCoordinationTodoArchive, + LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA, { + role: "agent", operation_id: "archive-maintenance", + }], ] as const) { test(`promoted ${name} waits behind the bootstrap and rollback maintenance lock`, async (t) => { const root = await mkdtemp(join(tmpdir(), "loopx-native-maintenance-race-")); diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 313ab8ac75..3a033b62c5 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -75,6 +75,7 @@ "tests/control_plane_ts/coordination_projection.test.ts", "tests/control_plane_ts/coordination_state_contract.test.ts", "tests/control_plane_ts/local_authority_runtime.test.ts", + "tests/control_plane_ts/local_archive_attempt.test.ts", "tests/control_plane_ts/local_authority_shadow_outbox.test.ts", "tests/control_plane_ts/authority_store_conformance.ts", "tests/control_plane_ts/nokv_authority_store.test.ts", From 28354cc7b00271050cffa25700cbe57338f3f120 Mon Sep 17 00:00:00 2001 From: Steven <96154058+steven-kid@users.noreply.github.com> Date: Wed, 9 Sep 2026 00:59:16 +0800 Subject: [PATCH 2/3] fix(todos): allow archive retry to recover missing display Signed-off-by: Steven <96154058+steven-kid@users.noreply.github.com> --- .../todo-archive-delivery-recovery.md | 10 +++++ loopx/control_plane/todos/path_resolution.py | 3 +- .../todos/provider_terminal_lifecycle.py | 3 ++ .../test_archive_retry_delivery.py | 38 ++++++++++++++++++- 4 files changed, 51 insertions(+), 3 deletions(-) diff --git a/docs/architecture/todo-archive-delivery-recovery.md b/docs/architecture/todo-archive-delivery-recovery.md index 90710b6451..87467ab102 100644 --- a/docs/architecture/todo-archive-delivery-recovery.md +++ b/docs/architecture/todo-archive-delivery-recovery.md @@ -49,6 +49,12 @@ This does not promise recovery of a CLI response lost after acknowledgement or discover attempts made before correlation records existed. Callers requiring that stronger guarantee need an explicit caller-retained request identity. +Missing Markdown can be rebuilt from canonical Todos during archive retry; +successful recovery still requires the same exact-attempt acknowledgement. +An explicit `todo project-markdown` rebuild repairs the display without +acknowledging an archive attempt. The next archive with the same role and limit +first replays and acknowledges that prior batch; it does not archive a new one. + A changed promoted archive uses four Python/TypeScript request-responses: authority read, archive transaction, projection readback, and acknowledgement. The previous path used three; empty and preview paths add no acknowledgement. @@ -75,6 +81,10 @@ crash/retry conformance coverage. 原 operation、receipt 和 store。过期确认不能清除新操作;结果不确定时保留重试身份, 确定未提交的拒绝允许后续修正请求继续。 +Markdown 缺失时,归档重试可从 canonical Todos 重建投影,成功后仍须确认原 attempt。 +单独运行 `todo project-markdown` 只修复显示,不确认归档 attempt;之后相同 role/limit +的归档调用会先重放并确认旧批次,不会归档新一批。 + 恢复保证截至投影确认成功;之后相同 CLI 调用代表新一批归档,不承诺恢复确认之后才 丢失的 stdout,也不追溯发现旧版本未记录的调用。发生归档的路径从三次跨运行时调用 增加为四次,额外确认和持久化是明确的正确性成本。canonical journal 消费者同时接管 diff --git a/loopx/control_plane/todos/path_resolution.py b/loopx/control_plane/todos/path_resolution.py index f7e868d985..f09f88dc66 100644 --- a/loopx/control_plane/todos/path_resolution.py +++ b/loopx/control_plane/todos/path_resolution.py @@ -14,6 +14,7 @@ def resolve_todo_state_path( goal_id: str, project: Path | None = None, state_file: Path | None = None, + require_existing: bool = True, ) -> tuple[Path | None, Path]: registry = load_registry(registry_path) goal, resolved_project, resolved_state_file = resolve_goal_state( @@ -24,6 +25,6 @@ def resolve_todo_state_path( ) if goal is None: raise ValueError(f"goal {goal_id!r} is not present in the registry") - if not resolved_state_file.exists(): + if require_existing and not resolved_state_file.exists(): raise ValueError(f"active state file does not exist: {resolved_state_file}") return resolved_project, resolved_state_file diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index 2a989308f8..179c4463ee 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -74,6 +74,9 @@ def _route_terminal_call(command: str, call: Mapping[str, Any]) -> dict[str, Any goal_id=goal_id, project=call.get("project"), state_file=call.get("state_file"), + # Canonical archive can recover its display after committing or + # replaying. The legacy fallback still requires an existing file. + require_existing=False, ) return archive_canonical_todos_if_promoted( registry_path=registry_path, diff --git a/tests/control_plane/test_archive_retry_delivery.py b/tests/control_plane/test_archive_retry_delivery.py index c847216ee1..f8bc96319c 100644 --- a/tests/control_plane/test_archive_retry_delivery.py +++ b/tests/control_plane/test_archive_retry_delivery.py @@ -138,8 +138,9 @@ def _archive_cli( ) +@pytest.mark.parametrize("missing_display", [False, True], ids=["existing", "missing"]) def test_archive_cli_recovers_committed_result_after_process_exit( - tmp_path: Path, + tmp_path: Path, missing_display: bool, ) -> None: registry, runtime, state = _fixture(tmp_path) original_markdown = state.read_bytes() @@ -151,6 +152,8 @@ def test_archive_cli_recovers_committed_result_after_process_exit( assert committed["moved_count"] == 1 assert committed["moved_todo_ids"] == ["todo_completed"] assert state.read_bytes() == original_markdown + if missing_display: + state.unlink() retry = _archive_cli(registry) assert retry.returncode == 0, retry.stderr or retry.stdout @@ -162,6 +165,10 @@ def test_archive_cli_recovers_committed_result_after_process_exit( assert recovered["provider_revision"] == committed["provider_revision"] assert recovered["changed"] is False assert recovered["projection_delivery"] in {"delivered", "current"} + assert state.exists() + assert recovered["archive_delivery_ack"]["status"] == "acknowledged" + if missing_display: + assert recovered["projection_outbox"]["recovery_scope"] == "todo_sections_only" canonical = read_canonical_todos_if_promoted( runtime_root=runtime, goal_id="archive-goal" ) @@ -198,12 +205,16 @@ def test_archive_cli_recovers_committed_result_after_process_exit( ) +@pytest.mark.parametrize("missing_display", [False, True], ids=["existing", "missing"]) def test_pending_projection_and_preview_preserve_archive_retry( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + missing_display: bool, ) -> None: registry, runtime, state = _fixture(tmp_path) original_markdown = state.read_text(encoding="utf-8") + if missing_display: + state.unlink() def reject_projection(*_args, **_kwargs): raise OSError("injected projection write failure") @@ -222,7 +233,10 @@ def reject_projection(*_args, **_kwargs): preview = _archive_cli(registry, execute=False) assert preview.returncode == 0, preview.stderr assert json.loads(preview.stdout)["moved_count"] == 0 - assert state.read_text(encoding="utf-8") == original_markdown + if missing_display: + assert not state.exists() + else: + assert state.read_text(encoding="utf-8") == original_markdown canonical = read_canonical_todos_if_promoted( runtime_root=runtime, goal_id="archive-goal" ) @@ -237,6 +251,7 @@ def reject_projection(*_args, **_kwargs): assert replay["status"] == "replayed" assert replay["original_receipt"] == committed["original_receipt"] assert replay["projection_delivery"] in {"delivered", "current"} + assert replay["archive_delivery_ack"]["status"] == "acknowledged" def test_archive_rejects_snapshot_drift_before_selection( @@ -317,3 +332,22 @@ def unavailable_ack(method, params): ) assert canonical is not None assert canonical["provider_revision"] == committed["provider_revision"] + + +def test_unpromoted_archive_does_not_rebuild_missing_display(tmp_path: Path) -> None: + from loopx.control_plane.coordination.legacy_writer_fence import ( + legacy_coordination_writer_fence_path, + ) + + registry, runtime, state = _fixture(tmp_path) + state.unlink() + legacy_coordination_writer_fence_path( + runtime_root=runtime, goal_id="archive-goal" + ).unlink() + authority = runtime / "authority" / "file-v0" + before = {path: path.read_bytes() for path in authority.rglob("*") if path.is_file()} + result = _archive_cli(registry) + assert result.returncode == 1 + assert "active state file does not exist" in json.loads(result.stdout)["error"] + assert not state.exists() + assert {path: path.read_bytes() for path in authority.rglob("*") if path.is_file()} == before From 6c575caefdb34a8614578d0007c1c15a8f3f4e39 Mon Sep 17 00:00:00 2001 From: Steven <96154058+steven-kid@users.noreply.github.com> Date: Wed, 9 Sep 2026 22:19:45 +0800 Subject: [PATCH 3/3] fix(qualification): scan decoded scheduler facts before actor dispatch Signed-off-by: Steven <96154058+steven-kid@users.noreply.github.com> --- .../testing/model_behavior_packet_safety.py | 137 +++++++++++++ .../testing/model_behavior_qualification.py | 3 +- .../test_model_behavior_qualification.py | 185 ++++++++++++++++++ .../model_behavior_scheduler_transport.json | 10 + 4 files changed, 334 insertions(+), 1 deletion(-) create mode 100644 loopx/control_plane/testing/model_behavior_packet_safety.py create mode 100644 tests/fixtures/control_plane/model_behavior_scheduler_transport.json diff --git a/loopx/control_plane/testing/model_behavior_packet_safety.py b/loopx/control_plane/testing/model_behavior_packet_safety.py new file mode 100644 index 0000000000..bc44f225b0 --- /dev/null +++ b/loopx/control_plane/testing/model_behavior_packet_safety.py @@ -0,0 +1,137 @@ +"""Expose typed scheduler wire contents to the actor packet's safety scanner.""" +from __future__ import annotations + +import base64 +import binascii +import json +import re +import zlib +from collections.abc import Mapping +from typing import Any + +_FACTS_FLAG = "--scheduler-host-facts-chunk" +# Match the existing native follow-up decoder, without invoking its effects. +_MAX_ENCODED_CHARS = 4_096 +_MAX_INFLATED_BYTES = 16_384 + + +def _unique_object(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + result: dict[str, Any] = {} + for key, value in pairs: + if key in result: + raise ValueError("duplicate JSON key") + result[key] = value + return result + + +def _reject_constant(_value: str) -> Any: + raise ValueError("non-finite JSON constant") + + +def _validation_copy(value: Any) -> Any: + # Unlike deepcopy, do not preserve aliases: validating a typed argv must + # never exempt an alias of that list exposed in an unrelated packet field. + if isinstance(value, Mapping): + return {key: _validation_copy(item) for key, item in value.items()} + if isinstance(value, list): + return [_validation_copy(item) for item in value] + return value + + +def _decode_facts(chunks: list[str]) -> dict[str, Any]: + try: + if sum(map(len, chunks)) > _MAX_ENCODED_CHARS: + raise ValueError("encoded boundary") + encoded = "".join(chunks) + if not encoded or not re.fullmatch(r"[A-Za-z0-9_-]+", encoded): + raise ValueError("encoded alphabet") + compressed = base64.b64decode(encoded + "=" * (-len(encoded) % 4), altchars=b"-_", validate=True) + if base64.urlsafe_b64encode(compressed).decode().rstrip("=") != encoded: + raise ValueError("non-canonical base64url") + inflater = zlib.decompressobj() + raw = inflater.decompress(compressed, _MAX_INFLATED_BYTES + 1) + if len(raw) > _MAX_INFLATED_BYTES or not inflater.eof or inflater.unused_data or inflater.unconsumed_tail: + raise ValueError("inflated boundary or incomplete/trailing stream") + payload = json.loads(raw.decode("utf-8"), object_pairs_hook=_unique_object, parse_constant=_reject_constant) + if not isinstance(payload, dict) or payload.get("schema_version") != "loopx_scheduler_host_followup_hint_v0": + raise ValueError("hint schema") + facts = payload.get("host_facts") + if not isinstance(facts, dict) or facts.get("schema_version") != "loopx_scheduler_heartbeat_host_facts_v0": + raise ValueError("facts schema") + if not isinstance(payload.get("before"), dict) or not isinstance(payload.get("use_current_hint"), bool): + raise ValueError("hint shape") + return payload + except (ValueError, binascii.Error, zlib.error, RecursionError): + # Do not echo untrusted bytes or parsed values into a public failure. + raise ValueError("scheduler host facts must be bounded canonical compressed JSON") from None + + +def _expose_chunks(args: list[Any], *, command: str, schema_valid: bool) -> None: + positions: list[int] = [] + chunks: list[str] = [] + index = 0 + encoded_chars = 0 + while index < len(args): + item = args[index] + prefix_length = 0 + if item == _FACTS_FLAG: + index += 1 + if index == len(args) or not isinstance(args[index], str): + raise ValueError("scheduler host facts chunk value is missing") + item = args[index] + elif isinstance(item, str) and item.startswith(_FACTS_FLAG + "="): + prefix_length = len(_FACTS_FLAG) + 1 + else: + index += 1 + continue + chunk_length = len(item) - prefix_length + if not chunk_length: + raise ValueError("scheduler host facts chunk value is missing") + encoded_chars += chunk_length + if encoded_chars > _MAX_ENCODED_CHARS: + raise ValueError("scheduler host facts exceed the encoded boundary") + positions.append(index) + chunks.append(item[prefix_length:]) + index += 1 + if not chunks: + return # Legacy, unencoded hints still receive ordinary recursive scanning. + if not schema_valid or args[:2] != ["quota", command]: + raise ValueError("scheduler host facts require the matching typed hint and command") + payload = _decode_facts(chunks) + # Only the validation view changes. Scan the entire decoded envelope once, + # including extensions; retain every non-transport argument for scanning. + for position in positions: + args[position] = "" + args[positions[0]] = payload + + +def scheduler_transport_validation_view(packet: Mapping[str, Any], *, arm: str) -> dict[str, Any]: + """Decode only real packet fields; dotted key names cannot imitate a path. + + This is a confidentiality view, not scheduler admission or execution. The + caller must recursively scan it, and must send the original packet onward. + """ + view: dict[str, Any] = _validation_copy(packet) + if arm == "full_packet": + scheduler = view.get("scheduler_hint") + else: + scheduler = view.get("scheduler") + if not isinstance(scheduler, Mapping): + return view + codex_app = scheduler.get("codex_app") + if not isinstance(codex_app, Mapping): + return view + if arm == "full_packet": + for kind, command in (("ack", "scheduler-ack-current"), ("failure", "scheduler-fail-current")): + hint = codex_app.get(kind + "_hint") + args = hint.get("cli_args") if isinstance(hint, Mapping) else None + if isinstance(hint, Mapping) and isinstance(args, list): + _expose_chunks(args, command=command, + schema_valid=scheduler.get("schema_version") == "scheduler_hint_v0" + and hint.get("schema_version") == f"codex_app_scheduler_{kind}_hint_v0") + else: + args = codex_app.get("ack_cli_args") + if isinstance(args, list): + _expose_chunks(args, command="scheduler-ack-current", + schema_valid=packet.get("schema_version") == "loopx_turn_envelope_v0") + return view diff --git a/loopx/control_plane/testing/model_behavior_qualification.py b/loopx/control_plane/testing/model_behavior_qualification.py index c5e3049bdd..8b662c6496 100644 --- a/loopx/control_plane/testing/model_behavior_qualification.py +++ b/loopx/control_plane/testing/model_behavior_qualification.py @@ -17,6 +17,7 @@ SECRET_LIKE_SURFACE_PATTERN, ) from ..work_items.interaction_contract import INTERACTION_RESPONSE_PLAN_SCHEMA_VERSION +from .model_behavior_packet_safety import scheduler_transport_validation_view MODEL_BEHAVIOR_QUALIFICATION_SCHEMA_VERSION = "model_behavior_qualification_v0" @@ -339,7 +340,7 @@ def build_model_behavior_actor_request( ) packet_schema_version = _packet_schema(normalized_packet, arm=arm) _validated_packet_response_plan(normalized_packet, arm=arm) - _reject_private_or_secret_material(normalized_packet) + _reject_private_or_secret_material(scheduler_transport_validation_view(normalized_packet, arm=arm)) return { "schema_version": MODEL_BEHAVIOR_ACTOR_REQUEST_SCHEMA_VERSION, "qualification_id": _token(qualification_id, field="qualification_id"), diff --git a/tests/control_plane/test_model_behavior_qualification.py b/tests/control_plane/test_model_behavior_qualification.py index 48e3aa92fe..ad50a28461 100644 --- a/tests/control_plane/test_model_behavior_qualification.py +++ b/tests/control_plane/test_model_behavior_qualification.py @@ -1,12 +1,16 @@ from __future__ import annotations +import base64 import json +import zlib from collections.abc import Mapping +from pathlib import Path from typing import Any import pytest from loopx.control_plane.quota.turn_envelope import build_turn_envelope +from loopx.control_plane.runtime.public_safety import SECRET_LIKE_SURFACE_PATTERN from loopx.control_plane.testing.model_behavior_qualification import ( MODEL_BEHAVIOR_ACTOR_RESULT_SCHEMA_VERSION, MODEL_BEHAVIOR_ARM_TERMINAL_RECEIPT_SCHEMA_VERSION, @@ -15,6 +19,7 @@ build_model_behavior_actor_request, compare_model_behavior_receipts, model_behavior_semantic_contract_from_packet, + normalize_model_behavior_actor_request, run_model_behavior_qualification_arm, run_model_behavior_qualification_pair, ) @@ -161,6 +166,186 @@ def test_actor_request_rejects_private_or_secret_material( ) +_HOST_FACTS_FLAG = "--scheduler-host-facts-chunk" + + +def _scheduler_wire(operation: str = "ack") -> bytes: + fixture = Path(__file__).parents[1] / "fixtures/control_plane/model_behavior_scheduler_transport.json" + return bytes.fromhex(json.loads(fixture.read_text())[operation]["zlib_hex"]) + + +def _scheduler_transport_packet( + arm: str, operation: str, *, inline: bool = False, compressed: bytes | None = None, +) -> tuple[dict[str, Any], list[str]]: + if compressed is None: + # Frozen public synthetic wire bytes keep the collision independent of + # the platform's zlib encoder. Hex storage is not a secret-shaped value. + compressed = _scheduler_wire(operation) + encoded = base64.urlsafe_b64encode(compressed).decode().rstrip("=") + command = "scheduler-ack-current" if operation == "ack" else "scheduler-fail-current" + args = ["quota", command, "--goal-id", "goal-native-followup", "--agent-id", "agent-native-followup"] + for offset in range(0, len(encoded), 384): + chunk = encoded[offset:offset + 384] + args.extend([f"{_HOST_FACTS_FLAG}={chunk}"] if inline else [_HOST_FACTS_FLAG, chunk]) + packet = _full_packet() + if arm == "full_packet": + kind = "ack" if operation == "ack" else "failure" + packet["scheduler_hint"] = {"schema_version": "scheduler_hint_v0", "codex_app": { + f"{kind}_hint": {"schema_version": f"codex_app_scheduler_{kind}_hint_v0", "cli_args": args}, + }} + else: + packet = build_turn_envelope(packet) + packet["scheduler"] = {"codex_app": {"ack_cli_args": args}} + return packet, args + + +@pytest.mark.parametrize("arm, operation", [ + ("full_packet", "ack"), ("full_packet", "host_failure"), ("candidate_packet", "ack"), +]) +@pytest.mark.parametrize("inline", [False, True]) +def test_actor_request_scans_decoded_scheduler_facts_without_changing_wire( + arm: str, operation: str, inline: bool, +) -> None: + packet, args = _scheduler_transport_packet(arm, operation, inline=inline) + assert any(SECRET_LIKE_SURFACE_PATTERN.search(arg) for arg in args) + before = json.dumps(packet, sort_keys=True) + + request = build_model_behavior_actor_request(packet, qualification_id="public-wire-collision", arm=arm) + + assert json.dumps(request["packet"], sort_keys=True) == before + assert json.dumps(packet, sort_keys=True) == before + assert normalize_model_behavior_actor_request(request) == request + + +@pytest.mark.parametrize("arm, operation", [ + ("full_packet", "ack"), ("full_packet", "host_failure"), ("candidate_packet", "ack"), +]) +@pytest.mark.parametrize("location", ["before", "host_facts", "extension"]) +def test_actor_rejects_private_material_inside_encoded_scheduler_facts( + arm: str, operation: str, location: str, +) -> None: + payload = json.loads(zlib.decompress(_scheduler_wire(operation))) + if location == "before": + payload[location]["api_key"] = "synthetic-private-value" + elif location == "host_facts": + payload[location]["note"] = "/" + "Users/example/private.txt" + else: + payload[location] = [{"nested": ["token" + "=abcdefghijklmnop"]}] + packet, _ = _scheduler_transport_packet(arm, operation, inline=True, + compressed=zlib.compress(json.dumps(payload).encode())) + before = json.dumps(packet, sort_keys=True) + with pytest.raises(ValueError, match="credential-shaped field|local absolute path|credential-like value"): + build_model_behavior_actor_request(packet, qualification_id="encoded-private", arm=arm) + assert json.dumps(packet, sort_keys=True) == before + + +@pytest.mark.parametrize("case", [ + "bad_json", "utf8", "root_array", "duplicate_key", "nonfinite", "hint_schema", "facts_schema", + "before_shape", "current_hint_shape", "truncated", "trailing", "concatenated", "bomb", "encoded_limit", +]) +def test_actor_rejects_malformed_or_unbounded_scheduler_wire(case: str) -> None: + payload = json.loads(zlib.decompress(_scheduler_wire())) + if case == "hint_schema": + payload["schema_version"] = "future_hint" + elif case == "facts_schema": + payload["host_facts"]["schema_version"] = "future_facts" + elif case == "before_shape": + payload["before"] = [] + elif case == "current_hint_shape": + payload["use_current_hint"] = "true" + elif case == "nonfinite": + payload["extension"] = float("nan") + elif case == "bomb": + payload["extension"] = "a" * 16_385 + raw = json.dumps(payload).encode() + if case == "bad_json": + raw = b"not-json" + elif case == "utf8": + raw = b"\xff" + elif case == "root_array": + raw = b"[]" + elif case == "duplicate_key": + raw = raw[:-1] + b', "before": {}}' + compressed = zlib.compress(raw) + if case == "truncated": + compressed = compressed[:-1] + elif case == "trailing": + compressed += b"hidden tail" + elif case == "concatenated": + compressed += zlib.compress(b"{}") + elif case == "encoded_limit": + compressed = b"x" * 3_073 + packet, _ = _scheduler_transport_packet("full_packet", "ack", compressed=compressed) + with pytest.raises(ValueError, match="scheduler host facts"): + build_model_behavior_actor_request(packet, qualification_id="invalid-wire", arm="full_packet") + + +@pytest.mark.parametrize("case", [ + "hint_schema", "scheduler_schema", "command", "missing", "empty_inline", "alphabet", "padding", "pad_bits", +]) +def test_actor_rejects_unrecognized_or_malformed_scheduler_arguments(case: str) -> None: + packet, args = _scheduler_transport_packet("full_packet", "ack") + if case == "hint_schema": + packet["scheduler_hint"]["codex_app"]["ack_hint"]["schema_version"] = "future_hint" + elif case == "scheduler_schema": + packet["scheduler_hint"]["schema_version"] = "future_scheduler" + elif case == "command": + args[1] = "scheduler-fail-current" + elif case == "missing": + args.append(_HOST_FACTS_FLAG) + elif case == "empty_inline": + args.append(_HOST_FACTS_FLAG + "=") + elif case == "alphabet": + args[7] = "+invalid" + elif case == "padding": + args[-1] += "=" + else: + # One byte canonically encodes as eA; eB has nonzero unused pad bits. + args[6:] = [_HOST_FACTS_FLAG, "eB"] + with pytest.raises(ValueError, match="scheduler host facts"): + build_model_behavior_actor_request(packet, qualification_id="invalid-args", arm="full_packet") + + +@pytest.mark.parametrize("case", ["ordinary_arg", "unrelated", "dotted_key", "alias"]) +def test_actor_does_not_exempt_untyped_fields_or_other_arguments(case: str) -> None: + packet, args = _scheduler_transport_packet("full_packet", "ack") + if case == "ordinary_arg": + args.extend(["--reason-summary", "token" + "=abcdefghijklmnop"]) + elif case == "unrelated": + packet["diagnostic"] = {"cli_args": args.copy()} + elif case == "dotted_key": + packet["scheduler_hint.codex_app.ack_hint.cli_args"] = args.copy() + else: + packet["diagnostic"] = args # Same object as the valid typed argv. + with pytest.raises(ValueError, match="credential-like value"): + build_model_behavior_actor_request(packet, qualification_id="untyped-wire", arm="full_packet") + + +def test_actor_normalization_rescans_tampered_encoded_material() -> None: + packet, _ = _scheduler_transport_packet("candidate_packet", "ack") + request = build_model_behavior_actor_request(packet, qualification_id="tampered-wire", arm="candidate_packet") + payload = json.loads(zlib.decompress(_scheduler_wire())) + payload["extension"] = {"password": "synthetic-private-value"} + tainted, _ = _scheduler_transport_packet("candidate_packet", "ack", + compressed=zlib.compress(json.dumps(payload).encode())) + request["packet"] = tainted + with pytest.raises(ValueError, match="credential-shaped field"): + normalize_model_behavior_actor_request(request) + + +def test_actor_accepts_unencoded_legacy_hint_and_inflated_limit() -> None: + packet, args = _scheduler_transport_packet("full_packet", "ack") + args[6:] = ["--reset-token", "public-reset"] + assert build_model_behavior_actor_request(packet, qualification_id="legacy-hint", arm="full_packet")["packet"] == packet + payload = json.loads(zlib.decompress(_scheduler_wire())) + payload["padding"] = "" + payload["padding"] = "a" * (16_384 - len(json.dumps(payload).encode())) + raw = json.dumps(payload).encode() + assert len(raw) == 16_384 + packet, _ = _scheduler_transport_packet("full_packet", "ack", compressed=zlib.compress(raw)) + assert build_model_behavior_actor_request(packet, qualification_id="at-limit", arm="full_packet")["packet"] == packet + + def test_qualification_receipt_is_compact_and_drops_raw_conversation() -> None: receipt = run_model_behavior_qualification_arm( _full_packet(), diff --git a/tests/fixtures/control_plane/model_behavior_scheduler_transport.json b/tests/fixtures/control_plane/model_behavior_scheduler_transport.json new file mode 100644 index 0000000000..10304fe807 --- /dev/null +++ b/tests/fixtures/control_plane/model_behavior_scheduler_transport.json @@ -0,0 +1,10 @@ +{ + "ack": { + "public_fixture_seed": 327, + "zlib_hex": "78da8d53db6e133110fd173f6faacd6dd304f5818722552a95a80a1254c81a7bc78915676d7c691a55fd77c6dea42910a06febd99939171f3f3181ca7a648b27268c956b6c39c8a86dc783b48eea5d32a662121c086d74dc718f0eb4e7608cdd62cb160a4cc08aa15248830fb89f670bd659bf01c37dea587538b468a8c7ef8ef3d1271aff916c84cc625fe7c1d818d862d810b8ddb81489cbb062b9cc37baa373fe3ba58ac32e1edaeb8a6d75d7da2d5fd9e4a93079ae9847690be69fe07bf2011472b17310c23fffad69f7c1928046fdcd8c40e8a62dcaf7fa4284ac801183a51606c991adf5ebe040e2e92d447c6543e48afc0cc519b9e61d627b340d9659baa673ff39e820dfc040d9bc29390201e78c263bbd4f26e37fb8bdfc74f1f1eae6f3dde5f5d77757377797b75fde5f5f0ca7fbdedd6f10125aec88a234a43fe3f4779cb9d3043e3abaf437ae57a04df2f88b89c41a3de40d10697c548f9a417d3e18cdeeea6631ae1775fd8d269796925364e6af132a8b511b8872c5ad08e81f8e0234f18f39b7412f6930e5ac938a09ce053442526aa762d2d6a215cdacc69918a1c0e95c4839c67a329b4d27f373054a4c019a218e140ce7ed5011e4018717ecffe81fe7019795f62f83ae922acedba5c710f26b234ff0b1e4f775f525e8f794f4715d35f5f71ce7809147bbc6bcaa9c4e5812e40a37c029eea1c734d6ba479ecb2d71f57c85e0a34088fc1833fe50e7517a3a32ab298ff2d5484fbb4499af71471dafd6e92e9e494b2a38e5e8ac34a964b8a021abb263217902c97b5fdad8f35b88167a7b6505a7a79902c532799fdf40aef637fefc13190fb750" + }, + "host_failure": { + "public_fixture_seed": 2395, + "zlib_hex": "78da8d534b6f133110fe2f3e6faa6d92dd3c500f1c8a54a954a22a488090356b8f132bcedaf8d134aafadf193b9b368202bdd9e3197f8f9979641d2aeb912d1f5967acd8a0e420a2b63d0fc23a8af7c9988a0970d069a3e39e7b74a03d0763ec0e255b2a30012b864a2115dee350cf96acb77e0b86fbd4b3ea78916828c7ef5feaa34f54fe33d90899c510e7c1d818d8f2bc2570bb75291297f38ae530dfea9eeef9b5a188c33e1ed3eb8aed742fed8eaf6df214983e55cca3b005f34ff0817c0085bcdb3b08e19f6f1bfafb684940a3fe66462074238bf2415f8890153062b0d29d41726467fd263810f8fa2f447c6d43e48afc0cc519b1e13da27c310d5659baa6fbe138ea217760a46cfe29390201e78c263bbd4f26e37fb8bdfc74f1f1eae6f3dde5f5d77757377797b75fde5f5f9c3743eefe370801127ba2280ce9cf38871e67ee54810f8e9afec6ef1568933c0e26b2a8b76853a407a28e1ef23710e9615c8fdb513d1f8d677775bb9cd4cbbafe96b32c8d4fd19a4faf482d6e6d218a35b75d407f7fd2104d2a629edea05754997ce905a8859cb462ba80c914652717b339346321e762368349a3daa6192b558bf1bc966ddbc9e90c9a1908b9c0a6992f24611e817801ff8f0b935ce0b2d4c37e0ced2daed093f376e53184bc7c64113e94713e8d3ecffd771afc495db5f58f3cdd01238f7683f9cf727bc59c20d6b8054ed31f0ee0c65af7c073581269cfd7083e760891bf4c1dbfaf73296d92c8b2ca8e9e96e4c4e424b5ee44469974bec13d559ce4ea3e9e094baa388dd9594952c9f08ea6da2a95eb9227d08cf39cc69ede42bcd01d94169c03ed14686a93f7794572f4b856bf0028b1c2db" + } +}