diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 6de929c984..4e320120ad 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -534,13 +534,22 @@ but optimizing it must retain identity consumption, replay and conflict checks. Current local facades reuse their generated id within managed-runtime retries. Separate CLI invocations are not implicitly one attempt: claim exposes -`--claim-operation-id`, while create and text/note update do not currently expose -an equivalent cross-process recovery key. That is a caller-recovery limitation, +`--claim-operation-id` and update exposes `--update-operation-id`; create does not +currently expose an equivalent cross-process recovery key. That is a caller-recovery limitation, not proof of duplicate business effects or universal exactly-once execution. Any extension must define the retry boundary and distinguish retries from new intent before adding keys or durable attempt tracking. Test lost responses and intervening writes; a source-level ban on UUID construction proves neither. +The canonical Todo commands now share one TS receipt recovery owner. Their +local result contract retains unresolved post-commit readback as `ambiguous` +and names the original operation for recovery; it does not infer no-write from +an unavailable receipt or retry the CAS automatically. See the +[command recovery checkpoint](typescript-control-plane-migration-v0.md#command-receipt-and-recovery-ownership) +for intentional diagnostic changes and the complete fixture matrix. Historical +receipt identity and one-way Markdown delivery remain intact. This is command +recovery qualification, not storage retention, service availability or promotion. + For every request, the authority performs this sequence: 1. load the aggregate and provider generation; diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index e1602677f3..02215cc4bd 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -454,12 +454,19 @@ Operation identity 标识调用方的一次逻辑尝试,不是参数组合。 存储成本,但优化时必须保留 identity consumption、重放与冲突校验。 当前本地 facade 在 managed-runtime retry 内复用生成的 id。两次独立 CLI 调用不会 -自动视为同一尝试:claim 提供 `--claim-operation-id`,create 和 text/note update -目前没有等价的跨进程恢复 key。这是 caller recovery 的限制,不证明业务效果重复, +自动视为同一尝试:claim 提供 `--claim-operation-id`,update 提供 +`--update-operation-id`,create 目前没有等价的跨进程恢复 key。这是 caller recovery 的限制,不证明业务效果重复, 也不能宣称通用 exactly-once。扩展前应先定义重试边界、区分 retry 与新 intent,再 决定是否需要 key 或耐久 attempt tracking。用丢响应与中间插入其他写入来验证, 而不是用禁止 UUID 构造的源码扫描代替语义测试。 +Canonical Todo 命令现在共用一个 TS 回执恢复 owner。本地结果合同将提交后未能 +确认的回读保留为 `ambiguous`,并指出恢复所需的原 operation;不能从回执不可用 +推断未写入,也不会自动重试 CAS。明确的诊断变化及完整 fixture 矩阵见 +[命令恢复检查点](typescript-control-plane-migration-v0.zh-CN.md#命令回执与恢复的统一所有者)。 +历史回执身份和单向 Markdown 投递保持不变。这是命令恢复验证,不代表存储保留、 +服务可用性或 promotion 已获资格。 + 对每个 request,authority 执行以下顺序: 1. load aggregate 与 provider generation; diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 0d69816274..d1ab5f5fcf 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -429,6 +429,49 @@ rewrite or new writer admission is implied. General add/update admission and a generic repair action for every invalid condition remain separate scopes; this is not a claim of zero behavior change or full Todo writer closure. +#### Command receipt and recovery ownership + +At baseline `bfd1ec8db`, create, claim, update, complete/supersede, archive and +Monitor poll repeated envelope matching, result projection and post-CAS +readback. `coordination/command_receipt.ts` now owns those shared semantics; +command modules retain request normalization/digests, admission, payload +validation and mutations. `coordination/todo_archive.ts` owns retention, +separate from terminal validation and lease release. Internal callers import +that owner directly; the old module does not retain an unused re-export. + +Intentional observable changes on these canonical command paths: + +- An applied/ambiguous commit followed by unreadable receipt remains + `ambiguous`, with `recovery.operation_id` and + `retry_with_same_operation_id=true`. A read failure cannot erase possible + durable acceptance. A thrown commit response receives one receipt lookup, + never an automatic second write. +- A conclusive CAS conflict or failed commit remains that result if diagnostic + readback fails. An exact historical receipt still takes precedence. An + applied response with a missing receipt remains a protocol failure. +- Malformed create/Monitor result objects and update/terminal/archive change decisions + fail with `invalid_coordination_command_receipt`; they cannot become successful + replay/no-op through coercion or escape as an unchecked decoder error. Claim + retains its existing receipt-error code and historical omitted-change codec. +- Read failures consistently include `changed=false`; this with an `ambiguous` + status means no proven successful result, **not** proof that nothing was written. + Identity-conflict and missing-receipt messages use shared coordination wording; + their existing reason codes remain stable. + +Existing request digests, receipt schemas, success payloads, no-op consumption, +lease/grant checks, permanent Markdown delivery and provider defaults remain +compatible. The production-scale fixture now drives all seven command operations +through normal, lost-response, unreadable-readback and thrown-response cases, +including replay after an intervening commit. Real File/SQLite/PostgreSQL and +NoKV transport conformance share that matrix. The three-arm read-only source +rehearsal also loses archive responses on actual File/PostgreSQL commits. + +This removes duplicated TS transaction authority, not Python business writers: +no new bridge or RPC is introduced, and cross-runtime calls are unchanged. +T1 metadata/effect closure, T2 retained Monitor leases and D1–D3 qualification +remain separate work; the compatibility editor and other command protocols +retain their distinct receipt contracts. No Goal promotion is implied. + #### Execution cards after the current stack This is a **conditional execution plan**, not a merged-status declaration. diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index 7bdfa4780a..ff82fb149e 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -334,6 +334,42 @@ generation fence、claim/exclusion、capacity 和 PR 等待语义保持。非法 准入。普通 add/update 准入及覆盖全部非法条件的通用修复动作仍是独立范围;不能宣称 全量零行为变化或全部 Todo writer 已闭合。 +#### 命令回执与恢复的统一所有者 + +在基线 `bfd1ec8db`,create、claim、update、complete/supersede、archive 与 +Monitor poll 分别重复 envelope 匹配、结果投影和 CAS 后回读。 +现在由 `coordination/command_receipt.ts` 统一这些语义;各命令继续拥有请求 +规范化/摘要、准入、回执业务载荷校验及状态变更。 +`coordination/todo_archive.ts` 单独拥有归档保留事务,与终态校验和 lease +释放分离。内部调用方直接导入新 owner,旧模块不保留无实际用途的 re-export。 + +这些 canonical 命令路径有以下明确的可观察变化: + +- 提交已 applied 或 ambiguous、但回执不可读时,结果保留为 `ambiguous`, + 携带 `recovery.operation_id` 和 `retry_with_same_operation_id=true`。 + 读取失败不能抹掉可能已经持久化的事实。提交响应抛异常后只查一次回执, + 不自动再次写入。 +- 明确的 CAS conflict 或提交失败,不再被随后的诊断读取失败覆盖。精确的 + 历史回执仍优先返回;applied 响应却缺少回执,仍然是协议失败。 +- create/Monitor 结果对象、update/terminal/archive 的变更判定损坏时,返回 + `invalid_coordination_command_receipt`,不能通过隐式转换变成成功 replay/ + no-op,也不能直接逸出为未处理的解码异常。Claim 保留原回执错误码和历史 + 省略 changed 字段的兼容解析。 +- 读取失败统一携带 `changed=false`;当 status 为 `ambiguous` 时,它表示 + 尚无成功结果证明,**不表示**已证明没有写入。身份冲突与回执缺失的错误消息 + 采用统一 coordination 措辞,原 reason code 保持不变。 + +原请求摘要、回执 schema、成功载荷、no-op 身份消耗、lease/grant 校验、永久 +Markdown 投递和默认 provider 保持兼容。完整生产规模 fixture 现在覆盖七种 +命令的正常提交、响应丢失、回读不可用、响应抛异常,以及插入其他提交后的 +历史重放。真实 File/SQLite/PostgreSQL 和 NoKV transport conformance +共用该矩阵。三路只读源演练还在真实 File/PostgreSQL 归档提交后丢弃响应。 + +本次删除重复的 TS 事务权威,不宣称删除 Python 业务 writer;没有新增 bridge +或 RPC,跨运行时调用数不变。T1 的 metadata/effect 闭合、T2 的带 lease +Monitor 和 D1–D3 资格验证仍待后续;兼容编辑器及其他命令保留各自的回执合同。 +本次不代表 Goal promotion。 + #### 当前 stack 合入后的执行卡 这是**条件式执行规划**,不是所有阶段已完成的声明。2026-09-09 核查时,#4053、#4117、 diff --git a/examples/control_plane/authority-three-arm-rehearsal.py b/examples/control_plane/authority-three-arm-rehearsal.py index 4cc28e7a03..399982d692 100644 --- a/examples/control_plane/authority-three-arm-rehearsal.py +++ b/examples/control_plane/authority-three-arm-rehearsal.py @@ -43,7 +43,7 @@ import {Pool} from 'pg'; import {FileAuthorityStore} from '__FILE_STORE__'; import {PostgreSqlAuthorityStore, installPostgreSqlAuthorityStoreSchema} from '__PG_STORE__'; -import {executeCoordinationTodoArchiveCompleted} from '__TERMINAL__'; +import {executeCoordinationTodoArchiveCompleted} from '__ARCHIVE__'; import {canonicalAuthorityBytes} from '__CODEC__'; import {evaluateTodoResumeConditions} from '__RESUME__'; @@ -117,17 +117,35 @@ next_projection: request.initial, }); assert.equal(initialized.status, 'applied', `${name} initialization failed`); - const archived = await executeCoordinationTodoArchiveCompleted(store, { + // Lose the response after the actual backend commit. Recovery must read + // the original receipt rather than attempt the archive a second time. + let commitCount = 0; + const responseLostStore = { + storeIdentity: () => store.storeIdentity(), + loadAuthority: () => store.loadAuthority(), + readReceipt: (id) => store.readReceipt(id), + scanCommitted: (cursor, limit) => store.scanCommitted(cursor, limit), + commitAuthority: async (request) => { + commitCount++; + const committed = await store.commitAuthority(request); + assert.equal(committed.status, 'applied'); + return {status: 'ambiguous', reason_code: 'synthetic_response_loss', + reason: 'isolated rehearsal discarded the commit response'}; + }, + }; + const archiveRequest = { goal_id: request.goal_id, role: request.role, max_active_done: request.max_active_done, operation_id: `${name}-three-arm-archive`, dry_run: false, now: new Date('2026-01-01T00:00:00Z'), - }); + }; + const archived = await executeCoordinationTodoArchiveCompleted(responseLostStore, archiveRequest); + assert.equal(commitCount, 1); assert.equal( archived.status, - 'applied', + 'recovered', `${name} archive failed (${String(archived.reason_code ?? 'unknown')}:` + `${archiveFailureCategory(archived)})`, ); @@ -149,6 +167,10 @@ assert.deepEqual(end, {status: 'page', transactions: [], next_cursor: finalPage.next_cursor, has_more: false}); assert.deepEqual(await store.loadAuthority(), loaded, 'journal reads changed authority'); + const replay = await executeCoordinationTodoArchiveCompleted(store, archiveRequest); + assert.equal(replay.status, 'replayed'); + assert.equal(replay.cursor, archived.cursor); + assert.deepEqual(await store.loadAuthority(), loaded); results[name] = {archived, head: loaded.head}; } @@ -290,10 +312,10 @@ def _node_script(repository: Path) -> str: ), ) .replace( - "__TERMINAL__", + "__ARCHIVE__", _module_uri( repository, - "loopx/control_plane/coordination/todo_terminal_lifecycle.ts", + "loopx/control_plane/coordination/todo_archive.ts", ), ) .replace( diff --git a/loopx/control_plane/coordination/command_receipt.ts b/loopx/control_plane/coordination/command_receipt.ts new file mode 100644 index 0000000000..fd2791728c --- /dev/null +++ b/loopx/control_plane/coordination/command_receipt.ts @@ -0,0 +1,106 @@ +/** Durable command recovery, owned by coordination rather than by a provider. + * A receipt proves a historical decision, never current execution authority. + * Business planners and request hashes remain with their command owners. */ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult, + AuthorityStoreReceiptResult} from "./authority_store.ts"; +import {AuthorityStoreProtocolError, canonicalAuthorityObject} from "./authority_store_codec.ts"; +import {projectionDelivery} from "../todos/projection_delivery.ts"; + +export type ReceiptPhase = "applied" | "recovered" | "replayed"; +interface ReceiptPayload { + fields: JsonObject; + changed: boolean; +} +interface CommandReceiptContract { + result_schema: S; + identity: JsonObject & {schema_version: string; operation_id: string; goal_id: string; request_sha256: string}; + /** Decode only the command's historical payload, without reading current state. */ + decode(original: JsonObject, phase: ReceiptPhase): ReceiptPayload; + failure(code: string, reason: string): JsonObject & {schema_version: S}; +} +type Result = JsonObject & {schema_version: S}; + +/** One envelope identity, result projection and post-commit state machine for + * canonical Todo commands. Existing wire schemas and request digests are retained. */ +export class CoordinationCommandReceipt { + readonly contract: CommandReceiptContract; + constructor(contract: CommandReceiptContract) { this.contract = contract; } + + private project(receipt: AuthorityStoreReceiptResult, phase: ReceiptPhase): Result | null { + if (receipt.status === "missing") return null; + const {result_schema, identity, decode, failure} = this.contract; + if (receipt.status !== "found") return {schema_version: result_schema, ...receipt, changed: false}; + const original = receipt.receipts[0]; + if (receipt.receipts.length !== 1 || !original || + Object.entries(identity).some(([key, value]) => original[key] !== value)) { + return failure("coordination_operation_identity_mismatch", + "operation id already names a different coordination request"); + } + let payload: ReceiptPayload; + try { + payload = decode(original, phase); + } catch (error) { + if (!(error instanceof AuthorityStoreProtocolError)) throw error; + return failure("invalid_coordination_command_receipt", error.message); + } + return {...payload.fields, schema_version: result_schema, + status: phase === "applied" && !payload.changed ? "no_change" : phase, + changed: phase !== "replayed" && payload.changed, + provider_revision: receipt.provider_revision, cursor: receipt.cursor, + projection_delivery: projectionDelivery(payload.changed), + projection_source: "committed_authority_journal"}; + } + + async read(store: AuthorityStore): Promise | null> { + return this.project(await store.readReceipt(this.contract.identity.operation_id), "replayed"); + } + + async commit(store: AuthorityStore, commit: AuthorityStoreCommit): Promise> { + const {identity, result_schema, failure} = this.contract; + if (commit.operation_id !== identity.operation_id) { + throw new AuthorityStoreProtocolError("commit and receipt operation identities differ"); + } + // Catch only the effect whose response can be lost after durable acceptance. + // Never retry the write here, and never turn an exception into no-write proof. + let committed: AuthorityStoreCommitResult; + try { committed = await store.commitAuthority(commit); } + catch { + committed = {status: "ambiguous", reason_code: "coordination_commit_response_lost", + reason: "commit response was not received; recover the original operation"}; + } + let receipt: AuthorityStoreReceiptResult; + try { receipt = await store.readReceipt(identity.operation_id); } + catch { + receipt = {status: "unavailable", reason_code: "coordination_receipt_read_failed", + reason: "durable receipt read did not complete"}; + } + if (receipt.status === "found") { + return this.project(receipt, committed.status === "applied" ? "applied" : "recovered")!; + } + if (receipt.status === "missing" && committed.status === "applied") { + return failure("coordination_commit_readback_mismatch", "applied command lacks its durable receipt"); + } + if (committed.status === "ambiguous" || + (committed.status === "applied" && receipt.status !== "missing")) { + return {schema_version: result_schema, status: "ambiguous", changed: false, + reason_code: "coordination_receipt_recovery_required", + reason: "commit may be durable; recover its receipt using the same operation id", + commit_status: committed.status, receipt_status: receipt.status, + recovery: {operation_id: identity.operation_id, retry_with_same_operation_id: true}, + ...(receipt.status === "missing" ? {} : {readback_reason_code: receipt.reason_code})}; + } + // A conclusive rejection is not erased by a later diagnostic read failure. + return {schema_version: result_schema, ...committed, changed: false}; + } +} + +/** Terminal/archive results persist an explicit change decision. Missing or + * malformed decisions must not be coerced into a successful no-op. */ +export function commandReceiptResult(original: JsonObject): ReceiptPayload { + const fields = canonicalAuthorityObject(original.result, "command receipt result"); + if (typeof fields.changed !== "boolean") { + throw new AuthorityStoreProtocolError("command receipt result.changed must be boolean"); + } + return {fields: {...fields, original_receipt: original}, changed: fields.changed}; +} diff --git a/loopx/control_plane/coordination/local_archive_attempt.ts b/loopx/control_plane/coordination/local_archive_attempt.ts index b5a09655af..161f9c350e 100644 --- a/loopx/control_plane/coordination/local_archive_attempt.ts +++ b/loopx/control_plane/coordination/local_archive_attempt.ts @@ -16,7 +16,7 @@ import { COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA, executeCoordinationTodoArchiveCompleted, type CoordinationTodoArchiveInput, -} from "./todo_terminal_lifecycle.ts"; +} from "./todo_archive.ts"; const ATTEMPT_SCHEMA = "loopx_local_todo_archive_attempt_v0"; export const LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA = diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index 912789e6be..7ba94c105e 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -1,3 +1,4 @@ +import {COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA} from "./todo_archive.ts"; import {executeTodoContinuation} from "./todo_continuation.ts"; import { withFileMutationLock } from "../effect_runtime_io.ts"; import { ShadowManagementError, requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath } from "./shadow_management.ts"; @@ -59,7 +60,6 @@ import { executeCoordinationTodoUpdate, } from "./todo_update.ts"; import { - COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, executeCoordinationTodoTerminalLifecycle, } from "./todo_terminal_lifecycle.ts"; diff --git a/loopx/control_plane/coordination/todo_archive.ts b/loopx/control_plane/coordination/todo_archive.ts new file mode 100644 index 0000000000..1097be272d --- /dev/null +++ b/loopx/control_plane/coordination/todo_archive.ts @@ -0,0 +1,185 @@ +/** Archive is a retention transaction, separate from completion and its lease + * and validation effects. Both compose the same durable receipt owner. */ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit} from "./authority_store.ts"; +import {AuthorityStoreProtocolError, canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {indexCoordinationProjection, prepareCoordinationProjectionCommit, validateCoordinationTodoReadModel, + type CoordinationProjectionMutation} from "./coordination_projection.ts"; +import {CoordinationCommandReceipt, commandReceiptResult} from "./command_receipt.ts"; +import {selectCoordinationTodoArchive} from "./todo_archive_selection.ts"; + +export const COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA = + "loopx_coordination_todo_archive_result_v0"; +export const COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA = + "loopx_coordination_todo_archive_receipt_v0"; + +export interface CoordinationTodoArchiveInput { + readonly goal_id: string; + readonly role: "agent" | "user"; + readonly max_active_done: number; + readonly operation_id: string; + readonly expected_provider_revision?: string; + readonly dry_run: boolean; + readonly now: Date; +} + +export type CoordinationTodoArchiveResult = JsonObject & { + readonly schema_version: typeof COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA; +}; + +function archiveFailure(code: string, reason: string): CoordinationTodoArchiveResult { + return { + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + status: "failed", + changed: false, + reason_code: code, + reason, + }; +} + +function normalizeArchiveInput(raw: CoordinationTodoArchiveInput): CoordinationTodoArchiveInput { + if (!Number.isSafeInteger(raw.max_active_done) || raw.max_active_done < 0) { + throw new AuthorityStoreProtocolError("max_active_done must be a non-negative safe integer"); + } + if (raw.role !== "agent" && raw.role !== "user") { + throw new AuthorityStoreProtocolError("role must be one of: agent, user"); + } + if (typeof raw.dry_run !== "boolean") throw new AuthorityStoreProtocolError("dry_run must be a boolean"); + if (!(raw.now instanceof Date) || Number.isNaN(raw.now.valueOf())) { + throw new AuthorityStoreProtocolError("now must be a valid Date"); + } + return { + ...raw, + goal_id: requireAuthorityStoreId(raw.goal_id, "goal id"), + role: raw.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", + ), + }), + }; +} + +function archiveReceipt(input: CoordinationTodoArchiveInput, requestSha: string) { + return new CoordinationCommandReceipt({result_schema: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + identity: {schema_version: COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA, + operation_id: input.operation_id, goal_id: input.goal_id, request_sha256: requestSha}, + failure: archiveFailure, + decode(original) { + const result = commandReceiptResult(original); + return {...result, fields: {...result.fields, operation_id: input.operation_id}}; + }}); +} + +/** Archive the oldest canonical completed records while preserving standing decisions. */ +export async function executeCoordinationTodoArchiveCompleted( + store: AuthorityStore, + rawInput: CoordinationTodoArchiveInput, +): Promise { + let input: CoordinationTodoArchiveInput; + try { + input = normalizeArchiveInput(rawInput); + } catch (error) { + return archiveFailure( + "invalid_coordination_todo_archive", + error instanceof Error ? error.message : "invalid Todo archive request", + ); + } + const requestSha = canonicalAuthoritySha256({ + goal_id: input.goal_id, + 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, + }), + }); + // 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 = await archiveReceipt(input, requestSha).read(store); + 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); + validateCoordinationTodoReadModel(head.head, input.goal_id); + } catch (error) { + return archiveFailure( + "invalid_coordination_projection", + error instanceof Error ? error.message : "invalid coordination projection", + ); + } + const selection = selectCoordinationTodoArchive({ + role: input.role, + max_active_done: input.max_active_done, + todos: projection.todo_ids.map((todoId) => projection.todos.get(todoId)!), + }); + const moved = selection.moved_todo_ids.map((todoId) => projection.todos.get(todoId)!); + 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, + max_active_done: selection.max_active_done, + moved_count: selection.moved_count, + moved_todo_ids: selection.moved_todo_ids, + retained_standing_decision_count: selection.retained_standing_decision_count, + }; + if (input.dry_run) { + return { + ...result, + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + status: moved.length > 0 ? "planned" : "no_change", + dry_run: true, + provider_revision: head.provider_revision, + cursor: head.cursor, + }; + } + if (moved.length === 0) { + return { + ...result, + schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, + status: "no_change", + dry_run: false, + provider_revision: head.provider_revision, + cursor: head.cursor, + }; + } + const mutations: CoordinationProjectionMutation[] = moved.map((todo) => ({ + kind: "todo_upsert", + todo: {...todo, archive_state: "archive", updated_at: updatedAt}, + })); + const commit: AuthorityStoreCommit = prepareCoordinationProjectionCommit({ + goal_id: input.goal_id, + operation_id: input.operation_id, + expected_provider_revision: head.provider_revision, + projection: head.head, + mutations, + }); + commit.receipts = [{ + schema_version: COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA, + operation_id: input.operation_id, + goal_id: input.goal_id, + request_sha256: requestSha, + result, + }]; + return archiveReceipt(input, requestSha).commit(store, commit); +} diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index 48e2ec7a9f..552d93dcb2 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -1,6 +1,6 @@ import {leaseOwnerRejection as ownerRejection} from "../work_items/task_lease_eligibility.ts"; import type { JsonObject } from "../effect_program.ts"; -import type { AuthorityStore, AuthorityStoreCommit, AuthorityStoreReceiptResult } from "./authority_store.ts"; +import type { AuthorityStore, AuthorityStoreCommit } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityObject, @@ -8,7 +8,7 @@ import { requireAuthorityStoreId, } from "./authority_store_codec.ts"; import {validateContinuationNote, computeContinuationTodoFacts} from "./continuation_note.ts"; -import {projectionDelivery} from "../todos/projection_delivery.ts"; +import {CoordinationCommandReceipt} from "./command_receipt.ts"; import {normalizeRegisteredTodoAgents, normalizeTodoAgent} from "./todo_agents.ts"; import { prepareCoordinationProjectionCommit, @@ -406,25 +406,13 @@ export async function executeCoordinationTodoClaim( dry_run: input.dry_run, ...(leaseRequest === null ? {} : {lease_request: leaseRequest}), }); - const replay = ( - receipt: AuthorityStoreReceiptResult, - status: "replayed" | "applied" | "recovered", - ): CoordinationTodoClaimResult | null => { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") { - return { schema_version: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, ...receipt }; - } - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || - original?.schema_version !== COORDINATION_TODO_CLAIM_RECEIPT_SCHEMA || - original.operation_id !== input.operation_id || original.goal_id !== input.goal_id || - original.request_sha256 !== requestSha) { - return failure("coordination_operation_identity_mismatch", - "operation id already names a different coordination request"); - } - let result: JsonObject; - try { - result = canonicalAuthorityObject(original.result, "original claim result"); + const receipt = new CoordinationCommandReceipt({result_schema: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, + identity: {schema_version: COORDINATION_TODO_CLAIM_RECEIPT_SCHEMA, + operation_id: input.operation_id, goal_id: input.goal_id, request_sha256: requestSha}, + failure: (code, reason) => failure(code === "invalid_coordination_command_receipt" + ? "invalid_coordination_todo_claim_receipt" : code, reason), + decode(original) { + const result = canonicalAuthorityObject(original.result, "original claim result"); if (result.todo_id !== input.todo_id || result.claimed_by !== input.claimed_by || (result.changed !== undefined && typeof result.changed !== "boolean") || (result.changed !== false && typeof result.updated_at !== "string") || @@ -439,23 +427,10 @@ export async function executeCoordinationTodoClaim( throw new AuthorityStoreProtocolError("original claim lease identity is invalid"); } } - } catch (error) { - return failure("invalid_coordination_todo_claim_receipt", - error instanceof Error ? error.message : "invalid claim receipt"); - } - return { - ...result, - schema_version: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, - status: status === "applied" && result.changed === false ? "no_change" : status, - changed: status !== "replayed" && result.changed !== false, - provider_revision: receipt.provider_revision, - cursor: receipt.cursor, - original_receipt: original, - projection_delivery: projectionDelivery(result.changed !== false), - projection_source: "committed_authority_journal", - }; - }; - const existing = replay(await store.readReceipt(input.operation_id), "replayed"); + // Pre-change claim receipts may omit changed; those encoded a real mutation. + return {fields: {...result, original_receipt: original}, changed: result.changed !== false}; + }}); + const existing = await receipt.read(store); if (existing !== null) return existing; const head = await store.loadAuthority(); @@ -726,11 +701,5 @@ export async function executeCoordinationTodoClaim( request_sha256: requestSha, result, }]; - const committed = await store.commitAuthority(commit); - const readback = replay(await store.readReceipt(input.operation_id), - committed.status === "applied" ? "applied" : "recovered"); - if (readback !== null) return readback; - return committed.status === "applied" - ? failure("coordination_commit_readback_mismatch", "applied claim lacks its durable receipt") - : { schema_version: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, ...committed, changed: false }; + return receipt.commit(store, commit); } diff --git a/loopx/control_plane/coordination/todo_create.ts b/loopx/control_plane/coordination/todo_create.ts index c985daa1be..2663bfe8c4 100644 --- a/loopx/control_plane/coordination/todo_create.ts +++ b/loopx/control_plane/coordination/todo_create.ts @@ -1,6 +1,6 @@ import type { JsonObject } from "../effect_program.ts"; -import {projectionDelivery} from "../todos/projection_delivery.ts"; -import type { AuthorityStore, AuthorityStoreReceiptResult } from "./authority_store.ts"; +import {CoordinationCommandReceipt} from "./command_receipt.ts"; +import type { AuthorityStore } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityObject, @@ -48,38 +48,16 @@ function failure(code: string, reason: string, detail: JsonObject = {}): Coordin }; } -function replayCreate( - receipt: AuthorityStoreReceiptResult, - input: CoordinationTodoCreateInput, - requestSha: string, - status: "replayed" | "applied" | "recovered", -): CoordinationTodoCreateResult | null { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") { - return {schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, ...receipt}; - } - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || - original?.schema_version !== COORDINATION_TODO_CREATE_RECEIPT_SCHEMA || - original.operation_id !== input.operation_id || original.goal_id !== input.goal_id || - original.request_sha256 !== requestSha || original.todo_id !== input.todo.todo_id) { - return failure( - "coordination_operation_identity_mismatch", - "operation id already names a different Todo create request", - ); - } - return { - schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, - status, - changed: status !== "replayed", - todo_id: original.todo_id, - todo: original.todo, - provider_revision: receipt.provider_revision, - cursor: receipt.cursor, - original_receipt: original, - projection_delivery: projectionDelivery(true), - projection_source: "committed_authority_journal", - }; +function createReceipt(input: CoordinationTodoCreateInput, requestSha: string) { + return new CoordinationCommandReceipt({result_schema: COORDINATION_TODO_CREATE_RESULT_SCHEMA, + identity: {schema_version: COORDINATION_TODO_CREATE_RECEIPT_SCHEMA, + operation_id: input.operation_id, goal_id: input.goal_id, todo_id: String(input.todo.todo_id), + request_sha256: requestSha}, failure, + decode(original) { + const todo = canonicalAuthorityObject(original.todo, "created Todo receipt"); + if (todo.todo_id !== input.todo.todo_id) throw new AuthorityStoreProtocolError("created Todo receipt identity mismatch"); + return {fields: {todo_id: original.todo_id, todo, original_receipt: original}, changed: true}; + }}); } function normalizeCreateInput(rawInput: CoordinationTodoCreateInput): CoordinationTodoCreateInput { @@ -204,15 +182,7 @@ async function commitCreate( todo_id: todoId, todo: created, }]; - const committed = await store.commitAuthority(commit); - const readback = replayCreate( - await store.readReceipt(input.operation_id), input, requestSha, - committed.status === "applied" ? "applied" : "recovered", - ); - if (readback !== null) return readback; - return committed.status === "applied" - ? failure("coordination_commit_readback_mismatch", "applied create lacks its durable receipt") - : {schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, ...committed, changed: false}; + return createReceipt(input, requestSha).commit(store, commit); } /** Create one canonical work item through the provider transaction and outbox. */ @@ -236,9 +206,7 @@ export async function executeCoordinationTodoCreate( actor_agent_id: input.actor_agent_id, dry_run: input.dry_run, }); - const existing = replayCreate( - await store.readReceipt(input.operation_id), input, requestSha, "replayed", - ); + const existing = await createReceipt(input, requestSha).read(store); if (existing !== null) return existing; const head = await store.loadAuthority(); diff --git a/loopx/control_plane/coordination/todo_monitor_poll.ts b/loopx/control_plane/coordination/todo_monitor_poll.ts index b5129d7c3a..cc9027e143 100644 --- a/loopx/control_plane/coordination/todo_monitor_poll.ts +++ b/loopx/control_plane/coordination/todo_monitor_poll.ts @@ -1,8 +1,8 @@ /** One canonical transaction for an observation and its independent successors. * Network polling, quota settlement and display delivery are separate effects. */ import type {JsonObject} from "../effect_program.ts"; -import type {AuthorityStore, AuthorityStoreReceiptResult} from "./authority_store.ts"; -import {canonicalAuthorityObject, canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import type {AuthorityStore} from "./authority_store.ts"; +import {AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; import {normalizeRegisteredTodoAgents, normalizeTodoAgent} from "./todo_agents.ts"; import {indexCoordinationProjection, prepareCoordinationProjectionCommit, validateCoordinationTodoReadModel} from "./coordination_projection.ts"; import {TODO_DOMAIN_ITEM_SCHEMA} from "./coordination_state_contract.ts"; @@ -11,7 +11,7 @@ import {planMonitorMetadata, TODO_MONITOR_METADATA_REQUEST_SCHEMA} from "../todo import {planMonitorSuccessor, selectMonitorTodo, MONITOR_SUCCESSOR_REQUEST_SCHEMA} from "../scheduler/monitor_successor.ts"; import {optionalNonEmptyString, requireBoolean} from "../runtime_decode.ts"; import {planTodoAuthoringScope, TODO_AUTHORING_SCOPE_REQUEST_SCHEMA} from "../todos/authoring_scope.ts"; -import {projectionDelivery} from "../todos/projection_delivery.ts"; +import {CoordinationCommandReceipt} from "./command_receipt.ts"; export const COORDINATION_MONITOR_POLL_REQUEST_SCHEMA = "loopx_coordination_monitor_poll_request_v0"; export const COORDINATION_MONITOR_POLL_RESULT_SCHEMA = "loopx_coordination_monitor_poll_result_v0"; @@ -27,23 +27,23 @@ export interface CoordinationMonitorPollInput { intent: JsonObject; } -function failure(reason_code: string, reason: string): JsonObject { +function failure(reason_code: string, reason: string): JsonObject & {schema_version: typeof COORDINATION_MONITOR_POLL_RESULT_SCHEMA} { return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status: "failed", changed: false, reason_code, reason}; } -function replay(receipt: AuthorityStoreReceiptResult, input: CoordinationMonitorPollInput, - hash: string, status: "applied" | "recovered" | "replayed"): JsonObject | null { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...receipt}; - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || original?.schema_version !== RECEIPT_SCHEMA || - original.goal_id !== input.goal_id || original.operation_id !== input.operation_id || original.request_sha256 !== hash) { - return failure("coordination_operation_identity_mismatch", "operation id already names a different Monitor observation or successor intent"); - } - return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status, - changed: status !== "replayed", provider_revision: receipt.provider_revision, cursor: receipt.cursor, - writeback: {...canonicalAuthorityObject(original.writeback, "Monitor writeback"), provider_replayed: status === "replayed"}, - projection_delivery: projectionDelivery(true), projection_source: "committed_authority_journal"}; +function monitorReceipt(input: CoordinationMonitorPollInput, hash: string) { + return new CoordinationCommandReceipt({result_schema: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, + identity: {schema_version: RECEIPT_SCHEMA, goal_id: input.goal_id, + operation_id: input.operation_id, request_sha256: hash}, failure, + decode(original, phase) { + const writeback = canonicalAuthorityObject(original.writeback, "Monitor writeback"); + if (writeback.schema_version !== "monitor_poll_todo_writeback_v0" || + writeback.goal_id !== input.goal_id || writeback.monitor_effect_id !== input.operation_id || + !Array.isArray(writeback.next_todos)) { + throw new AuthorityStoreProtocolError("Monitor receipt writeback identity or successors invalid"); + } + return {fields: {writeback: {...writeback, provider_replayed: phase === "replayed"}}, changed: true}; + }}); } function normalize(raw: CoordinationMonitorPollInput): CoordinationMonitorPollInput { @@ -156,7 +156,8 @@ export async function executeCoordinationMonitorPoll(store: AuthorityStore, // Original wire identity, before any normalization/default route inference. const hash = canonicalAuthoritySha256({goal_id: input.goal_id, observation: input.observation, intent: input.intent, actor_agent_id: input.actor_agent_id, dry_run: input.dry_run}); - const previous = replay(await store.readReceipt(input.operation_id), input, hash, "replayed"); + const receipt = monitorReceipt(input, hash); + const previous = await receipt.read(store); if (previous) return previous; const head = await store.loadAuthority(); if (head.status !== "loaded") return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...head}; @@ -169,10 +170,5 @@ export async function executeCoordinationMonitorPoll(store: AuthorityStore, expected_provider_revision: head.provider_revision, projection: head.head, mutations: plan.mutations}); commit.receipts = [{schema_version: RECEIPT_SCHEMA, goal_id: input.goal_id, operation_id: input.operation_id, request_sha256: hash, writeback: plan.writeback}]; - const committed = await store.commitAuthority(commit); - const result = replay(await store.readReceipt(input.operation_id), input, hash, - committed.status === "applied" ? "applied" : "recovered"); - return result ?? (committed.status === "applied" - ? failure("coordination_commit_readback_mismatch", "applied Monitor transaction lacks its durable receipt") - : {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...committed, changed: false}); + return receipt.commit(store, commit); } diff --git a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts index 55a0d57aa5..b0abdfb327 100644 --- a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts +++ b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts @@ -1,11 +1,10 @@ import { createHash } from "node:crypto"; -import {projectionDelivery} from "../todos/projection_delivery.ts"; +import {CoordinationCommandReceipt, commandReceiptResult} from "./command_receipt.ts"; import type { JsonObject } from "../effect_program.ts"; import type { AuthorityStore, AuthorityStoreCommit, - AuthorityStoreReceiptResult, } from "./authority_store.ts"; import { AuthorityStoreProtocolError, @@ -46,7 +45,6 @@ import { normalizeAgent, normalizeWriteScopes, } from "../work_items/task_lease_acquire.ts"; -import { selectCoordinationTodoArchive } from "./todo_archive_selection.ts"; import { userTodoScopeConflict, USER_TODO_TASK_CLASSES } from "../todos/authoring_scope.ts"; import { deriveCoordinationTodoSuccessorProposals, @@ -57,11 +55,6 @@ export const COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA = "loopx_coordination_todo_terminal_lifecycle_result_v0"; export const COORDINATION_TODO_TERMINAL_LIFECYCLE_RECEIPT_SCHEMA = "loopx_coordination_todo_terminal_lifecycle_receipt_v0"; -export const COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA = - "loopx_coordination_todo_archive_result_v0"; -export const COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA = - "loopx_coordination_todo_archive_receipt_v0"; - const TERMINAL_COMMANDS = ["complete", "supersede"] as const; const TODO_ROLES = ["agent", "user"] as const; const DECISION_OUTCOMES = ["approve", "reject", "cancel"] as const; @@ -105,23 +98,9 @@ export interface CoordinationTodoTerminalLifecycleInput { readonly now: Date; } -export interface CoordinationTodoArchiveInput { - readonly goal_id: string; - readonly role: TodoRole; - readonly max_active_done: number; - readonly operation_id: string; - readonly expected_provider_revision?: string; - readonly dry_run: boolean; - readonly now: Date; -} - export type CoordinationTodoTerminalLifecycleResult = JsonObject & { readonly schema_version: typeof COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA; }; -export type CoordinationTodoArchiveResult = JsonObject & { - readonly schema_version: typeof COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA; -}; - type CoordinationTodoTerminalFailureKind = | "decision_rejection" | "protocol_failure"; @@ -368,16 +347,6 @@ function terminalFailure( }; } -function archiveFailure(code: string, reason: string): CoordinationTodoArchiveResult { - return { - schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, - status: "failed", - changed: false, - reason_code: code, - reason, - }; -} - function terminalRequestSha(input: CoordinationTodoTerminalLifecycleInput): string { const completionPolicyIdentity = input.completion_policy_request === null ? null @@ -416,45 +385,14 @@ function terminalRequestSha(input: CoordinationTodoTerminalLifecycleInput): stri }); } -function replayTerminal( - receipt: AuthorityStoreReceiptResult, - input: CoordinationTodoTerminalLifecycleInput, - requestSha: string, - status: "replayed" | "applied" | "recovered", -): CoordinationTodoTerminalLifecycleResult | null { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") { - return { - schema_version: COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, - ...receipt, - changed: false, - }; - } - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || - original?.schema_version !== COORDINATION_TODO_TERMINAL_LIFECYCLE_RECEIPT_SCHEMA || - original.operation_id !== input.operation_id || original.goal_id !== input.goal_id || - original.todo_id !== input.todo_id || original.command !== input.command || - original.request_sha256 !== requestSha) { - return terminalFailure( - "coordination_operation_identity_mismatch", - "operation id already names a different Todo terminal lifecycle request", - {}, - "decision_rejection", - ); - } - const result = canonicalAuthorityObject(original.result, "terminal lifecycle receipt result"); - return { - ...result, - schema_version: COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, - status: status === "applied" && result.changed === false ? "no_change" : status, - changed: status !== "replayed" && result.changed === true, - provider_revision: receipt.provider_revision, - cursor: receipt.cursor, - original_receipt: original, - projection_delivery: projectionDelivery(result.changed === true), - projection_source: "committed_authority_journal", - }; +function terminalReceipt(input: CoordinationTodoTerminalLifecycleInput, requestSha: string) { + return new CoordinationCommandReceipt({result_schema: COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, + identity: {schema_version: COORDINATION_TODO_TERMINAL_LIFECYCLE_RECEIPT_SCHEMA, + operation_id: input.operation_id, goal_id: input.goal_id, todo_id: input.todo_id, + command: input.command, request_sha256: requestSha}, + failure: (code, reason) => terminalFailure(code, reason, {}, + code === "coordination_operation_identity_mismatch" ? "decision_rejection" : "protocol_failure"), + decode: commandReceiptResult}); } async function commitTerminalResult( @@ -499,24 +437,7 @@ async function commitTerminalResult( request_sha256: requestSha, result, }]; - const committed = await store.commitAuthority(commit); - const readback = replayTerminal( - await store.readReceipt(input.operation_id), - input, - requestSha, - committed.status === "applied" ? "applied" : "recovered", - ); - if (readback !== null) return readback; - return committed.status === "applied" - ? terminalFailure( - "coordination_commit_readback_mismatch", - "applied terminal lifecycle mutation lacks its durable receipt", - ) - : { - schema_version: COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, - ...committed, - changed: false, - }; + return terminalReceipt(input, requestSha).commit(store, commit); } function todoFact(todo: JsonObject): JsonObject { @@ -742,12 +663,7 @@ export async function executeCoordinationTodoTerminalLifecycle( ); } const requestSha = terminalRequestSha(input); - const replay = replayTerminal( - await store.readReceipt(input.operation_id), - input, - requestSha, - "replayed", - ); + const replay = await terminalReceipt(input, requestSha).read(store); if (replay !== null) return replay; const head = await store.loadAuthority(); @@ -1083,184 +999,3 @@ export async function executeCoordinationTodoTerminalLifecycle( ] : []; return commitTerminalResult(store, input, requestSha, head, result, mutations); } - -function normalizeArchiveInput(raw: CoordinationTodoArchiveInput): CoordinationTodoArchiveInput { - if (!Number.isSafeInteger(raw.max_active_done) || raw.max_active_done < 0) { - throw new AuthorityStoreProtocolError("max_active_done must be a non-negative safe integer"); - } - return { - ...raw, - 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"), - }; -} - -function replayArchive( - receipt: AuthorityStoreReceiptResult, - input: CoordinationTodoArchiveInput, - requestSha: string, - status: "replayed" | "applied" | "recovered", -): CoordinationTodoArchiveResult | null { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") { - return {schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...receipt, changed: false}; - } - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || - original?.schema_version !== COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA || - original.operation_id !== input.operation_id || original.goal_id !== input.goal_id || - original.request_sha256 !== requestSha) { - return archiveFailure( - "coordination_operation_identity_mismatch", - "operation id already names a different Todo archive request", - ); - } - const result = canonicalAuthorityObject(original.result, "Todo archive receipt result"); - 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, - cursor: receipt.cursor, - original_receipt: original, - projection_delivery: projectionDelivery(result.changed === true), - projection_source: "committed_authority_journal", - }; -} - -/** Archive the oldest canonical completed records while preserving standing decisions. */ -export async function executeCoordinationTodoArchiveCompleted( - store: AuthorityStore, - rawInput: CoordinationTodoArchiveInput, -): Promise { - let input: CoordinationTodoArchiveInput; - try { - input = normalizeArchiveInput(rawInput); - } catch (error) { - return archiveFailure( - "invalid_coordination_todo_archive", - error instanceof Error ? error.message : "invalid Todo archive request", - ); - } - const requestSha = canonicalAuthoritySha256({ - goal_id: input.goal_id, - 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, - }), - }); - // 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); - validateCoordinationTodoReadModel(head.head, input.goal_id); - } catch (error) { - return archiveFailure( - "invalid_coordination_projection", - error instanceof Error ? error.message : "invalid coordination projection", - ); - } - const selection = selectCoordinationTodoArchive({ - role: input.role, - max_active_done: input.max_active_done, - todos: projection.todo_ids.map((todoId) => projection.todos.get(todoId)!), - }); - const moved = selection.moved_todo_ids.map((todoId) => projection.todos.get(todoId)!); - 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, - max_active_done: selection.max_active_done, - moved_count: selection.moved_count, - moved_todo_ids: selection.moved_todo_ids, - retained_standing_decision_count: selection.retained_standing_decision_count, - }; - if (input.dry_run) { - return { - ...result, - schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, - status: moved.length > 0 ? "planned" : "no_change", - dry_run: true, - provider_revision: head.provider_revision, - cursor: head.cursor, - }; - } - if (moved.length === 0) { - return { - ...result, - schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, - status: "no_change", - dry_run: false, - provider_revision: head.provider_revision, - cursor: head.cursor, - }; - } - const mutations: CoordinationProjectionMutation[] = moved.map((todo) => ({ - kind: "todo_upsert", - todo: {...todo, archive_state: "archive", updated_at: updatedAt}, - })); - const commit: AuthorityStoreCommit = prepareCoordinationProjectionCommit({ - goal_id: input.goal_id, - operation_id: input.operation_id, - expected_provider_revision: head.provider_revision, - projection: head.head, - mutations, - }); - commit.receipts = [{ - schema_version: COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA, - operation_id: input.operation_id, - goal_id: input.goal_id, - request_sha256: requestSha, - result, - }]; - const committed = await store.commitAuthority(commit); - const readback = replayArchive( - await store.readReceipt(input.operation_id), - input, - requestSha, - committed.status === "applied" ? "applied" : "recovered", - ); - if (readback !== null) return readback; - return committed.status === "applied" - ? archiveFailure( - "coordination_commit_readback_mismatch", - "applied Todo archive mutation lacks its durable receipt", - ) - : {schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...committed, changed: false}; -} diff --git a/loopx/control_plane/coordination/todo_update.ts b/loopx/control_plane/coordination/todo_update.ts index cd8b2de078..768b9ed3db 100644 --- a/loopx/control_plane/coordination/todo_update.ts +++ b/loopx/control_plane/coordination/todo_update.ts @@ -1,7 +1,7 @@ import type { JsonObject } from "../effect_program.ts"; import { TODO_WORK_REQUIREMENT_FIELDS } from "../todos/work_requirements.ts"; import { TODO_OWNERSHIP_INTENT_FIELDS } from "../todos/authoring_scope.ts"; -import type { AuthorityStore, AuthorityStoreCommit, AuthorityStoreReceiptResult } from "./authority_store.ts"; +import type { AuthorityStore, AuthorityStoreCommit } from "./authority_store.ts"; import { AuthorityStoreProtocolError, canonicalAuthorityBytes, @@ -23,7 +23,7 @@ import { evaluateCoordinationTerminalFence, COORDINATION_TERMINAL_FENCE_REQUEST_ import { leaseEpoch } from "../work_items/task_lease_acquire.ts"; import { parseIsoTimestamp } from "../runtime_timestamp.ts"; import { normalizeNativePlanningIntent, planNativeTodoUpdate } from "../todos/native_update_plan.ts"; -import { projectionDelivery } from "../todos/projection_delivery.ts"; +import {CoordinationCommandReceipt} from "./command_receipt.ts"; export const COORDINATION_TODO_UPDATE_REQUEST_SCHEMA = "loopx_local_coordination_todo_update_request_v0"; @@ -114,30 +114,15 @@ function normalizeInput(raw: CoordinationTodoUpdateInput): CoordinationTodoUpdat patch, clear_fields: clearFields}; } -function replayUpdate( - receipt: AuthorityStoreReceiptResult, input: CoordinationTodoUpdateInput, - requestSha: string, status: "replayed" | "applied" | "recovered", -): CoordinationTodoUpdateResult | null { - if (receipt.status === "missing") return null; - if (receipt.status !== "found") { - return {schema_version: COORDINATION_TODO_UPDATE_RESULT_SCHEMA, ...receipt, changed: false}; - } - const original = receipt.receipts[0]; - if (receipt.receipts.length !== 1 || - original?.schema_version !== COORDINATION_TODO_UPDATE_RECEIPT_SCHEMA || - original.operation_id !== input.operation_id || original.goal_id !== input.goal_id || - original.todo_id !== input.todo_id || original.request_sha256 !== requestSha || - typeof original.changed !== "boolean") { - return failure("coordination_operation_identity_mismatch", - "operation id already names a different Todo update request"); - } - return {schema_version: COORDINATION_TODO_UPDATE_RESULT_SCHEMA, - status: status === "applied" && !original.changed ? "no_change" : status, - changed: status !== "replayed" && original.changed, - todo_id: input.todo_id, provider_revision: receipt.provider_revision, - cursor: receipt.cursor, original_receipt: original, - projection_delivery: projectionDelivery(original.changed), - projection_source: "committed_authority_journal"}; +function updateReceipt(input: CoordinationTodoUpdateInput, requestSha: string) { + return new CoordinationCommandReceipt({result_schema: COORDINATION_TODO_UPDATE_RESULT_SCHEMA, + identity: {schema_version: COORDINATION_TODO_UPDATE_RECEIPT_SCHEMA, + operation_id: input.operation_id, goal_id: input.goal_id, todo_id: input.todo_id, + request_sha256: requestSha}, failure, + decode(original) { + if (typeof original.changed !== "boolean") throw new AuthorityStoreProtocolError("update receipt changed must be boolean"); + return {fields: {todo_id: input.todo_id, original_receipt: original}, changed: original.changed}; + }}); } function updateRequestSha(input: CoordinationTodoUpdateInput): string { @@ -309,7 +294,8 @@ export async function executeCoordinationTodoUpdate( error instanceof Error ? error.message : "invalid Todo update"); } const requestSha = updateRequestSha(input); - const replay = replayUpdate(await store.readReceipt(input.operation_id), input, requestSha, "replayed"); + const receipt = updateReceipt(input, requestSha); + const replay = await receipt.read(store); if (replay !== null) return replay; if (input.actor_agent_id === null || !input.registered_agents.includes(input.actor_agent_id)) { return failure("actor_not_registered", "Todo update requires a registered actor"); @@ -343,11 +329,5 @@ export async function executeCoordinationTodoUpdate( commit.receipts = [{schema_version: COORDINATION_TODO_UPDATE_RECEIPT_SCHEMA, operation_id: input.operation_id, goal_id: input.goal_id, todo_id: input.todo_id, request_sha256: requestSha, changed}]; - const committed = await store.commitAuthority(commit); - const readback = replayUpdate(await store.readReceipt(input.operation_id), input, requestSha, - committed.status === "applied" ? "applied" : "recovered"); - if (readback !== null) return readback; - return committed.status === "applied" - ? failure("coordination_commit_readback_mismatch", "applied update lacks its durable receipt") - : {schema_version: COORDINATION_TODO_UPDATE_RESULT_SCHEMA, ...committed, changed: false}; + return receipt.commit(store, commit); } diff --git a/tests/control_plane/test_local_authority_shadow_cli_e2e.py b/tests/control_plane/test_local_authority_shadow_cli_e2e.py index 6d575a80e6..e0c112464b 100644 --- a/tests/control_plane/test_local_authority_shadow_cli_e2e.py +++ b/tests/control_plane/test_local_authority_shadow_cli_e2e.py @@ -13,6 +13,9 @@ REPO_ROOT = Path(__file__).resolve().parents[2] +# Keep the crash-gap observation window aligned with the CLI subprocess timeout; +# importing the Python CLI can exceed a short local polling budget on CI. +CLI_TIMEOUT_SECONDS = 30 def _workspace(tmp_path: Path, *, goal_id: str) -> tuple[Path, Path, Path]: @@ -82,7 +85,7 @@ def _cli(registry: Path, runtime_root: Path, *args: str) -> dict[str, object]: check=True, capture_output=True, text=True, - timeout=30, + timeout=CLI_TIMEOUT_SECONDS, ) return json.loads(completed.stdout) @@ -312,12 +315,22 @@ def test_product_cli_loses_capture_between_commit_and_observer_then_refreshes_sn stderr=subprocess.PIPE, text=True, ) - deadline = time.monotonic() + 5.0 + deadline = time.monotonic() + CLI_TIMEOUT_SECONDS while first_text not in state.read_text(encoding="utf-8"): + returncode = process.poll() + if returncode is not None: + stdout, stderr = process.communicate(timeout=5) + raise AssertionError( + "primary Todo CLI exited before its commit became visible " + f"(returncode={returncode}, stdout={stdout!r}, stderr={stderr!r})" + ) if time.monotonic() >= deadline: process.kill() - process.communicate(timeout=5) - raise AssertionError("primary Todo commit did not become visible") + stdout, stderr = process.communicate(timeout=5) + raise AssertionError( + "primary Todo commit did not become visible within " + f"{CLI_TIMEOUT_SECONDS}s (stdout={stdout!r}, stderr={stderr!r})" + ) time.sleep(0.01) assert process.poll() is None process.kill() diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 582ed8d9bc..88670a9378 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,7 +1,9 @@ import {registerAuthorityScanConformance} from "./authority_scan_conformance.ts"; +import {executeCoordinationTodoArchiveCompleted} from "../../loopx/control_plane/coordination/todo_archive.ts"; import assert from "node:assert/strict"; import { createHash } from "node:crypto"; import test from "node:test"; +import {registerCoordinationReceiptConformance} from "./coordination_receipt_conformance.ts"; import {registerNativePlanningUpdateConformance} from "./native_planning_update_conformance.ts"; import type { @@ -35,7 +37,6 @@ import {projectStandingDecisions} from "../../loopx/control_plane/todos/standing import {evaluateTodoResumeConditions} from "../../loopx/control_plane/todos/resume_condition.ts"; import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; import { - executeCoordinationTodoArchiveCompleted, executeCoordinationTodoTerminalLifecycle, } from "../../loopx/control_plane/coordination/todo_terminal_lifecycle.ts"; import { editCoordinationTodo, TODO_COMPATIBILITY_EDIT_SCHEMA } from "../../loopx/control_plane/coordination/todo_compatibility_edit.ts"; @@ -207,6 +208,7 @@ export function registerAuthorityStoreConformance( ): void { registerAuthorityScanConformance(providerName, factory); registerNativePlanningUpdateConformance(providerName, factory); + registerCoordinationReceiptConformance(providerName, factory); for (const native of [false, true]) test(`${providerName} conformance: standing revocation survives canonical ordering and archive (${native ? "native" : "legacy"})`, async (t) => { const {store} = await factory(t); const goal = "goal-standing"; diff --git a/tests/control_plane_ts/command_receipt.test.ts b/tests/control_plane_ts/command_receipt.test.ts new file mode 100644 index 0000000000..ef487643fb --- /dev/null +++ b/tests/control_plane_ts/command_receipt.test.ts @@ -0,0 +1,118 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {AuthorityStore, AuthorityStoreCommit, AuthorityStoreCommitResult, AuthorityStoreReceiptResult} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {CoordinationCommandReceipt, commandReceiptResult} from "../../loopx/control_plane/coordination/command_receipt.ts"; + +const identity = {schema_version: "test_receipt_v0", operation_id: "original-operation", + goal_id: "goal-a", request_sha256: "intent-digest"}; +const original = {...identity, result: {changed: true, todo_id: "todo-a"}}; +const found: AuthorityStoreReceiptResult = {status: "found", provider_revision: "historical-revision", + cursor: "historical-cursor", receipts: [original]}; +const commit: AuthorityStoreCommit = {operation_id: identity.operation_id, + expected_provider_revision: "before", next_projection: {}, events: [], receipts: [original]}; +const applied: AuthorityStoreCommitResult = {status: "applied", provider_revision: "historical-revision", cursor: "historical-cursor"}; +const unavailable: AuthorityStoreReceiptResult = {status: "unavailable", reason_code: "disconnected", reason: "synthetic disconnect"}; +const receipt = () => new CoordinationCommandReceipt({result_schema: "test_result_v0", identity, + decode: commandReceiptResult, failure: (reason_code, reason) => + ({schema_version: "test_result_v0", status: "failed", changed: false, reason_code, reason})}); +function store(result: AuthorityStoreCommitResult, readback: AuthorityStoreReceiptResult): AuthorityStore { + return {storeIdentity: async () => {throw new Error("receipt recovery cannot inspect identity");}, + loadAuthority: async () => {throw new Error("receipt recovery cannot reinterpret the current head");}, + scanCommitted: async () => {throw new Error("receipt recovery cannot scan unrelated transactions");}, + commitAuthority: async () => result, readReceipt: async id => { + assert.equal(id, identity.operation_id); return readback; + }}; +} + +test("receipt identity rejects every cross-command, cross-goal and changed-intent alias", async () => { + for (const field of Object.keys(identity)) { + const result = await receipt().read(store(applied, {...found, receipts: [{...original, [field]: "different"}]})); + assert.equal(result?.reason_code, "coordination_operation_identity_mismatch", field); + } + for (const receipts of [[], [original, original]]) { + assert.equal((await receipt().read(store(applied, {...found, receipts})))?.reason_code, + "coordination_operation_identity_mismatch"); + } +}); + +test("receipt read is a trust boundary: malformed result cannot become success or no-op", async () => { + for (const result of [null, [], "invalid", {}, {changed: "false"}, {changed: 0}]) { + const readback = {...found, receipts: [{...identity, result}]} as AuthorityStoreReceiptResult; + const response = await receipt().read(store(applied, readback)); + assert.equal(response?.status, "failed"); + assert.equal(response?.reason_code, "invalid_coordination_command_receipt"); + } +}); + +test("historical no-op stays no-op on acceptance and never gains projection delivery", async () => { + const readback = {...found, receipts: [{...identity, result: {changed: false}}]}; + const target = store(applied, readback); + const result = await receipt().commit(target, commit); + assert.equal(result.status, "no_change"); + assert.equal(result.projection_delivery, "not_required"); + assert.equal(result.changed, false); + assert.equal((await receipt().read(target))?.changed, false); + assert.equal((await receipt().read(target))?.status, "replayed"); +}); + +test("pre-write receipt failure blocks planning and is not post-commit uncertainty", async () => { + assert.equal((await receipt().read(store(applied, unavailable)))?.status, "unavailable"); + assert.equal(await receipt().read(store(applied, {status: "missing"})), null); +}); + +test("only an exact durable receipt resolves a lost or rejected commit response", async () => { + const results: AuthorityStoreCommitResult[] = [applied, + {status: "ambiguous", reason_code: "lost", reason: "response lost"}, + {status: "failed", reason_code: "rejected", reason: "write rejected"}, + {status: "conflict", conflict_kind: "operation_id_exists", current_provider_revision: "later", current_cursor: "later"}]; + for (const result of results) { + const response = await receipt().commit(store(result, found), commit); + assert.equal(response.status, result.status === "applied" ? "applied" : "recovered"); + assert.equal(response.provider_revision, found.provider_revision); + assert.equal(response.cursor, found.cursor); + assert.deepEqual(response.original_receipt, original); + } +}); + +test("readback failure never erases a conclusive CAS conflict or rejection", async () => { + for (const result of [ + {status: "conflict", conflict_kind: "provider_revision_mismatch", current_provider_revision: "other", current_cursor: "other"}, + {status: "failed", reason_code: "denied", reason: "write denied"}, + ] satisfies AuthorityStoreCommitResult[]) { + const response = await receipt().commit(store(result, unavailable), commit); + assert.equal(response.status, result.status); + assert.equal(response.changed, false); + } +}); + +test("successful response with absent durable receipt is a protocol violation", async () => { + const response = await receipt().commit(store(applied, {status: "missing"}), commit); + assert.equal(response.status, "failed"); + assert.equal(response.reason_code, "coordination_commit_readback_mismatch"); +}); + +test("unresolved commit exposes same-operation recovery even when readback throws", async () => { + const target = store(applied, unavailable); + target.commitAuthority = async () => {throw new Error("synthetic post-commit response loss");}; + target.readReceipt = async () => {throw new Error("synthetic read failure");}; + const response = await receipt().commit(target, commit); + assert.equal(response.status, "ambiguous"); + assert.deepEqual(response.recovery, {operation_id: identity.operation_id, retry_with_same_operation_id: true}); + assert.equal(response.changed, false); + assert.equal(response.original_receipt, undefined); +}); + +test("a programmer error in a command decoder is not swallowed as corrupt storage", async () => { + const invalid = new CoordinationCommandReceipt({result_schema: "test_result_v0", identity, + decode() {throw new TypeError("decoder bug");}, + failure: (reason_code, reason) => ({schema_version: "test_result_v0", reason_code, reason})}); + await assert.rejects(invalid.read(store(applied, found)), /decoder bug/); +}); + +test("commit identity mismatch fails before the effect", async () => { + let writes = 0; + const target = store(applied, found); + target.commitAuthority = async () => {writes++; return applied;}; + await assert.rejects(receipt().commit(target, {...commit, operation_id: "other"}), /identities differ/); + assert.equal(writes, 0); +}); diff --git a/tests/control_plane_ts/coordination_receipt_conformance.ts b/tests/control_plane_ts/coordination_receipt_conformance.ts new file mode 100644 index 0000000000..bd87d98d7f --- /dev/null +++ b/tests/control_plane_ts/coordination_receipt_conformance.ts @@ -0,0 +1,126 @@ +import {executeCoordinationTodoArchiveCompleted} from "../../loopx/control_plane/coordination/todo_archive.ts"; +/** Command recovery over the complete production-scale head. Faults wrap the + * transport result only; all successful writes still reach the selected backend. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommitResult} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {TODO_DOMAIN_ITEM_SCHEMA} from "../../loopx/control_plane/coordination/coordination_state_contract.ts"; +import {executeCoordinationTodoCreate} from "../../loopx/control_plane/coordination/todo_create.ts"; +import {executeCoordinationTodoClaim} from "../../loopx/control_plane/coordination/todo_claim.ts"; +import {executeCoordinationTodoUpdate} from "../../loopx/control_plane/coordination/todo_update.ts"; +import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; +import {executeCoordinationTodoTerminalLifecycle} from "../../loopx/control_plane/coordination/todo_terminal_lifecycle.ts"; +import {productionScaleCoordinationFixture, PRODUCTION_SCALE_VALIDATION_DECLARATION} from "./production_scale_coordination_fixture.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; + +type Command = "create" | "claim" | "update" | "complete" | "supersede" | "archive" | "monitor"; +const commands: readonly Command[] = ["create", "claim", "update", "complete", "supersede", "archive", "monitor"]; +type Fault = "none" | "lost_response" | "unreadable_receipt" | "ambiguous_unreadable" | "thrown_response"; +const faults: readonly Fault[] = ["none", "lost_response", "unreadable_receipt", "ambiguous_unreadable", "thrown_response"]; + +async function commandFixture(store: AuthorityStore, command: Command) { + const goal_id = "goal-a"; + const fixture = productionScaleCoordinationFixture(goal_id); + const projection = fixture.projection; + const todos = projection.todos as JsonObject[]; + const leased = new Set((projection.leases as JsonObject[]).map(row => row.todo_id)); + const claimTodo = todos.find(row => row.role === "agent" && row.status === "open" && + row.task_class === "advancement_task" && !leased.has(row.todo_id))!; + const monitor = todos.find(row => row.role === "agent" && row.status === "open" && + row.task_class === "continuous_monitor" && !leased.has(row.todo_id))!; + assert.ok(claimTodo); assert.ok(monitor); + const operation_id = `recover-${command}`; + const common = {goal_id, operation_id, registered_agents: fixture.registered_agents, + dry_run: false, now: new Date("2026-09-07T07:00:00Z")}; + assert.equal((await store.commitAuthority({operation_id: "seed-recovery", expected_provider_revision: null, + events: [], receipts: [], next_projection: projection})).status, "applied"); + const invoke = (target: AuthorityStore, identity = operation_id): Promise => { + const request = {...common, operation_id: identity}; + if (command === "create") return executeCoordinationTodoCreate(target, {...request, + actor_agent_id: "agent-a", todo: {schema_version: TODO_DOMAIN_ITEM_SCHEMA, todo_id: "todo_recovery_created", + role: "agent", status: "open", done: false, archive_state: "active", text: "Recover the accepted create"}}); + if (command === "claim") return executeCoordinationTodoClaim(target, {...request, + actor_agent_id: String(claimTodo.claimed_by), claimed_by: String(claimTodo.claimed_by), + todo_id: String(claimTodo.todo_id), expected_role: "agent", + lease_request: {idempotency_key: "recovery-lease", expected_version: 0, ttl_seconds: 2700}}); + if (command === "update") return executeCoordinationTodoUpdate(target, {...request, + actor_agent_id: "agent-a", todo_id: fixture.completion_todo_id, expected_role: "agent", + patch: {text: "Recover the accepted copy edit"}, clear_fields: [], + lease_idempotency_key: fixture.completion_lease_idempotency_key, + lease_expected_version: fixture.completion_lease_expected_version}); + if (command === "archive") return executeCoordinationTodoArchiveCompleted(target, {...request, + role: "agent", max_active_done: 5}); + if (command === "monitor") return executeCoordinationMonitorPoll(target, {...request, + actor_agent_id: String(monitor.claimed_by), + observation: {todo_id: monitor.todo_id, result_hash: "new-recovery-evidence", material_change: true, + generated_at: "2026-09-07T07:00:00Z"}, + intent: {next_agent_todo: "Advance the recovered material change", next_action_kind: "implement"}}); + return executeCoordinationTodoTerminalLifecycle(target, {...request, command, + todo_id: fixture.completion_todo_id, actor_agent_id: "agent-a", expected_role: "agent", + lifecycle_grants: [], authority_reason: null, decision_outcome: null, + lease_idempotency_key: fixture.completion_lease_idempotency_key, + lease_expected_version: fixture.completion_lease_expected_version, + allow_user_gate_auto_acquire: false, requested_no_followup: command === "complete", + requested_completion_turn_key: null, requested_completion_identity_source: null, + linked_successor_todo_ids: [], successor_intents: [], note: null, evidence: "synthetic validation", + reason: "Replace the scoped work", clear_claim: false, completion_policy_request: null, + validation_declaration: command === "complete" ? PRODUCTION_SCALE_VALIDATION_DECLARATION : null, + validation_receipt: command === "complete" ? {schema_version: "issue_fix_validation_command_v0", command_label: "production-scale fixture validation", + exit_code: 0, passed: true, status: "passed", summary: "synthetic validation passed", + stdout_captured: false, stderr_captured: false, local_path_captured: false} : null}); + }; + return {invoke, operation_id, initial: await store.loadAuthority()}; +} + +export function registerCoordinationReceiptConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + for (const command of commands) for (const fault of faults) { + test(`${provider}: ${command} receipt recovery / ${fault}`, async t => { + const {store} = await factory(t); + const {invoke, operation_id, initial} = await commandFixture(store, command); + let attempted = false, commits = 0; + const wrapped: AuthorityStore = { + storeIdentity: () => store.storeIdentity(), loadAuthority: () => store.loadAuthority(), + scanCommitted: (cursor, limit) => store.scanCommitted(cursor, limit), + readReceipt: id => attempted && (fault === "unreadable_receipt" || fault === "ambiguous_unreadable") + ? Promise.resolve({status: "unavailable", reason_code: "synthetic_disconnect", reason: "receipt transport disconnected"}) + : store.readReceipt(id), + commitAuthority: async commit => { + commits++; attempted = true; + const result = await store.commitAuthority(commit); + assert.equal(result.status, "applied", JSON.stringify(result)); + if (fault === "thrown_response") throw new Error("synthetic response lost after commit"); + return (fault === "lost_response" || fault === "ambiguous_unreadable") + ? {status: "ambiguous", reason_code: "synthetic_timeout", reason: "commit response lost"} satisfies AuthorityStoreCommitResult + : result; + }, + }; + const result = await invoke(wrapped); + const uncertain = fault === "unreadable_receipt" || fault === "ambiguous_unreadable"; + assert.equal(result.status, uncertain ? "ambiguous" : fault === "none" ? "applied" : "recovered", JSON.stringify(result)); + if (uncertain) { + assert.equal((result.recovery as JsonObject).operation_id, operation_id); + assert.equal((result.recovery as JsonObject).retry_with_same_operation_id, true); + assert.equal(result.changed, false); + } + assert.equal(commits, 1, "a response fault never triggers a second commit"); + const after = await store.loadAuthority(); + assert.notDeepEqual(after, initial); + const receipt = await store.readReceipt(operation_id); + assert.equal(receipt.status, "found"); + assert.equal((await invoke(store)).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), after, "retry never changes head, leases or successors"); + // An unrelated later transaction must not redirect recovery to the latest head. + assert.equal(after.status, "loaded"); + if (after.status !== "loaded" || receipt.status !== "found") return; + await store.commitAuthority({operation_id: "later-write", expected_provider_revision: after.provider_revision, + next_projection: after.head, events: [], receipts: []}); + const replay = await invoke(store); + assert.equal(replay.cursor, receipt.cursor); + assert.equal(replay.provider_revision, receipt.provider_revision); + assert.equal(canonicalAuthoritySha256(replay.original_receipt ?? (replay.writeback as JsonObject)), + canonicalAuthoritySha256(result.original_receipt ?? replay.original_receipt ?? replay.writeback)); + }); + } +} diff --git a/tests/control_plane_ts/local_archive_attempt.test.ts b/tests/control_plane_ts/local_archive_attempt.test.ts index aa39dd6aa0..d7a9eab566 100644 --- a/tests/control_plane_ts/local_archive_attempt.test.ts +++ b/tests/control_plane_ts/local_archive_attempt.test.ts @@ -19,7 +19,7 @@ import { 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"; +import { executeCoordinationTodoArchiveCompleted } from "../../loopx/control_plane/coordination/todo_archive.ts"; function todo(id: string): JsonObject { return {schema_version: "todo_item_v0", todo_id: id, role: "agent",