From 301b40d9add46abe23a006464c4d1304b6240f9a Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 14:09:11 +0800 Subject: [PATCH 1/3] fix(coordination): recover observed commands and verify current claim execution Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/command_receipt.ts | 18 ++++- .../coordination/lease_acquisition_proof.ts | 65 +++++++++++++++++++ .../coordination/task_lease_acquire.ts | 51 +++------------ .../coordination/task_lease_lifecycle.ts | 4 +- .../coordination/todo_archive.ts | 9 ++- .../control_plane/coordination/todo_claim.ts | 48 +++++++++----- .../control_plane/coordination/todo_create.ts | 7 +- .../coordination/todo_monitor_poll.ts | 4 +- .../coordination/todo_terminal_lifecycle.ts | 14 ++-- .../control_plane/coordination/todo_update.ts | 4 +- .../goals/acceptance_authority.ts | 13 ++-- .../work_items/team_plan_authority.ts | 4 +- 12 files changed, 159 insertions(+), 82 deletions(-) create mode 100644 loopx/control_plane/coordination/lease_acquisition_proof.ts diff --git a/loopx/control_plane/coordination/command_receipt.ts b/loopx/control_plane/coordination/command_receipt.ts index fd2791728c..013b1806a7 100644 --- a/loopx/control_plane/coordination/command_receipt.ts +++ b/loopx/control_plane/coordination/command_receipt.ts @@ -3,7 +3,7 @@ * 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"; + AuthorityStoreReceiptResult, AuthorityStoreLoadResult} from "./authority_store.ts"; import {AuthorityStoreProtocolError, canonicalAuthorityObject} from "./authority_store_codec.ts"; import {projectionDelivery} from "../todos/projection_delivery.ts"; @@ -21,6 +21,12 @@ interface CommandReceiptContract { } type Result = JsonObject & {schema_version: S}; +/** A head and a receipt are separate reads. The receipt observed after the head + * wins: that head may already include this operation's effects. */ +export type CommandObservation = + | {kind: "receipt"; result: Result} + | {kind: "authority"; authority: AuthorityStoreLoadResult}; + /** 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 { @@ -56,6 +62,16 @@ export class CoordinationCommandReceipt { return this.project(await store.readReceipt(this.contract.identity.operation_id), "replayed"); } + /** Call after the initial historical lookup and command-specific source gates. + * Recheck before interpreting the head (including unavailable/stale heads). + * This is not a transaction or a write retry: a later competing commit remains + * subject to the provider CAS and normal post-commit receipt recovery. */ + async observe(store: AuthorityStore): Promise> { + const authority = await store.loadAuthority(); + const replay = await this.read(store); + return replay === null ? {kind: "authority", authority} : {kind: "receipt", result: replay}; + } + async commit(store: AuthorityStore, commit: AuthorityStoreCommit): Promise> { const {identity, result_schema, failure} = this.contract; if (commit.operation_id !== identity.operation_id) { diff --git a/loopx/control_plane/coordination/lease_acquisition_proof.ts b/loopx/control_plane/coordination/lease_acquisition_proof.ts new file mode 100644 index 0000000000..75e31790aa --- /dev/null +++ b/loopx/control_plane/coordination/lease_acquisition_proof.ts @@ -0,0 +1,65 @@ +/** Historical acquisition and current permission are different facts. Both + * standalone acquire and atomic claim/acquire return this current proof. */ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore, AuthorityStoreLoadResult} from "./authority_store.ts"; +import {canonicalAuthorityObject} from "./authority_store_codec.ts"; +import {indexCoordinationProjection, validateCoordinationTodoReadModel} from "./coordination_projection.ts"; +import {canonicalTaskLease, canonicalTaskLeaseAcquireFacts} from "./task_lease_state.ts"; +import {HANDOFF_MODES} from "./handoff_mode_policy.ts"; +import {requireStringLiteral} from "../runtime_decode.ts"; +import {leaseOwnerRejection} from "../work_items/task_lease_eligibility.ts"; +import {leaseEpoch, leaseVersion, leaseIsActive} from "../work_items/task_lease_acquire.ts"; +import {acceptanceWorkGuard} from "../goals/acceptance_contract.ts"; + +interface AcquisitionIdentity { + goal_id: string; todo_id: string; owner: string; idempotency_key: string; + registered_agents: readonly string[]; now: Date; +} + +/** Decode only a loaded decision snapshot. Receipt precedence must be settled + * before calling this: an already committed operation needs no new admission. */ +export function acquisitionFacts(head: AuthorityStoreLoadResult, input: AcquisitionIdentity) { + if (head.status !== "loaded") return {head, facts: null, mode: null}; + validateCoordinationTodoReadModel(head.head, input.goal_id); + const index = indexCoordinationProjection(head.head, input.goal_id); + return {head, facts: canonicalTaskLeaseAcquireFacts(index, input.goal_id, input.todo_id, input.registered_agents, input.now), + mode: requireStringLiteral(head.head.handoff_mode ?? "legacy", HANDOFF_MODES, "canonical handoff_mode")}; +} + +export async function currentLeaseAcquisitionProof(store: AuthorityStore, + input: AcquisitionIdentity & {operation_id: string; required_handoff_mode?: "hard_lease"}, + result: JsonObject & {schema_version: S}, + failed: (code: string, reason: string, detail?: JsonObject) => JsonObject & {schema_version: S}, +): Promise { + if (!["applied", "no_change", "replayed", "recovered"].includes(String(result.status))) return result; + const recovery = {operation_id: input.operation_id, retry_with_same_operation_id: true}; + const unavailable = () => ({...failed("canonical_acquire_readback_required", + "acquisition receipt is durable but current authority is unavailable; retry the same request"), + status: "ambiguous", original_receipt: result.original_receipt, recovery}); + let loaded: AuthorityStoreLoadResult; + try { loaded = await store.loadAuthority(); } + catch { return unavailable(); } + const {head, facts, mode} = acquisitionFacts(loaded, input); + if (head.status !== "loaded" || facts === null) return unavailable(); + const original = canonicalTaskLease(canonicalAuthorityObject(result.lease, "original acquire lease"), input.goal_id, input.todo_id); + const current = facts.current; + const details = {handoff_mode: mode, original_receipt: result.original_receipt, + current_provider_revision: head.provider_revision, current_cursor: head.cursor}; + const rejection = mode === "soft_claim" ? "handoff_mode_forbids_lease" + : input.required_handoff_mode !== undefined && mode !== input.required_handoff_mode ? "claim_lease_requires_hard_lease" + : leaseOwnerRejection(facts.todo, input.owner, input.registered_agents); + if (rejection) return failed(rejection, `current authority rejects lease acquire replay: ${rejection}`, details); + if (!current || !leaseIsActive(current, input.now) || current.owner !== input.owner || + current.idempotency_key !== input.idempotency_key || leaseEpoch(current) !== leaseEpoch(original) || + leaseVersion(current) < leaseVersion(original)) { + return failed("idempotency_key_reuse", "acquire receipt belongs to a retired execution; use a new execution key", details); + } + const acceptance = acceptanceWorkGuard(head.head, input.goal_id, input.todo_id); + if (acceptance !== null && !acceptance.allowed) { + return failed(String(acceptance.reason_code), `${String(acceptance.reason)} Inspect Goal acceptance and ask the owner to configure or rebind this Todo.`, + {...details, goal_acceptance_guard: acceptance}); + } + // A renewal advances version/expiry within this execution. Never rewrite the + // immutable acquisition receipt to make it look like the renewed decision. + return {...result, ...details, lease: current}; +} diff --git a/loopx/control_plane/coordination/task_lease_acquire.ts b/loopx/control_plane/coordination/task_lease_acquire.ts index 272b0be016..ec57ff6560 100644 --- a/loopx/control_plane/coordination/task_lease_acquire.ts +++ b/loopx/control_plane/coordination/task_lease_acquire.ts @@ -4,12 +4,10 @@ import type {JsonObject} from "../effect_program.ts"; import type {AuthorityStore} from "./authority_store.ts"; import {AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthoritySha256} from "./authority_store_codec.ts"; import {CoordinationCommandReceipt, commandReceiptResult} from "./command_receipt.ts"; -import {indexCoordinationProjection, validateCoordinationTodoReadModel, prepareCoordinationProjectionCommit} from "./coordination_projection.ts"; -import {canonicalTaskLease, canonicalTaskLeaseAcquireFacts} from "./task_lease_state.ts"; -import {HANDOFF_MODES} from "./handoff_mode_policy.ts"; -import {requireStringLiteral} from "../runtime_decode.ts"; +import {prepareCoordinationProjectionCommit} from "./coordination_projection.ts"; +import {canonicalTaskLease} from "./task_lease_state.ts"; import {decideTaskLeaseAcquire, materializeTaskLeaseAcquire} from "../work_items/task_lease_acquire_decision.ts"; -import {leaseOwnerRejection} from "../work_items/task_lease_eligibility.ts"; +import {acquisitionFacts, currentLeaseAcquisitionProof} from "./lease_acquisition_proof.ts"; import {acceptanceWorkGuard} from "../goals/acceptance_contract.ts"; import {normalizeGoalId, normalizeTodoId, normalizeOwner, normalizeIdempotencyKey, normalizeWriteScopes, normalizeTtl, leaseEpoch, leaseVersion, leaseIsActive, @@ -24,7 +22,7 @@ export interface CanonicalTaskLeaseAcquireInput { export async function executeCanonicalTaskLeaseAcquire(store: AuthorityStore, raw: CanonicalTaskLeaseAcquireInput, beforeCommit?: (lease: JsonObject) => Promise): Promise { const schema = "loopx_canonical_task_lease_acquire_result_v0"; - const failed = (code: string, reason: string, detail: JsonObject = {}) => + const failed = (code: string, reason: string, detail: JsonObject = {}): JsonObject & {schema_version: typeof schema} => ({schema_version: schema, status: "failed", changed: false, reason_code: code, reason, failure_stage: "validation", ...detail}); let input: CanonicalTaskLeaseAcquireInput & {ttl_seconds: number}; try { @@ -62,47 +60,16 @@ export async function executeCanonicalTaskLeaseAcquire(store: AuthorityStore, ra return {...payload, fields: {...payload.fields, operation_id: identity.operation_id}}; }}); - const readFacts = async () => { - const head = await store.loadAuthority(); - if (head.status !== "loaded") return {head, facts: null, mode: null}; - validateCoordinationTodoReadModel(head.head, input.goal_id); - const index = indexCoordinationProjection(head.head, input.goal_id); - return {head, facts: canonicalTaskLeaseAcquireFacts(index, input.goal_id, input.todo_id, input.registered_agents, input.now), - mode: requireStringLiteral(head.head.handoff_mode ?? "legacy", HANDOFF_MODES, "canonical handoff_mode")}; - }; - // Do not grant execution from a receipt that outlived its execution generation. - const currentProof = async (result: JsonObject): Promise => { - if (!["applied", "no_change", "replayed", "recovered"].includes(String(result.status))) return result; - const {head, facts, mode} = await readFacts(); - if (head.status !== "loaded" || facts === null) return {...failed("canonical_acquire_readback_required", - "acquisition receipt is durable but current authority is unavailable; retry the same request"), - status: "ambiguous", original_receipt: result.original_receipt, recovery: {operation_id: identity.operation_id, retry_with_same_operation_id: true}}; - const original = canonicalTaskLease(canonicalAuthorityObject(result.lease, "original acquire lease"), input.goal_id, input.todo_id); - const current = facts.current; - const details = {handoff_mode: mode, original_receipt: result.original_receipt, - current_provider_revision: head.provider_revision, current_cursor: head.cursor}; - const rejection = mode === "soft_claim" ? "handoff_mode_forbids_lease" : leaseOwnerRejection(facts.todo, input.owner, input.registered_agents); - if (rejection) return failed(rejection, `current authority rejects lease acquire replay: ${rejection}`, details); - if (!current || !leaseIsActive(current, input.now) || current.owner !== input.owner || - current.idempotency_key !== input.idempotency_key || leaseEpoch(current) !== leaseEpoch(original) || - leaseVersion(current) < leaseVersion(original)) { - return failed("idempotency_key_reuse", "acquire receipt belongs to a retired execution; use a new execution key", details); - } - const acceptance = acceptanceWorkGuard(head.head, input.goal_id, input.todo_id); - if (acceptance !== null && !acceptance.allowed) { - return failed(String(acceptance.reason_code), `${String(acceptance.reason)} Inspect Goal acceptance and ask the owner to configure or rebind this Todo.`, - {...details, goal_acceptance_guard: acceptance}); - } - // Renewal may advance version/expiry within this execution. Return current - // usable proof while the immutable receipt preserves the original decision. - return {...result, ...details, lease: current}; - }; + const currentProof = (result: JsonObject & {schema_version: typeof schema}) => + currentLeaseAcquisitionProof(store, {...input, operation_id: identity.operation_id}, result, failed); let committed = false; try { const replay = await receipt.read(store); if (replay !== null) return await currentProof(replay); - const {head, facts, mode} = await readFacts(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return await currentProof(observation.result); + const {head, facts, mode} = acquisitionFacts(observation.authority, input); if (head.status !== "loaded" || facts === null) return failed("canonical_lease_authority_unavailable", "canonical lease authority is unavailable; restore the selected provider before retrying", {...head}); const decision = decideTaskLeaseAcquire({handoff_mode: mode!, registered_agents: input.registered_agents, diff --git a/loopx/control_plane/coordination/task_lease_lifecycle.ts b/loopx/control_plane/coordination/task_lease_lifecycle.ts index 3b8e64f209..95c02f61b5 100644 --- a/loopx/control_plane/coordination/task_lease_lifecycle.ts +++ b/loopx/control_plane/coordination/task_lease_lifecycle.ts @@ -109,7 +109,9 @@ export async function executeCanonicalTaskLeaseLifecycle(store: AuthorityStore, try { const replay = await receipt.read(store); if (replay !== null) return replay; - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") return {schema_version: contract.result, ...head, failure_stage: "validation", changed: false}; const index = indexCoordinationProjection(head.head, input.goal_id); validateCoordinationTodoReadModel(head.head, input.goal_id); diff --git a/loopx/control_plane/coordination/todo_archive.ts b/loopx/control_plane/coordination/todo_archive.ts index 1097be272d..fe3f086896 100644 --- a/loopx/control_plane/coordination/todo_archive.ts +++ b/loopx/control_plane/coordination/todo_archive.ts @@ -97,11 +97,16 @@ export async function executeCoordinationTodoArchiveCompleted( }); // Preview observes the current snapshot without consuming or replaying a // durable operation identity. Historical receipts precede current-head CAS. + const receipt = archiveReceipt(input, requestSha); if (!input.dry_run) { - const replay = await archiveReceipt(input, requestSha).read(store); + const replay = await receipt.read(store); if (replay !== null) return replay; } - const head = await store.loadAuthority(); + const observation = input.dry_run + ? {kind: "authority" as const, authority: await store.loadAuthority()} + : await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") { return {schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...head, changed: false}; } diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index 9b27da6348..a906a1c2b8 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -1,4 +1,5 @@ import {AUTHORITY_SOURCE_CHANGED, uncheckedAuthoritySource, type AuthoritySourceCheck} from "./authority_source.ts"; +import {currentLeaseAcquisitionProof} from "./lease_acquisition_proof.ts"; import {canonicalTaskLeaseAcquireFacts} from "./task_lease_state.ts"; import {evaluateTaskLeaseAcquireDecision, materializeTaskLeaseAcquire} from "../work_items/task_lease_acquire_decision.ts"; import type { JsonObject } from "../effect_program.ts"; @@ -395,26 +396,39 @@ export async function executeCoordinationTodoClaim( // 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 (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason, - existing === null ? {} : {original_receipt: existing.original_receipt}, "decision_rejection"); - if (existing !== null) { - // A historical claim receipt cannot grant work after its acceptance binding - // changed. Preserve the original response when the contract is absent. - const current = await store.loadAuthority(); - if (current.status === "loaded") { - const guard = acceptanceWorkGuard(current.head, input.goal_id, input.todo_id); - if (guard !== null && !guard.allowed) { - return failure(String(guard.reason_code), `${String(guard.reason)} Inspect Goal acceptance and ask the owner to configure or rebind this Todo.`, + const currentProof = async (result: CoordinationTodoClaimResult): Promise => { + if (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason, + {original_receipt: result.original_receipt}, "decision_rejection"); + if (leaseRequest !== null) { + try { + result = await currentLeaseAcquisitionProof(store, {...input, owner: input.claimed_by, + idempotency_key: leaseRequest.idempotency_key, required_handoff_mode: "hard_lease"}, result, + (code, reason, detail = {}) => failure(code, reason, detail, "decision_rejection")); + } catch (error) { + return failure("invalid_coordination_task_lease", error instanceof Error ? error.message : "invalid canonical task lease", + {original_receipt: result.original_receipt}); + } + } else if (["replayed", "recovered"].includes(String(result.status))) { + // A plain historical claim has no lease proof to grant, but an enabled + // acceptance contract still governs adoption of that work. + const current = await store.loadAuthority(); + if (current.status === "loaded") { + const guard = acceptanceWorkGuard(current.head, input.goal_id, input.todo_id); + if (guard !== null && !guard.allowed) return failure(String(guard.reason_code), + `${String(guard.reason)} Inspect Goal acceptance and ask the owner to configure or rebind this Todo.`, {goal_acceptance_guard: guard}, "decision_rejection"); } } if (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason, - {original_receipt: existing.original_receipt}, "decision_rejection"); - return existing; - } - - const head = await store.loadAuthority(); + {original_receipt: result.original_receipt}, "decision_rejection"); + return result; + }; + const existing = await receipt.read(store); + if (existing !== null) return currentProof(existing); + if (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason, {}, "decision_rejection"); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return currentProof(observation.result); + const head = observation.authority; if (head.status !== "loaded") { return { schema_version: COORDINATION_TODO_CLAIM_RESULT_SCHEMA, @@ -635,5 +649,5 @@ export async function executeCoordinationTodoClaim( request_sha256: requestSha, result, }]; - return receipt.commit(store, commit); + return currentProof(await receipt.commit(store, commit)); } diff --git a/loopx/control_plane/coordination/todo_create.ts b/loopx/control_plane/coordination/todo_create.ts index 0dc343b4fc..9a9c3e045a 100644 --- a/loopx/control_plane/coordination/todo_create.ts +++ b/loopx/control_plane/coordination/todo_create.ts @@ -216,11 +216,14 @@ export async function executeCoordinationTodoCreate( actor_agent_id: input.actor_agent_id, dry_run: input.dry_run, }); - const existing = await createReceipt(input, requestSha).read(store); + const receipt = createReceipt(input, requestSha); + const existing = await receipt.read(store); if (existing !== null) return existing; if (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason); - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") { return {schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, ...head}; } diff --git a/loopx/control_plane/coordination/todo_monitor_poll.ts b/loopx/control_plane/coordination/todo_monitor_poll.ts index 6d483a9da5..fea14c2645 100644 --- a/loopx/control_plane/coordination/todo_monitor_poll.ts +++ b/loopx/control_plane/coordination/todo_monitor_poll.ts @@ -177,7 +177,9 @@ export async function executeCoordinationMonitorPoll(store: AuthorityStore, const previous = await receipt.read(store); if (previous) return previous; if (!await authoritySourcesCurrent()) return failure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason); - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...head}; let plan: ReturnType; try { plan = planWriteback(input, head.head); } diff --git a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts index 7ee2795db5..cd522b90ba 100644 --- a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts +++ b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts @@ -790,12 +790,15 @@ export async function executeCoordinationTodoTerminalLifecycle( const requestSha = terminalRequestSha(input); // Terminal recovery reports history, not present execution authority. It // precedes review/validation freshness even if a Monitor has since reopened. - const replay = await terminalReceipt(input, requestSha).read(store); + const receipt = terminalReceipt(input, requestSha); + const replay = await receipt.read(store); if (replay !== null) return replay; if (!await authoritySourcesCurrent()) return terminalFailure(AUTHORITY_SOURCE_CHANGED.code, AUTHORITY_SOURCE_CHANGED.reason, {}, "decision_rejection"); - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") { return { schema_version: COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA, @@ -830,13 +833,6 @@ export async function executeCoordinationTodoTerminalLifecycle( todo_id: input.todo_id, }, "decision_rejection"); } - // The first receipt read can precede a peer commit while this head already - // observes it. Recover only the matching operation receipt; never discard a - // validation receipt to manufacture a terminal replay from Todo state alone. - if (input.command === "complete" && todo.status === "done" && input.validation_receipt !== null) { - const committedReplay = await terminalReceipt(input, requestSha).read(store); - if (committedReplay !== null) return committedReplay; - } if (input.expected_role !== null && todo.role !== input.expected_role) { return terminalFailure( "todo_role_mismatch", diff --git a/loopx/control_plane/coordination/todo_update.ts b/loopx/control_plane/coordination/todo_update.ts index 3274431071..a61f14fae9 100644 --- a/loopx/control_plane/coordination/todo_update.ts +++ b/loopx/control_plane/coordination/todo_update.ts @@ -166,7 +166,9 @@ export async function executeCoordinationTodoUpdate( !input.registered_agents.includes(input.actor_agent_id)) { return failure("actor_not_registered", "Todo update requires a registered actor"); } - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") { return {schema_version: COORDINATION_TODO_UPDATE_RESULT_SCHEMA, ...head, changed: false}; } diff --git a/loopx/control_plane/goals/acceptance_authority.ts b/loopx/control_plane/goals/acceptance_authority.ts index f66a555591..b3869fd51c 100644 --- a/loopx/control_plane/goals/acceptance_authority.ts +++ b/loopx/control_plane/goals/acceptance_authority.ts @@ -2,7 +2,7 @@ * Only trusted local owner/host adapters may invoke these mutation exports. */ import {isAbsolute} from "node:path"; import type {JsonObject} from "../effect_program.ts"; -import type {AuthorityStore, AuthorityStoreHead} from "../coordination/authority_store.ts"; +import type {AuthorityStore, AuthorityStoreHead, AuthorityStoreLoadResult} from "../coordination/authority_store.ts"; import {authorityStoreSourceAuthority} from "../coordination/authority_store.ts"; import {AuthorityStoreProtocolError, canonicalAuthorityObject, canonicalAuthoritySha256, requireAuthorityStoreId} from "../coordination/authority_store_codec.ts"; @@ -52,8 +52,7 @@ function receiptFor(request: JsonObject, kind: "configure" | "verify") { decode: original => ({fields: canonicalAuthorityObject(original.result, "acceptance operation result"), changed: true})}); return {identity, receipt}; } -async function current(store: AuthorityStore, request: JsonObject): Promise { - const loaded = await store.loadAuthority(); +function current(loaded: AuthorityStoreLoadResult, request: JsonObject): AuthorityStoreHead | JsonObject { if (loaded.status !== "loaded") return {...loaded}; if (loaded.provider_revision !== request.expected_provider_revision) return { status: "conflict", changed: false, reason_code: "goal_acceptance_provider_revision_mismatch", @@ -92,7 +91,9 @@ export async function configureGoalAcceptance(store: AuthorityStore, value: Json const command = receiptFor(request, "configure"); const replay = await command.receipt.read(store); if (replay) return source(store, {...replay, projection_delivery: "not_required"}); - const head = await current(store, request); + const observation = await command.receipt.observe(store); + if (observation.kind === "receipt") return source(store, {...observation.result, projection_delivery: "not_required"}); + const head = current(observation.authority, request); if (!loaded(head)) return source(store, head); const previous = readGoalAcceptance(head.head, String(request.goal_id)); let state: AcceptanceState | null = previous; @@ -124,7 +125,9 @@ export async function commitGoalAcceptanceVerification(store: AuthorityStore, va const command = receiptFor(request, "verify"); const replay = await command.receipt.read(store); if (replay) return source(store, {...replay, projection_delivery: "not_required"}); - const head = await current(store, request); + const observation = await command.receipt.observe(store); + if (observation.kind === "receipt") return source(store, {...observation.result, projection_delivery: "not_required"}); + const head = current(observation.authority, request); if (!loaded(head)) return source(store, head); const goalId = String(request.goal_id); const state = readGoalAcceptance(head.head, goalId); diff --git a/loopx/control_plane/work_items/team_plan_authority.ts b/loopx/control_plane/work_items/team_plan_authority.ts index 93cde7a74b..a3cdcdea5e 100644 --- a/loopx/control_plane/work_items/team_plan_authority.ts +++ b/loopx/control_plane/work_items/team_plan_authority.ts @@ -23,7 +23,9 @@ export async function commitTeamPlan(store: AuthorityStore, request: JsonObject) }}); const previous = await receipt.read(store); if (previous) return previous; - const head = await store.loadAuthority(); + const observation = await receipt.observe(store); + if (observation.kind === "receipt") return observation.result; + const head = observation.authority; if (head.status !== "loaded") return {...head}; // The canonical revision is included at preview. Stale work cannot be // admitted merely because its Markdown projection has not caught up yet. From 2de18b06a9ef03179bb5683e8f434d0d5aa29d76 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 14:09:11 +0800 Subject: [PATCH 2/3] test(coordination): qualify replay races and acquisition proof across providers Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../test_canonical_claim_execution_proof.py | 79 +++++ .../authority_store_conformance.ts | 4 + .../claim_acquisition_proof_conformance.ts | 289 ++++++++++++++++++ .../command_observation_conformance.ts | 125 ++++++++ .../command_observation_scenarios.ts | 118 +++++++ .../control_plane_ts/command_receipt.test.ts | 61 ++++ .../local_authority_runtime.test.ts | 9 +- 7 files changed, 683 insertions(+), 2 deletions(-) create mode 100644 tests/control_plane/test_canonical_claim_execution_proof.py create mode 100644 tests/control_plane_ts/claim_acquisition_proof_conformance.ts create mode 100644 tests/control_plane_ts/command_observation_conformance.ts create mode 100644 tests/control_plane_ts/command_observation_scenarios.ts diff --git a/tests/control_plane/test_canonical_claim_execution_proof.py b/tests/control_plane/test_canonical_claim_execution_proof.py new file mode 100644 index 0000000000..83151b7285 --- /dev/null +++ b/tests/control_plane/test_canonical_claim_execution_proof.py @@ -0,0 +1,79 @@ +"""The public claim command must never grant execution from historical proof.""" +import json +import subprocess +import sys +from pathlib import Path + +import pytest +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection + +REPO = Path(__file__).resolve().parents[2] + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("retirement", ["release", "transfer"]) +def test_claim_replay_uses_current_execution(tmp_path, monkeypatch, provider, retirement): + isolate_sqlite_runtime(tmp_path, monkeypatch) + runtime, state, registry = tmp_path / "runtime", tmp_path / "state.md", tmp_path / "registry.json" + goal, target = "claim-execution", "todo_claim_execution" + state.write_text("# Claim execution\n\n## Agent Todo\n") + registry.write_text(json.dumps({"common_runtime_root": str(runtime), "goals": [{ + "id": goal, "repo": str(tmp_path), "state_file": state.name, + "coordination": {"registered_agents": ["agent-a", "agent-b"]}, + }]})) + projection = build_todo_runtime_shadow_projection(goal_id=goal, handoff_mode="hard_lease", leases=[], todos=[{ + "schema_version": "todo_item_v0", "todo_id": target, "role": "agent", "status": "open", "done": False, + "text": "Adopt work and obtain current execution", "archive_state": "active", "source_section": "Agent Todo", + "index": 1, "task_class": "advancement_task", "required_write_scopes": ["src/**"], + }]) + initialize_canonical_authority(runtime, goal, projection, state_path=state, provider=provider) + state.unlink() # Recovery must not rebuild authority from missing Markdown. + + def cli(command, *args, expected_exit=0): + process = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), "--format", "json", + *command, "--goal-id", goal, "--todo-id", target, *args], cwd=REPO, + capture_output=True, text=True, timeout=90, check=False) + assert process.returncode == expected_exit, process.stdout + process.stderr + return json.loads(process.stdout) + + request = ["--claimed-by", "agent-a", "--agent-id", "agent-a", "--claim-operation-id", "claim-original", + "--task-lease-idempotency-key", "execution-a", "--task-lease-expected-version", "0"] + lease_args = ["--owner", "agent-a", "--idempotency-key", "execution-a"] + try: + preview = cli(["todo", "claim"], *request, "--dry-run") + assert preview["status"] == "planned" + assert not cli(["task-lease", "inspect"])["active"] + first = cli(["todo", "claim"], *request) + assert first["status"] == "applied" + assert first["lease"]["version"] == first["lease"]["lease_epoch"] == 1 + assert first["lease"]["write_scopes"] == ["src/**"] + renewed = cli(["task-lease", "renew"], *lease_args, "--expected-version", "1", "--ttl-seconds", "600") + replay = cli(["todo", "claim"], *request) + assert replay["status"] == "replayed" + assert replay["changed"] is False + assert replay["lease"] == renewed["lease"] + assert replay["lease"]["version"] == 2 + assert replay["original_receipt"] == first["original_receipt"] + assert replay["provider_revision"] == first["provider_revision"] + assert replay["current_provider_revision"] != first["provider_revision"] + transition = ["--expected-version", "2"] + if retirement == "transfer": + transition += ["--new-owner", "agent-b", "--new-idempotency-key", "execution-b", "--transfer-claim"] + retired = cli(["task-lease", retirement], *lease_args, *transition) + assert retired["ok"] + state.unlink(missing_ok=True) # Maintenance may refresh display; remove it before recovery. + before = cli(["task-lease", "inspect"]) + rejected = cli(["todo", "claim"], *request, expected_exit=1) + assert rejected["ok"] is False + assert rejected["error_code"] == ("idempotency_key_reuse" if retirement == "release" else "owner_conflicts_with_claim") + assert not rejected.get("lease"), "a rejected old execution cannot return an active lease" + assert cli(["task-lease", "inspect"]) == before + assert not state.exists(), "rejected recovery cannot revive a legacy writer" + readback = cli(["todo", "list"]) + assert first["source_authority"] == provider + "_v0" + assert readback["todos"][0]["claimed_by"] == ("agent-b" if retirement == "transfer" else "agent-a") + assert not (runtime / "goals" / goal / "task-leases" / f"{target}.json").exists() + finally: + subprocess.run([sys.executable, "-c", "from loopx.control_plane.effect_runtime import effect_runtime_result; effect_runtime_result('runtime.shutdown',{},retry_safe=False)"], + cwd=REPO, capture_output=True, text=True, timeout=30, check=True) diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 24b946abde..7506abe3ab 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,3 +1,5 @@ +import {registerClaimAcquisitionProofConformance} from "./claim_acquisition_proof_conformance.ts"; +import {registerCommandObservationConformance} from "./command_observation_conformance.ts"; import {registerPeriodicReportConformance} from "./periodic_report_conformance.ts"; import {registerIssueFixMonitorReconciliationConformance} from "./issue_fix_monitor_reconciliation_conformance.ts"; import {registerTodoConsumerScopeConformance} from "./todo_consumer_scope_conformance.ts"; @@ -251,6 +253,8 @@ export function registerAuthorityStoreConformance( registerLeaseLifecycleConformance(providerName, factory); registerClaimTransferConformance(providerName, factory); registerLeaseAcquisitionConformance(providerName, factory); + registerCommandObservationConformance(providerName, factory); + registerClaimAcquisitionProofConformance(providerName, factory); registerAuthorityScanConformance(providerName, factory); registerOwnershipObservationConformance(providerName, factory); registerSuccessionReadConformance(providerName, factory); diff --git a/tests/control_plane_ts/claim_acquisition_proof_conformance.ts b/tests/control_plane_ts/claim_acquisition_proof_conformance.ts new file mode 100644 index 0000000000..e0722e7d5f --- /dev/null +++ b/tests/control_plane_ts/claim_acquisition_proof_conformance.ts @@ -0,0 +1,289 @@ +/** Atomic adoption must meet the same present-tense execution proof as acquire. + * Receipt history is immutable even when current execution changes. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {executeCoordinationTodoClaim as claim} from "../../loopx/control_plane/coordination/todo_claim.ts"; +import {executeCanonicalTaskLeaseLifecycle as maintain} from "../../loopx/control_plane/coordination/task_lease_lifecycle.ts"; +import {prepareCoordinationProjectionCommit} from "../../loopx/control_plane/coordination/coordination_projection.ts"; +import {configureGoalAcceptance} from "../../loopx/control_plane/goals/acceptance_authority.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import {productionScaleLeaseAcquisitionFixture} from "./production_scale_coordination_fixture.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; + +async function head(store: AuthorityStore) { + const loaded = await store.loadAuthority(); + assert.equal(loaded.status, "loaded"); + if (loaded.status !== "loaded") throw new Error("fixture authority unavailable"); + return loaded; +} + +export function registerClaimAcquisitionProofConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + async function setup(t: test.TestContext, schema: "native" | "legacy") { + const {store, contender} = await factory(t); + const f = productionScaleLeaseAcquisitionFixture("goal-a", schema); + const projection = authorityProjectionFixture("goal-a", (f.projection.todos as JsonObject[]).map(todo => + todo.todo_id === f.target ? {...todo, required_write_scopes: ["lease-admission/**"]} : todo), + f.projection.leases as JsonObject[], schema, {handoff_mode: "hard_lease"}); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: projection, events: [], receipts: []})).status, "applied"); + const request = {goal_id: "goal-a", todo_id: f.target, claimed_by: "agent-a", actor_agent_id: "agent-a", + expected_role: "agent", registered_agents: f.registered_agents, operation_id: "claim-acquire", + lease_request: {idempotency_key: "execution-a", expected_version: 0, ttl_seconds: 600}, + dry_run: false, now: new Date(f.scenario.now)}; + const maintenance = {goal_id: request.goal_id, todo_id: request.todo_id, owner: "agent-a", + idempotency_key: "execution-a", expected_version: 1, ttl_seconds: 600, + registered_agents: request.registered_agents, now: new Date("2026-09-13T10:06:00Z")}; + return {store, contender, request, maintenance, projection}; + } + + for (const schema of ["native", "legacy"] as const) { + test(`${provider} claim proof: ${schema} renewal returns current proof and retains original receipt`, async t => { + const {store, contender, request, maintenance} = await setup(t, schema); + const first = await claim(store, request); + assert.equal(first.status, "applied", JSON.stringify(first)); + const original = await store.readReceipt(request.operation_id); + const renewed = await maintain(contender, {...maintenance, operation: "renew"}); + assert.equal(renewed.status, "applied", JSON.stringify(renewed)); + const before = await head(store); + const replay = await claim(store, {...request, now: maintenance.now}); + assert.equal(replay.status, "replayed", JSON.stringify(replay)); + assert.equal(replay.changed, false); + assert.equal(replay.provider_revision, first.provider_revision, "receipt revision remains historical"); + assert.equal(replay.current_provider_revision, before.provider_revision); + assert.equal((replay.lease as JsonObject).version, 2); + assert.equal((replay.lease as JsonObject).lease_epoch, 1); + assert.deepEqual(replay.lease, renewed.lease); + assert.deepEqual(replay.original_receipt, first.original_receipt); + assert.deepEqual(await store.readReceipt(request.operation_id), original); + assert.deepEqual(await head(store), before); + }); + + for (const retirement of ["release", "transfer", "expiry"] as const) { + test(`${provider} claim proof: ${schema}/${retirement} cannot replay current execution`, async t => { + const {store, contender, request, maintenance} = await setup(t, schema); + assert.equal((await claim(store, request)).status, "applied"); + const original = await store.readReceipt(request.operation_id); + if (retirement !== "expiry") { + const result = await maintain(contender, {...maintenance, operation: retirement, + ttl_seconds: retirement === "release" ? null : 600, + ...(retirement === "transfer" ? {new_owner: "agent-b", new_idempotency_key: "receiver-b", transfer_claim: true} : {})}); + assert.equal(result.status, "applied", JSON.stringify(result)); + } + const before = await head(store); + const result = await claim(store, {...request, + now: new Date(retirement === "expiry" ? "2026-09-13T10:15:00Z" : "2026-09-13T10:07:00Z")}); + assert.equal(result.status, "failed", JSON.stringify(result)); + assert.equal(result.reason_code, retirement === "transfer" ? "owner_conflicts_with_claim" : "idempotency_key_reuse"); + assert.equal(result.lease, undefined, "failure must not carry usable old proof"); + assert.deepEqual(await store.readReceipt(request.operation_id), original); + assert.deepEqual(await head(store), before); + }); + } + + for (const dimension of ["excluded", "unregistered", "archived", "done", "claim", "soft-mode", "legacy-mode"] as const) { + test(`${provider} claim proof: ${schema}/${dimension} current eligibility supersedes history`, async t => { + const {store, request} = await setup(t, schema); + assert.equal((await claim(store, request)).status, "applied"); + const before = await head(store); + const todo = (before.head.todos as JsonObject[]).find(row => row.todo_id === request.todo_id)!; + const patch: JsonObject = dimension === "excluded" ? {excluded_agents: ["agent-a"]} + : dimension === "archived" ? {archive_state: "archive"} + : dimension === "done" ? {status: "done", done: true} + : dimension === "claim" ? {claimed_by: "agent-b"} : {}; + if (Object.keys(patch).length) { + assert.equal((await store.commitAuthority(prepareCoordinationProjectionCommit({goal_id: request.goal_id, + operation_id: "eligibility-change", expected_provider_revision: before.provider_revision, + projection: before.head, mutations: [{kind: "todo_upsert", todo: {...todo, ...patch}}]}))).status, "applied"); + } else if (dimension.endsWith("mode")) { + assert.equal((await store.commitAuthority({operation_id: "mode-change", expected_provider_revision: before.provider_revision, + next_projection: {...before.head, handoff_mode: dimension === "soft-mode" ? "soft_claim" : "legacy"}, events: [], receipts: []})).status, "applied"); + } + const final = await head(store); + const replay = await claim(store, {...request, registered_agents: dimension === "unregistered" ? ["agent-b"] : request.registered_agents}); + assert.equal(replay.status, "failed", JSON.stringify(replay)); + assert.equal(replay.lease, undefined); + assert.deepEqual(await head(store), final); + }); + } + + test(`${provider} claim proof: ${schema} missing current head is ambiguous and same-operation recoverable`, async t => { + const {store, request} = await setup(t, schema); + const first = await claim(store, request); + assert.equal(first.status, "applied"); + const before = await head(store); + for (const throws of [false, true]) { + const unavailable = new Proxy(store, {get(target, key) { + if (key === "loadAuthority") return async () => { + if (throws) throw new Error("synthetic lost read response"); + return {status: "unavailable", reason_code: "offline", reason: "synthetic offline provider"}; + }; + const value = Reflect.get(target, key, target); + return typeof value === "function" ? value.bind(target) : value; + }}); + const result = await claim(unavailable, request); + assert.equal(result.status, "ambiguous", JSON.stringify(result)); + assert.equal(result.reason_code, "canonical_acquire_readback_required"); + assert.equal(result.lease, undefined); + assert.deepEqual(result.recovery, {operation_id: request.operation_id, retry_with_same_operation_id: true}); + } + const recovered = await claim(store, request); + assert.equal(recovered.status, "replayed"); + assert.deepEqual(recovered.lease, first.lease); + assert.deepEqual(await head(store), before); + }); + + test(`${provider} claim proof: ${schema} lost acknowledgment still requires current execution`, async t => { + const {store, contender, request, maintenance} = await setup(t, schema); + let writes = 0; + const intercepted = new Proxy(store, {get(target, key) { + if (key === "commitAuthority") return async (commit: AuthorityStoreCommit) => { + writes++; + assert.equal((await target.commitAuthority(commit)).status, "applied"); + assert.equal((await maintain(contender, {...maintenance, operation: "release", ttl_seconds: null})).status, "applied"); + throw new Error("response lost after commit and lease retirement"); + }; + const value = Reflect.get(target, key, target); + return typeof value === "function" ? value.bind(target) : value; + }}); + const result = await claim(intercepted, request); + assert.equal(writes, 1); + assert.equal(result.status, "failed", JSON.stringify(result)); + assert.equal(result.reason_code, "idempotency_key_reuse"); + const current = await head(store); + assert.equal(current.cursor, "3", "claim and independent release each commit once"); + assert.equal((await store.readReceipt(request.operation_id)).status, "found"); + assert.equal((current.head.leases as JsonObject[]).find(row => row.todo_id === request.todo_id)!.status, "released"); + }); + + test(`${provider} claim proof: ${schema} a no-op receipt also needs current proof`, async t => { + const {store, contender, request, maintenance} = await setup(t, schema); + assert.equal((await claim(store, request)).status, "applied"); + const noOp = {...request, operation_id: "adopt-existing-execution", + lease_request: {...request.lease_request, expected_version: 1}}; + const first = await claim(contender, noOp); + assert.equal(first.status, "no_change", JSON.stringify(first)); + assert.equal(first.changed, false); + assert.equal(first.projection_delivery, "not_required"); + const before = await head(store); + const replay = await claim(store, noOp); + assert.equal(replay.status, "replayed"); + assert.equal(replay.changed, false); + assert.deepEqual(replay.original_receipt, first.original_receipt); + assert.deepEqual(await head(store), before); + assert.equal((await maintain(contender, {...maintenance, operation: "release", ttl_seconds: null})).status, "applied"); + const released = await head(store); + assert.equal((await claim(store, noOp)).reason_code, "idempotency_key_reuse"); + assert.deepEqual(await head(store), released); + }); + + test(`${provider} claim proof: ${schema} source witness is checked after current proof`, async t => { + const {store, request} = await setup(t, schema); + const first = await claim(store, request); + assert.equal(first.status, "applied"); + const before = await head(store); + let checks = 0; + const result = await claim(store, request, async () => ++checks === 1); + assert.equal(checks, 2); + assert.equal(result.status, "failed", JSON.stringify(result)); + assert.equal(result.reason_code, "authority_source_changed"); + assert.equal(result.lease, undefined); + assert.deepEqual(await head(store), before); + }); + + + test(`${provider} claim proof: ${schema} changed request cannot reuse a successful claim identity`, async t => { + const {store, request} = await setup(t, schema); + assert.equal((await claim(store, request)).status, "applied"); + const before = await head(store); + const original = await store.readReceipt(request.operation_id); + for (const change of [ + {lease_request: {...request.lease_request, ttl_seconds: 900}}, + {lease_request: {...request.lease_request, expected_version: 1}}, + {lease_request: {...request.lease_request, idempotency_key: "different-execution"}}, + {claimed_by: "agent-b", actor_agent_id: "agent-b"}, + ]) { + const replay = await claim(store, {...request, ...change}); + assert.equal(replay.status, "failed"); + assert.equal(replay.reason_code, "coordination_operation_identity_mismatch"); + assert.equal(replay.lease, undefined); + assert.deepEqual(await head(store), before); + assert.deepEqual(await store.readReceipt(request.operation_id), original); + } + }); + + test(`${provider} claim proof: ${schema} acceptance rebinding is current authority`, async t => { + const {store, request} = await setup(t, schema); + const first = await claim(store, request); + assert.equal(first.status, "applied"); + const before = await head(store); + const document = {objective: "Validate the accepted task", non_goals: [], + criteria: [{id: "outcome", description: "Focused check passes", validation_argv: [process.execPath, "-e", "process.exit(0)"]}], + bindings: [{todo_id: request.todo_id, criterion_ids: ["outcome"]}]}; + assert.equal((await configureGoalAcceptance(store, {goal_id: request.goal_id, actor_agent_id: null, + operation_id: "bind-acceptance", expected_provider_revision: before.provider_revision, document})).status, "applied"); + assert.equal((await claim(store, request)).status, "replayed", "valid owner binding permits current execution"); + const bound = await head(store); + const todo = (bound.head.todos as JsonObject[]).find(row => row.todo_id === request.todo_id)!; + assert.equal((await store.commitAuthority(prepareCoordinationProjectionCommit({goal_id: request.goal_id, + operation_id: "change-bound-work", expected_provider_revision: bound.provider_revision, projection: bound.head, + mutations: [{kind: "todo_upsert", todo: {...todo, text: "Materially different work needs owner rebind"}}]}))).status, "applied"); + const changed = await head(store); + const result = await claim(store, request); + assert.equal(result.status, "failed"); + assert.equal((result.goal_acceptance_guard as JsonObject).allowed, false); + assert.equal(result.lease, undefined); + assert.deepEqual(result.original_receipt, first.original_receipt); + assert.deepEqual(await head(store), changed); + }); + + test(`${provider} claim proof: ${schema} new execution can replace a released generation`, async t => { + const {store, request, maintenance} = await setup(t, schema); + const first = await claim(store, request); + assert.equal(first.status, "applied"); + assert.equal((await maintain(store, {...maintenance, operation: "release", ttl_seconds: null})).status, "applied"); + const next = {...request, operation_id: "next-claim", now: maintenance.now, + lease_request: {idempotency_key: "execution-next", expected_version: 1, ttl_seconds: 600}}; + const result = await claim(store, next); + assert.equal(result.status, "applied", JSON.stringify(result)); + assert.equal(result.todo_changed, false, "the existing owner need not be claimed twice"); + assert.equal(result.lease_changed, true); + assert.equal((result.lease as JsonObject).lease_epoch, 2); + assert.equal((result.lease as JsonObject).version, 2); + const before = await head(store); + assert.equal((await claim(store, request)).reason_code, "idempotency_key_reuse"); + const replay = await claim(store, next); + assert.equal(replay.status, "replayed"); + assert.deepEqual(replay.lease, result.lease); + assert.deepEqual(await head(store), before); + }); + + + test(`${provider} claim proof: ${schema} lost acknowledgment recovers one still-current acquisition`, async t => { + const {store, request} = await setup(t, schema); + let writes = 0; + const interrupted = new Proxy(store, {get(target, key) { + if (key === "commitAuthority") return async (commit: AuthorityStoreCommit) => { + writes++; + assert.equal((await target.commitAuthority(commit)).status, "applied"); + throw new Error("synthetic acknowledgment loss"); + }; + const value = Reflect.get(target, key, target); + return typeof value === "function" ? value.bind(target) : value; + }}); + const recovered = await claim(interrupted, request); + assert.equal(recovered.status, "recovered", JSON.stringify(recovered)); + assert.equal(writes, 1); + assert.equal((recovered.lease as JsonObject).version, 1); + const committed = await head(store); + assert.equal(committed.cursor, "2"); + const replay = await claim(store, request); + assert.equal(replay.status, "replayed"); + assert.deepEqual(replay.original_receipt, recovered.original_receipt); + assert.deepEqual(replay.lease, recovered.lease); + assert.deepEqual(await head(store), committed); + }); + + } +} diff --git a/tests/control_plane_ts/command_observation_conformance.ts b/tests/control_plane_ts/command_observation_conformance.ts new file mode 100644 index 0000000000..bc7ff1f7d0 --- /dev/null +++ b/tests/control_plane_ts/command_observation_conformance.ts @@ -0,0 +1,125 @@ +/** Ordered concurrency, not timer luck: the first miss precedes a peer commit + * and the next head observes it. Every provider must recover that operation. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {commitTeamPlan} from "../../loopx/control_plane/work_items/team_plan_authority.ts"; +import {executeCoordinationTodoClaim} from "../../loopx/control_plane/coordination/todo_claim.ts"; +import {executeCanonicalTaskLeaseLifecycle} from "../../loopx/control_plane/coordination/task_lease_lifecycle.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; + +import {commandObservationCases, commandObservationScenario} from "./command_observation_scenarios.ts"; + +function deferred() { + let resolve!: () => void; + const promise = new Promise(done => {resolve = done;}); + return {promise, resolve}; +} + +/** Only the first absent receipt is delayed. All effects use the real backend. */ +export function pauseMissingReceipt(store: AuthorityStore) { + const missing = deferred(), release = deferred(); + let reads = 0, writes = 0; + const delayed = new Proxy(store, {get(target, key) { + if (key === "readReceipt") return async (operation: string) => { + const result = await target.readReceipt(operation); + if (++reads === 1) { + assert.equal(result.status, "missing"); + missing.resolve(); + await release.promise; + } + return result; + }; + if (key === "commitAuthority") return async (...args: Parameters) => { + writes++; + return target.commitAuthority(...args); + }; + const value = Reflect.get(target, key, target); + return typeof value === "function" ? value.bind(target) : value; + }}); + return {store: delayed, missing: missing.promise, release: release.resolve, writes: () => writes}; +} + +function teamRequest(): JsonObject { + return {goal_id: "goal-a", actor_agent_id: null, registered_agents: ["agent-a", "agent-b"], + supported_action_kinds: ["implement"], observed_at: "2026-09-13T10:05:00Z", + expected_state_fingerprint: "reviewed", current_state_fingerprint: "reviewed", + plan: {schema_version: "steward_team_plan_preview_v0", kind: "steward_team_plan_preview", + goal_id: "goal-a", proposal_id: "command-observation", objective: "Deliver both independent acceptance results", + quota_envelope: {slots: 2}, stop_condition: "Owner closes the request", + lanes: ["agent-a", "agent-b"].map(agent => ({lane_id: `lane-${agent}`, agent_id: agent, + acceptance: `Independent evidence from ${agent}`, first_todo: {text: "Same bounded work", + priority: "P1", task_class: "advancement_task", action_kind: "implement"}}))}}; +} + +export function registerCommandObservationConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + for (const schema of ["native", "legacy"] as const) { + for (const kind of commandObservationCases) { + test(`${provider} command observation: ${kind} recovers the ${schema} operation before new admission`, async t => { + const {store, contender} = await factory(t); + const {run} = await commandObservationScenario(kind, schema, store); + const paused = pauseMissingReceipt(store); + const outcome = run(paused.store).then(value => ({value}), error => ({error})); + await Promise.race([paused.missing, outcome.then(result => {throw new Error(`command returned before receipt lookup: ${JSON.stringify(result)}`);})]); + const winner = await run(contender).finally(paused.release); + assert.equal(winner.status, "applied", JSON.stringify(winner)); + const after = await contender.loadAuthority(); + const replay = await outcome; + assert.ok("value" in replay, `late receipt must precede ${kind} admission`); + assert.equal(replay.value.status, "replayed", JSON.stringify(replay.value)); + assert.equal(replay.value.provider_revision, winner.provider_revision); + assert.equal(replay.value.cursor, winner.cursor); + assert.equal(replay.value.changed, false); + assert.equal(paused.writes(), 0, "historical recovery must not submit a second transaction"); + assert.deepEqual(await store.loadAuthority(), after, "all unrelated Todos, decisions and leases survive replay"); + }); + } + test(`${provider} command observation: concurrent team replay preserves the full ${schema} graph`, async t => { + const {store, contender} = await factory(t); + const fixture = productionScaleCoordinationFixture("goal-a", schema); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: fixture.projection, events: [], receipts: []})).status, "applied"); + const input = teamRequest(), paused = pauseMissingReceipt(store); + const later = commitTeamPlan(paused.store, input); + // Attach a handler before releasing the barrier, including on failure. + const outcome = later.then(value => ({value}), error => ({error})); + await Promise.race([paused.missing, outcome.then(result => {throw new Error(`command returned before receipt lookup: ${JSON.stringify(result)}`);})]); + const first = await commitTeamPlan(contender, input).finally(paused.release); + assert.equal(first.status, "applied", JSON.stringify(first)); + const afterWinner = await contender.loadAuthority(); + const replay = await outcome; + assert.ok("value" in replay, "same-operation replay must not run create admission against the committed Todos"); + assert.equal(replay.value.status, "replayed", JSON.stringify(replay.value)); + assert.equal(paused.writes(), 0); + assert.equal(replay.value.provider_revision, first.provider_revision); + assert.deepEqual(await store.loadAuthority(), afterWinner); + }); + + test(`${provider} command observation: retired atomic claim lease cannot grant ${schema} execution`, async t => { + const {store} = await factory(t); + const projection = authorityProjectionFixture("goal-a", [{todo_id: "todo_claim_target", role: "agent", + text: "Adopt and execute a bounded task", status: "open", done: false, archive_state: "active", + required_write_scopes: ["command-observation/**"]}], [], schema, {handoff_mode: "hard_lease"}); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: projection, events: [], receipts: []})).status, "applied"); + const input = {goal_id: "goal-a", todo_id: "todo_claim_target", claimed_by: "agent-a", actor_agent_id: "agent-a", + expected_role: "agent" as const, registered_agents: ["agent-a", "agent-b"], operation_id: "claim-acquire", + lease_request: {idempotency_key: "execution-a", expected_version: 0, ttl_seconds: 600}, + dry_run: false, now: new Date("2026-09-13T10:05:00Z")}; + const first = await executeCoordinationTodoClaim(store, input); + assert.equal(first.status, "applied", JSON.stringify(first)); + const released = await executeCanonicalTaskLeaseLifecycle(store, {goal_id: "goal-a", todo_id: input.todo_id, + operation: "release", owner: "agent-a", idempotency_key: "execution-a", expected_version: 1, + ttl_seconds: null, registered_agents: input.registered_agents, now: input.now}); + assert.equal(released.status, "applied", JSON.stringify(released)); + const before = await store.loadAuthority(); + const replay = await executeCoordinationTodoClaim(store, input); + assert.equal(replay.status, "failed", "historical claim receipt must not return an active retired lease"); + assert.equal(replay.reason_code, "idempotency_key_reuse"); + assert.deepEqual(await store.loadAuthority(), before); + }); + } +} diff --git a/tests/control_plane_ts/command_observation_scenarios.ts b/tests/control_plane_ts/command_observation_scenarios.ts new file mode 100644 index 0000000000..6b7e66461d --- /dev/null +++ b/tests/control_plane_ts/command_observation_scenarios.ts @@ -0,0 +1,118 @@ +/** Real command entrypoints over the mixed production-scale graph. The scenario + * owns intent and expected behavior; the race harness owns only read ordering. */ +import assert from "node:assert/strict"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {executeCoordinationTodoCreate} from "../../loopx/control_plane/coordination/todo_create.ts"; +import {executeCoordinationTodoUpdate} from "../../loopx/control_plane/coordination/todo_update.ts"; +import {executeCoordinationTodoArchiveCompleted} from "../../loopx/control_plane/coordination/todo_archive.ts"; +import {executeCoordinationTodoTerminalLifecycle} from "../../loopx/control_plane/coordination/todo_terminal_lifecycle.ts"; +import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; +import {executeCanonicalTaskLeaseAcquire} from "../../loopx/control_plane/coordination/task_lease_acquire.ts"; +import {executeCanonicalTaskLeaseLifecycle} from "../../loopx/control_plane/coordination/task_lease_lifecycle.ts"; +import {executeCoordinationTodoClaim} from "../../loopx/control_plane/coordination/todo_claim.ts"; +import {configureGoalAcceptance, commitGoalAcceptanceVerification, inspectGoalAcceptance} from "../../loopx/control_plane/goals/acceptance_authority.ts"; +import {TODO_DOMAIN_ITEM_SCHEMA} from "../../loopx/control_plane/coordination/coordination_state_contract.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import {productionScaleCoordinationFixture, productionScaleLeasedMonitorFixture, + productionScaleLeaseAcquisitionFixture, productionScaleLeaseLifecycleFixture} from "./production_scale_coordination_fixture.ts"; + +export const commandObservationCases = ["create", "update", "claim", "claim-acquire", "archive", "supersede", + "monitor", "acquire", "renew", "release", "transfer", "acceptance-configure", "acceptance-verify"] as const; +export type CommandObservationCase = typeof commandObservationCases[number]; + +export async function commandObservationScenario(kind: CommandObservationCase, schema: "native" | "legacy", store: AuthorityStore) { + const goal = "goal-a", target = "todo_observation_target"; + const agents = ["agent-a", "agent-b"], now = new Date("2026-09-13T10:05:00Z"); + const f = productionScaleCoordinationFixture(goal, schema); + const todo = {todo_id: target, role: "agent", text: "Deliver an independently verifiable result", + task_class: "advancement_task", action_kind: "implement", status: "open", done: false, + archive_state: "active", required_write_scopes: ["command-observation/**"]}; + let projection = authorityProjectionFixture(goal, [...f.projection.todos as JsonObject[], todo], + f.projection.leases as JsonObject[], schema, {handoff_mode: "legacy"}); + const monitor = productionScaleLeasedMonitorFixture(goal, schema); + const acquire = productionScaleLeaseAcquisitionFixture(goal, schema); + const lifecycle = productionScaleLeaseLifecycleFixture(goal, schema); + if (kind === "monitor") projection = monitor.projection; + if (kind === "acquire") projection = acquire.projection; + if (["renew", "release", "transfer"].includes(kind)) projection = lifecycle.projection; + if (kind === "claim-acquire") projection.handoff_mode = "hard_lease"; + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: projection, events: [], receipts: []})).status, "applied"); + const initial = await store.loadAuthority(); + assert.equal(initial.status, "loaded"); + if (initial.status !== "loaded") throw new Error("fixture head missing"); + const common = {goal_id: goal, actor_agent_id: "agent-a", registered_agents: agents, + operation_id: `observe-${kind}`, dry_run: false, now}; + const expected_provider_revision = initial.provider_revision; + const document = {objective: "Deliver an independently verifiable result", non_goals: [], + criteria: [{id: "outcome", description: "The bounded validation passes", + validation_argv: [process.execPath, "-e", "process.exit(0)"]}], + bindings: [{todo_id: target, criterion_ids: ["outcome"]}]}; + let run: (backend: AuthorityStore) => Promise; + switch (kind) { + case "create": + run = backend => executeCoordinationTodoCreate(backend, {...common, + todo: {...todo, schema_version: TODO_DOMAIN_ITEM_SCHEMA, todo_id: "todo_observation_created", text: "A distinct independently deliverable task"}}); + break; + case "update": + run = backend => executeCoordinationTodoUpdate(backend, {...common, todo_id: target, expected_role: "agent", + patch: {text: "Refined acceptance", note: "Preserve independent evidence"}, clear_fields: [], expected_provider_revision}); + break; + case "claim": + case "claim-acquire": + run = backend => executeCoordinationTodoClaim(backend, {...common, todo_id: target, expected_role: "agent", + claimed_by: "agent-a", expected_provider_revision, + ...(kind === "claim-acquire" ? {lease_request: {idempotency_key: "execution-a", expected_version: 0, ttl_seconds: 600}} : {})}); + break; + case "archive": + run = backend => executeCoordinationTodoArchiveCompleted(backend, {...common, role: "agent", max_active_done: 0, + expected_provider_revision}); + break; + case "supersede": + run = backend => executeCoordinationTodoTerminalLifecycle(backend, {...common, todo_id: target, + command: "supersede", expected_role: "agent", lifecycle_grants: [], authority_reason: null, + decision_outcome: null, lease_idempotency_key: null, lease_expected_version: null, + allow_user_gate_auto_acquire: false, requested_no_followup: false, + requested_completion_turn_key: null, requested_completion_identity_source: null, + linked_successor_todo_ids: [], successor_intents: [], note: "Replacement owns continuation", + evidence: null, reason: "The owner retired this work", clear_claim: false, + validation_declaration: null, validation_receipt: null, completion_policy_request: null, + review_basis: {provider_revision: expected_provider_revision, registry_sha256: "a".repeat(64)}}); + break; + case "monitor": + run = backend => executeCoordinationMonitorPoll(backend, {...common, now: monitor.now, + observation: {todo_id: monitor.target, generated_at: monitor.now.toISOString(), result_hash: "new-evidence", material_change: true}, + intent: {next_agent_todo: "React to changed public evidence", next_action_kind: "implement", next_claimed_by: "agent-a"}, + lease_proof: monitor.proof}); + break; + case "acquire": + run = backend => executeCanonicalTaskLeaseAcquire(backend, {goal_id: goal, todo_id: acquire.target, owner: "agent-a", + idempotency_key: acquire.acquisition.execution_key, expected_version: 0, ttl_seconds: acquire.acquisition.ttl_seconds, + write_scopes: acquire.acquisition.write_scopes, registered_agents: agents, now: new Date(acquire.scenario.now)}); + break; + case "renew": + case "release": + case "transfer": + run = backend => executeCanonicalTaskLeaseLifecycle(backend, {goal_id: goal, todo_id: lifecycle.target, + operation: kind, owner: lifecycle.scenario.owner, idempotency_key: lifecycle.scenario.execution_key, + expected_version: lifecycle.scenario.version, ttl_seconds: kind === "release" ? null : 600, + ...(kind === "transfer" ? {new_owner: "agent-b", new_idempotency_key: "receiver-b"} : {}), + registered_agents: agents, now: new Date(lifecycle.scenario.now)}); + break; + case "acceptance-configure": + run = backend => configureGoalAcceptance(backend, {goal_id: goal, actor_agent_id: null, + operation_id: common.operation_id, expected_provider_revision, document}); + break; + case "acceptance-verify": { + assert.equal((await configureGoalAcceptance(store, {goal_id: goal, actor_agent_id: null, + operation_id: "configure", expected_provider_revision, document})).status, "applied"); + const basis = await inspectGoalAcceptance(store, goal); + run = backend => commitGoalAcceptanceVerification(backend, {goal_id: goal, operation_id: common.operation_id, + expected_provider_revision: basis.provider_revision, revision: basis.revision, contract_digest: basis.contract_digest, + results: [{criterion_id: "outcome", passed: true, exit_code: 0}]}); + break; + } + } + return {run, projection}; +} diff --git a/tests/control_plane_ts/command_receipt.test.ts b/tests/control_plane_ts/command_receipt.test.ts index ef487643fb..2d63e1e5e7 100644 --- a/tests/control_plane_ts/command_receipt.test.ts +++ b/tests/control_plane_ts/command_receipt.test.ts @@ -116,3 +116,64 @@ test("commit identity mismatch fails before the effect", async () => { await assert.rejects(receipt().commit(target, {...commit, operation_id: "other"}), /identities differ/); assert.equal(writes, 0); }); + +test("observation rechecks receipt after the head, including unavailable or unparseable heads", async () => { + const heads = [ + {status: "loaded" as const, provider_revision: "after", cursor: "2", head: {invalid_domain: true}}, + {status: "missing" as const}, + {status: "unavailable" as const, reason_code: "offline", reason: "synthetic read gap"}, + ]; + for (const authority of heads) { + const calls: string[] = []; + const target = store(applied, found); + target.loadAuthority = async () => {calls.push("head"); return authority;}; + target.readReceipt = async () => {calls.push("receipt"); return found;}; + target.commitAuthority = async () => {throw new Error("observation cannot write");}; + const observation = await receipt().observe(target); + assert.deepEqual(calls, ["head", "receipt"]); + assert.equal(observation.kind, "receipt"); + if (observation.kind !== "receipt") throw new Error("historical receipt lost"); + assert.equal(observation.result.status, "replayed"); + assert.equal(observation.result.provider_revision, "historical-revision"); + assert.equal(observation.result.changed, false); + } +}); + +test("observation does not turn a receipt read failure into missing-operation permission", async () => { + for (const readback of [unavailable, {...found, receipts: [{...original, request_sha256: "different"}]}]) { + const target = store(applied, readback); + target.loadAuthority = async () => ({status: "loaded", provider_revision: "current", cursor: "2", head: {}}); + const result = await receipt().observe(target); + assert.equal(result.kind, "receipt"); + if (result.kind !== "receipt") throw new Error("uncertain receipt admitted a new decision"); + assert.equal(result.result.status, readback.status === "unavailable" ? "unavailable" : "failed"); + assert.equal(result.result.changed, false); + } +}); + +test("only confirmed absence yields the original head for domain admission", async () => { + const target = store(applied, {status: "missing"}); + const authority = {status: "loaded" as const, provider_revision: "basis", cursor: "1", head: {retained: true}}; + target.loadAuthority = async () => authority; + const result = await receipt().observe(target); + assert.equal(result.kind, "authority"); + if (result.kind !== "authority") throw new Error("missing receipt became a result"); + assert.strictEqual(result.authority, authority, "do not reinterpret or mutate the decision snapshot"); +}); + +test("commit after observation still resolves through CAS receipt recovery without retrying the write", async () => { + let committed = false, writes = 0; + const target = store(applied, {status: "missing"}); + target.loadAuthority = async () => ({status: "loaded", provider_revision: "before", cursor: "1", head: {}}); + target.readReceipt = async () => committed ? found : {status: "missing"}; + assert.equal((await receipt().observe(target)).kind, "authority"); + committed = true; // A peer wins after the bounded observation, before this CAS. + target.commitAuthority = async () => { + writes++; + return {status: "conflict", conflict_kind: "provider_revision_mismatch", current_provider_revision: "after", current_cursor: "2"}; + }; + const result = await receipt().commit(target, commit); + assert.equal(result.status, "recovered"); + assert.equal(writes, 1); + assert.equal(result.provider_revision, "historical-revision"); +}); diff --git a/tests/control_plane_ts/local_authority_runtime.test.ts b/tests/control_plane_ts/local_authority_runtime.test.ts index 25727def35..ce999d6b5b 100644 --- a/tests/control_plane_ts/local_authority_runtime.test.ts +++ b/tests/control_plane_ts/local_authority_runtime.test.ts @@ -1423,7 +1423,7 @@ test("provider-first Todo claim validates authority and hard-lease ownership", a assert.equal((unchanged.todo as Record).claimed_by, undefined); }); -test("provider-first Todo claim atomically acquires its canonical hard lease", async () => { +test("provider-first atomic claim replays current execution and rejects expired receipt proof", async () => { const root = await mkdtemp(join(tmpdir(), "loopx-local-authority-claim-lease-")); const store = new FileAuthorityStore(join(root, "authority", "file-v0"), "goal-a"); assert.equal((await store.commitAuthority({ @@ -1487,10 +1487,15 @@ test("provider-first Todo claim atomically acquires its canonical hard lease", a const replay = await claimLocalCoordinationTodo({ ...request, - observed_at: "2026-09-05T06:00:00Z", + observed_at: "2026-09-05T04:45:00Z", }); assert.equal(replay.status, "replayed"); assert.deepEqual(replay.original_receipt, applied.original_receipt); + const expiredReplay = await claimLocalCoordinationTodo({...request, observed_at: "2026-09-05T05:15:00Z"}); + assert.equal(expiredReplay.status, "failed"); + assert.equal(expiredReplay.reason_code, "idempotency_key_reuse"); + assert.equal(expiredReplay.lease, undefined); + assert.deepEqual(expiredReplay.original_receipt, applied.original_receipt); const retiredGeneration = await claimLocalCoordinationTodo({ ...request, operation_id: "todo-claim:goal-a:todo_a:fresh-after-expiry", From 2decb65d507799cc19ba38662de79afc11d2e7a0 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 14:09:11 +0800 Subject: [PATCH 3/3] docs(authority): distinguish command history from current execution proof Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- ...shared-goal-authority-state-provider-v0.md | 3 +- .../typescript-control-plane-migration-v0.md | 12 ++++ ...script-control-plane-migration-v0.zh-CN.md | 9 +++ docs/reference/canonical-lease-renew.md | 56 +++++++++++++++++-- 4 files changed, 75 insertions(+), 5 deletions(-) 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 bac0923c6a..0abf60e4bc 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3210,7 +3210,7 @@ or moving a helper is not by itself a package exit. | --- | --- | --- | | A / L1: Monitor configuration (this slice) | Existing `todo update` config enters the TS planner/CAS/receipt; delete Python's duplicate intent field catalog. Separate authoring from observed hashes, times and generations. | Ordinary CLI/API, clear/omission, active lease proof, no-op/replay, failed display delivery, complete fixture and real providers. This does not complete delegated Chat or leased polling. | | A / L2: Complete public mutation admission | User completion updates share the TS edit/terminal transaction and reviewed Chat recovery; linked decision consumption/reject/cancel/resume now commit with the source, replacing Python followthrough rules. Continue the actual CLI/Turn/Chat inventory for remaining effect-owned decisions, delegated owner actions and Monitor lifecycle transitions; [caller contract](../../reference/canonical-todo-completion-update.md). | Build on merged T1 owners, not a generic raw patch. Prove permission rejection and exact caller response; remove replaced Python admission and name every remaining unsupported command. | -| A / L3: Canonical lease lifecycle | Standalone acquire/takeover, atomic claim lease admission and maintenance reuse TS facts/decision/materialization and one provider opening fence. Explicit claimed-work transfer now commits source-authorized Todo ownership and the new lease generation together; canonical request types exclude legacy held-fence fields. Acquire success verifies current execution proof; canonical completion can recover missing display. | Full-head scope conflict, archived/ineffective holders, exact create-CAS retry, stale execution, process loss and real CLI/four-arm rehearsal are covered. [Operation and remaining callers](../../reference/canonical-lease-renew.md). Executor-held external-effect fences remain explicit work; D1–D3/default holds remain. | +| A / L3: Canonical lease lifecycle | Standalone acquire/takeover, atomic claim lease admission and maintenance reuse TS facts/decision/materialization and one provider opening fence. Explicit claimed-work transfer now commits source-authorized Todo ownership and the new lease generation together; canonical request types exclude legacy held-fence fields. Standalone acquire and atomic claim/acquire share current execution proof, including renewal, retirement and uncertain readback. Canonical commands recheck the operation receipt after loading the head, before new admission; canonical completion can recover missing display. | Full-head scope conflict, archived/ineffective holders, exact create-CAS retry, stale execution, process loss and real CLI/four-arm rehearsal are covered. [Operation and remaining callers](../../reference/canonical-lease-renew.md). Executor-held external-effect fences remain explicit work; D1–D3/default holds remain. | | B / L4: Leased Monitor poll and settlement | Current execution proof now binds CLI intent, observation/generation/independent-successor CAS and historical business receipt. Quota pending admission is frozen before the business write; recovery preserves that decision after lease retirement. | Existing L3 lease lifecycle, real File/SQLite/PostgreSQL, mixed fixtures, process death between business/quota commits, competing renewal and unchanged polling. [Operation and snapshot rehearsal](../../reference/protocols/quota-monitor-observation-receipt-v0.md). Ordinary polls leave leases unchanged and spend no quota; separate authorities stay separate. The retained grouped-Monitor observation/reactivation caller now uses Todo update v4 and the shared Monitor planner, with unchanged-group display recovery. Canonical reactivation now atomically retires retained execution and reopens the observation cycle, sharing typed admission with polling; a fresh execution still needs explicit acquisition. Grouped reconciliation now acquires/revalidates/releases its own bounded execution, recovers interrupted cleanup, and plans the complete bucket set in TS; missing evidence and ambiguous/stale targets reject. This closes that retained caller across legacy/File/SQLite; native/imported mixed fixtures exercise the same effects on real PostgreSQL. Wider L2 admission, external-effect fences and D1–D3/default remain open. | | B / L5: Consumer and display closure | Reconcile #4316, audit Turn/quota/Dashboard/Chat source reads, and finish D1 freshness/recovery through the existing projection outbox. | CLI, Lark/Chat and packaged frontend read back their affected interactions; absent/stale display, empty canonical state, pending projection and data beyond UI limits. Delete post-promotion legacy fallbacks with each consumer. | | A–C / L6: Local durability qualification | Continue contributor-owned #4224/#4328 on the selected SQLite profile; reuse File/NoKV references and complete 7.2's ledger. | Capacity, real process/crash/restore/upgrade, retained receipts/scans, consumer lag, supported runtimes/OS and the separately authorized >=10-day synthetic soak. Missing measurements remain holds. | @@ -3233,6 +3233,7 @@ PRs**, conditional on the caller audit finding no additional missing effects: | L7 capture plus L8 integrated migration | 1–2 | Mixed-writer continuity, fenced whole-Goal rehearsal, export/rollback and cohort evidence. | | L9 default and bounded retirement | 1 | New-Goal onboarding/settings/install choose the qualified profile; remove final obsolete callers. | +The command-observation/current-proof closure removes a concrete L2/L3 concurrency hold. The retained-Monitor cycle and grouped executor closure remove concrete L4 holds, not an entire remaining package: the **5–8 PR planning range remains conditional**, rather than subtracting one for a lifecycle fix. Actual remaining executor/caller coverage, diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 9fb2ce0df2..75f5f6eaeb 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -22,6 +22,18 @@ Retain T0 caller/parity inventory, T1/T2 transaction/effect convergence, T3 comp ## Current implementation checkpoint +Canonical command observation now has one typed receipt/head boundary. Team, +Todo creation/edit/claim/terminal/archive, Monitor, lease maintenance and Goal +acceptance recheck receipts after the head read before interpreting new state. +This repairs same-operation races without provider API changes, write retries +or a Python copy of the decision. Standalone and atomic claim acquisition share +current lease proof; renewed proof is returned without rewriting history, while +retired execution and unavailable current authority cannot return stale success. +This closes a concurrency/current-proof slice of L2/L3, not whole-Goal migration, +default onboarding, contributor-owned SQLite D2 or T4 Python retirement. The +[operator contract](../../reference/canonical-lease-renew.md#commit-retry-and-readback) +distinguishes historical results from present execution. + Terminal review and validation now converge in the existing TS terminal owner. Agent completion and Monitor stop reuse Chat's canonical receipt-first recovery and display acknowledgement; v2 binds validation continuation to its source 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 8b1428e729..04306f6d3c 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 @@ -21,6 +21,15 @@ ## 当前实现检查点 +Canonical command 的 receipt/head 观察顺序统一归属 TS:团队规划、Todo 创建/ +修改/领取/终态/归档、Monitor、lease 维护和 Goal acceptance 在读 head 后复查原 +receipt,再执行新准入。这修复同 operation 并发竞争,不扩展 provider API、不 +自动重试写入,也不在 Python 复制决定。独立 acquire 与原子 claim/acquire 共用 +当前 lease 证明:续租返回新 proof 而保留原 receipt,退役执行或当前 authority +不可读均不能返回旧成功。它只关闭 L2/L3 的并发和当前证明缺口,不代表全 Goal +迁移、默认启用、contributor 的 SQLite D2 或 T4 Python 退役完成。见 +[操作和恢复合同](../../reference/canonical-lease-renew.md#commit-retry-and-readback)。 + 终结审核与验证已收敛到既有 TS terminal owner:Agent 完成、Monitor 停止复用 Chat 先恢复 canonical 回执再确认显示的路径;v2 把验证 continuation 绑定来源 revision, 准入/回放之后才请求私有声明。删除 Python 的终结操作审核分流和提前解析声明编排。 diff --git a/docs/reference/canonical-lease-renew.md b/docs/reference/canonical-lease-renew.md index 61747685de..2204822136 100644 --- a/docs/reference/canonical-lease-renew.md +++ b/docs/reference/canonical-lease-renew.md @@ -174,8 +174,49 @@ with the same owner/key/epoch. Replay after renewal returns the current version/expiry in `lease`, with the unchanged original decision in `original_receipt`; `current_provider_revision/current_cursor` identify that readback. A transferred, expired or released execution cannot be revived by its -old receipt. Claim receipts and maintenance receipts retain their historical -semantics; they are not acquire responses. +old receipt. Atomic `todo claim` with `--task-lease-idempotency-key` now uses the +same current-proof owner after commit/recovery, including receipt-only no-ops. +A plain claim without an acquisition request and maintenance receipts retain +historical semantics; they do not grant new execution. + +For atomic adoption, freeze both identities across an uncertain response: + +```bash +loopx --registry registry.json todo claim \ + --goal-id example-goal --todo-id todo_work \ + --claimed-by agent-a --agent-id agent-a \ + --claim-operation-id adopt-work-a \ + --task-lease-idempotency-key execution-a --task-lease-expected-version 0 +``` + +The operation id identifies the claim transaction; the execution key identifies +the lease generation. Repeating this exact command after renewal returns the +renewed lease and the original receipt. After release or expiry it fails with +`idempotency_key_reuse`; after a claim transfer it fails current-owner admission. +Inspect the current state before choosing a new execution key and operation id. +If the receipt exists but the current head cannot be read, the command returns +`ambiguous` with same-operation recovery, not the historical active lease. +Switching away from `hard_lease` also invalidates atomic claim/acquire success. + +| Readback | Meaning | Next action | +| --- | --- | --- | +| `replayed`, same epoch, newer lease version | Same execution was renewed | Use the returned current lease version. | +| `idempotency_key_reuse` | The receipt belongs to a retired execution | Inspect, then request a new execution with a new key. | +| `owner_conflicts_with_claim` | Current Todo ownership changed | Let the current owner continue or use an authorized handover. | +| `canonical_acquire_readback_required` | History is known, current proof is unavailable | Restore the provider and retry the same operation. | +| Acceptance/source rejection | Current control-plane authority changed | Resolve that boundary before attempting work. | + +A successful readback is still a point-in-time proof, not a lock over subsequent +external effects. Execution must retain its existing mutation fences; this +change does not close the remaining external-effect fencing work. + +Canonical commands recheck their receipt after reading the decision head. +This handles a peer committing the same operation between the first absent +receipt and the head read: recovery precedes duplicate-ID, stale-revision, +Monitor-generation and other new-admission checks. The second read does not +lock the head; commits after it still resolve through CAS and receipt recovery. +Archive preview remains a current-state preview with no receipt lookup. +No provider becomes the default and no legacy writer is re-enabled by this change. The canonical-only acquire and lifecycle requests are closed and versioned. Joint claim transfer uses `loopx_canonical_task_lease_claim_transfer_request_v0`, so an older @@ -299,8 +340,15 @@ expected-version 0 的原样重试可恢复回执,改变参数会拒绝。 Acquire 的成功还必须核对当前有效 owner/key/epoch 和资格。同一执行续约后, 重试返回 `lease` 中的当前版本/到期时间,以及 `original_receipt` 中不可变的原始 决定;`current_provider_revision/current_cursor` 标识当前读回。已转交、到期或释放 -的旧执行不能凭 receipt 复活。Todo claim 和维护 receipt 仍是历史语义,不能将其 -当成新的 acquire 响应。 +的旧执行不能凭 receipt 复活。携带 lease 请求的原子 Todo claim 也共享这项检查, +包括首次提交、丢响应恢复、历史重放和只保存 receipt 的 no-op。普通 claim 和维护 +receipt 仍是历史语义,不授予新执行权。当前 head 不可读时返回 ambiguous,要求 +沿用原 operation id 恢复;不能把原 receipt 中的 active lease 当作当前证明。 + +各 canonical 命令在读 head 后再次查原 receipt,解决另一调用恰在第一次查无回执后 +提交成功的竞争。回执优先于新一轮的重复 ID、陈旧 revision 和 Monitor generation +校验;之后仍由 CAS 防止覆盖并发更新。归档 dry-run 继续只看当前状态,不重放历史。 +公开参数、provider 默认值和 legacy 路径保持不变。 canonical acquire 与 lifecycle 各有封闭 wire,旧 renew wire 只接受 renew;旧 runtime 不识别新 acquire schema。fence、provider、注册源变化和 CAS 错误不回退