diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 7191f79bc7..b69aee08fc 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -1151,6 +1151,18 @@ debit. This closes the demonstrated T3 consumer gap, not D1–D3, provider promotion, or the remaining Python transaction adapters. See the [operating contract](../../quota-allocation.md#receipt-backed-settlement-progress). +**Canonical claim contention.** The TS claim command now retries a conclusive +provider revision CAS rejection at most twice, using the same operation and +lease keys. Every attempt rereads the receipt and complete authority and +revalidates source registration, Todo eligibility, acceptance and lease scopes. +An explicit provider revision or transfer grant stays pinned; ambiguous writes +retain existing receipt recovery. Independent claims can both finish while +same-Todo or overlapping-scope claims still admit one owner. This adopts the +shared-authority conflict contract at the canonical writer; it does not reserve +recommendations, change local writer serialization, or qualify sustained +multi-host throughput. CLI claim callers inherit the behavior; no frontend or +Lark action contract changes. + **Long-history transport boundary.** Replan history still has one TS decision owner. Small requests retain the inline codec; larger complete fact snapshots travel through a private, digest-bound local file reference. The same reducer 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 14e75dafe0..579f50f5af 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 @@ -838,6 +838,8 @@ T3/D1 reader,未完成全部 Todo writer、retention/compaction 或 promotion 配额准入与结算消费者现在从统一 Todo reader 读取完整来源,在显示压缩前解析显式 Todo 选择。它删除直接追加 Markdown 候选的路径,保留 promote 前的事件适配;promote 后权威为空或不可读都不能复活展示行。结算进度由现有 TS 回执链归约,Python 负责完整身份命令及 JSON/Markdown 展示。现有幂等 writer 可补齐缺失的 spend 回执而不再次扣款。这关闭已复现的 T3 消费者缺口,不代表 D1–D3、provider promotion 或剩余 Python 事务适配已完成。操作语义见[结算进度契约](../../quota-allocation.md#receipt-backed-settlement-progress)。 +**Canonical claim 争抢。** TS claim 命令现在对明确的 provider revision CAS 拒绝最多重试两次,保持同一 operation 与 lease key;每次重新读取回执和完整权威,复核来源注册、Todo 资格、acceptance 和 lease 写范围。显式 provider revision 或 transfer grant 保持固定,写入结果不明确时沿用回执恢复。独立认领可以同时完成,同 Todo 或重叠写范围仍只接受一方。这让 canonical writer 采用 shared-authority 冲突契约,不预留推荐项、不改变本机 writer 串行化,也不证明持续多主机吞吐。CLI claim 调用方继承该行为,frontend/Lark 动作契约不变。 + **长历史传输边界。** Replan 历史仍由一个 TS owner 决策。小请求保留 inline codec;较大的完整事实快照通过私有临时文件和摘要绑定的引用传递。同一 reducer 校验全部记录、agent 作用域内的 ACK 及重试身份;RPC 预算和展示窗口都不允许截断 diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index a906a1c2b8..a024eb6160 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -330,6 +330,23 @@ export async function executeCoordinationTodoClaim( store: AuthorityStore, rawInput: CoordinationTodoClaimInput, authoritySourcesCurrent: AuthoritySourceCheck = uncheckedAuthoritySource, +): Promise { + // Retry only a conclusive provider CAS rejection. Each attempt rereads the + // original receipt and the complete head, then rechecks source authorization, + // claim ownership, acceptance and lease/write-scope exclusion. Pinned revisions + // (including handoff grants) must return to their caller for a fresh observation. + for (let attempt = 0; ; attempt++) { + const result = await executeClaimAttempt(store, rawInput, authoritySourcesCurrent); + if (attempt >= 2 || rawInput.expected_provider_revision !== undefined || + rawInput.transfer_grant !== undefined || result.status !== "conflict" || + result.conflict_kind !== "provider_revision_mismatch") return result; + } +} + +async function executeClaimAttempt( + store: AuthorityStore, + rawInput: CoordinationTodoClaimInput, + authoritySourcesCurrent: AuthoritySourceCheck, ): Promise { let input: CoordinationTodoClaimInput; try { diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 8a2e382e2e..557e8eddf5 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,3 +1,4 @@ +import {registerClaimContentionConformance} from "./claim_contention_conformance.ts"; import {registerMonitorGateScopeConformance} from "./monitor_gate_scope_conformance.ts"; import {projectCoordinationSource, SOURCE_PROJECTION_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/source_projection.ts"; import {registerCanonicalSnapshotConformance} from "./canonical_snapshot_conformance.ts"; @@ -290,6 +291,7 @@ export function registerAuthorityStoreConformance( registerLeaseAcquisitionConformance(providerName, factory); registerCommandObservationConformance(providerName, factory); registerClaimAcquisitionProofConformance(providerName, factory); + registerClaimContentionConformance(providerName, factory); registerAuthorityScanConformance(providerName, factory); registerOwnershipObservationConformance(providerName, factory); registerSuccessionReadConformance(providerName, factory); @@ -1398,7 +1400,7 @@ export function registerAuthorityStoreConformance( assert.deepEqual(afterIdempotent.head, loaded.head); assert.equal((await store.readReceipt("claim-and-acquire-idempotent")).status, "found"); }); - test(`${providerName} conformance: competing ownership transactions cannot split claim and lease (${native ? "native" : "v0"})`, async (t) => { + test(`${providerName} conformance: competing ownership revalidates the winning claim and lease (${native ? "native" : "v0"})`, async (t) => { const {store, contender} = await factory(t); const goalId = "goal-competing-ownership"; const projection = { @@ -1435,8 +1437,9 @@ export function registerAuthorityStoreConformance( ])); assert.deepEqual( results.map((result) => result.status).sort(), - ["applied", "conflict"], + ["applied", "failed"], ); + assert.equal(results.find(result => result.status === "failed")?.reason_code, "claim_owner_mismatch"); const winnerIndex = results.findIndex((result) => result.status === "applied"); assert.notEqual(winnerIndex, -1); const winner = winnerIndex === 0 ? "agent-a" : "agent-b"; @@ -1543,7 +1546,7 @@ export function registerAuthorityStoreConformance( assert.deepEqual(await store.loadAuthority(), transferred); }); for (const fault of ["lease_replaced", "lost_response"] as const) { - test(`${providerName} conformance: hard-lease claim ${fault} (${native ? "native" : "v0"})`, async (t) => { + test(`${providerName} conformance: hard-lease claim revalidation ${fault} (${native ? "native" : "v0"})`, async (t) => { const {store, contender} = await factory(t); const goalId = "goal-claim"; const projection = {...todoClaimProjection(goalId, native), handoff_mode: "hard_lease"}; @@ -1596,7 +1599,8 @@ export function registerAuthorityStoreConformance( }, }; const result = await executeCoordinationTodoClaim(intercepted, request); - assert.equal(result.status, fault === "lease_replaced" ? "conflict" : "recovered"); + assert.equal(result.status, fault === "lease_replaced" ? "failed" : "recovered"); + if (fault === "lease_replaced") assert.equal(result.reason_code, "handoff_mode_requires_lease"); assert.equal((await store.readReceipt(request.operation_id)).status, fault === "lease_replaced" ? "missing" : "found"); const after = await store.loadAuthority(); diff --git a/tests/control_plane_ts/claim_contention_conformance.ts b/tests/control_plane_ts/claim_contention_conformance.ts new file mode 100644 index 0000000000..9df0b71c7d --- /dev/null +++ b/tests/control_plane_ts/claim_contention_conformance.ts @@ -0,0 +1,120 @@ +/** Exercise command retries through real provider CAS, with deterministic races. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {configureGoalAcceptance} from "../../loopx/control_plane/goals/acceptance_authority.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 {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; + +function beforeCommit(store: AuthorityStore, before: (input: AuthorityStoreCommit) => Promise): AuthorityStore { + return new Proxy(store, {get(target, key) { + if (key === "commitAuthority") return async (input: AuthorityStoreCommit) => { + await before(input); + return target.commitAuthority(input); + }; + const value = Reflect.get(target, key); + return typeof value === "function" ? value.bind(target) : value; + }}); +} + +export function registerClaimContentionConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + async function setup(t: test.TestContext, overlap = false) { + const {store, contender} = await factory(t); + const todos = ["todo_one", "todo_two"].map(todo_id => ({todo_id, role: "agent", status: "open", + done: false, archive_state: "active", task_class: "advancement_task", text: "Synthetic work", + claimed_by: null, required_write_scopes: [overlap ? "shared/**" : `${todo_id}/**`]})); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + next_projection: authorityProjectionFixture("goal-a", todos, [], "native", {handoff_mode: "hard_lease"}), + events: [], receipts: []})).status, "applied"); + const requests = ["agent-a", "agent-b"].map((agent, i) => ({goal_id: "goal-a", todo_id: todos[i].todo_id, + claimed_by: agent, actor_agent_id: agent, expected_role: "agent", registered_agents: ["agent-a", "agent-b"], + operation_id: `claim-${agent}`, lease_request: {idempotency_key: `execution-${agent}`, expected_version: 0, ttl_seconds: 600}, + dry_run: false, now: new Date("2026-09-30T10:00:00Z")})); + return {store, contender, requests}; + } + + for (const scenario of ["independent", "same-todo", "overlap"] as const) { + test(`${provider} claim contention: ${scenario} revalidates the latest head`, async t => { + const {store, contender, requests} = await setup(t, scenario === "overlap"); + if (scenario === "same-todo") requests[1].todo_id = requests[0].todo_id; + let arrivals = 0; + let release!: () => void; + const barrier = new Promise(resolve => { release = resolve; }); + const stores = [store, contender].map(value => beforeCommit(value, async () => { + if (++arrivals === 2) release(); + await barrier; + })); + const results = await Promise.all(requests.map((request, i) => claim(stores[i], request))); + assert.equal(results.filter(r => r.status === "applied").length, scenario === "independent" ? 2 : 1, + JSON.stringify(results)); + if (scenario !== "independent") { + const loser = results.findIndex(r => r.status !== "applied"); + assert.equal(results[loser].status, "failed", JSON.stringify(results[loser])); + assert.equal(results[loser].reason_code, scenario === "same-todo" ? "claim_owner_mismatch" : "write_scope_conflict"); + assert.equal((await store.readReceipt(requests[loser].operation_id)).status, "missing"); + } + for (let i = 0; i < results.length; i++) if (results[i].status === "applied") { + const receipt = await store.readReceipt(requests[i].operation_id); + assert.equal(receipt.status, "found"); + assert.equal((await claim(stores[i], requests[i])).status, "replayed"); + assert.deepEqual(await store.readReceipt(requests[i].operation_id), receipt); + } + assert.equal(arrivals, scenario === "independent" ? 3 : 2); + const final = await store.loadAuthority(); + if (final.status !== "loaded") throw new Error("missing final authority"); + const winners = requests.filter((_, i) => results[i].status === "applied"); + assert.equal((final.head.leases as JsonObject[]).length, winners.length); + for (const winner of winners) { + assert.equal((final.head.todos as JsonObject[]).find(todo => todo.todo_id === winner.todo_id)?.claimed_by, + winner.claimed_by, "rebasing an independent claim must retain the previous winner"); + } + }); + } + + test(`${provider} claim contention: newly required acceptance is rechecked`, async t => { + const {store, contender, requests} = await setup(t); + let commits = 0; + const wrapped = beforeCommit(store, async () => { + commits++; + const current = await contender.loadAuthority(); + if (current.status !== "loaded") throw new Error("missing fixture"); + assert.equal((await configureGoalAcceptance(contender, { + goal_id: "goal-a", actor_agent_id: null, operation_id: "require-acceptance", + expected_provider_revision: current.provider_revision, + document: {scope: {kind: "all_advancement"}, objective: "Validate the accepted work", non_goals: [], + criteria: [{id: "outcome", description: "Validation passes", validation_argv: ["python", "-V"]}], + bindings: []}, + })).status, "applied"); + }); + const result = await claim(wrapped, requests[0]); + assert.equal(commits, 1); + assert.equal(result.status, "failed"); + assert.equal((result.goal_acceptance_guard as JsonObject).allowed, false); + assert.equal((await store.readReceipt(requests[0].operation_id)).status, "missing"); + }); + + for (const scenario of ["pinned", "source-change", "exhausted"] as const) { + test(`${provider} claim contention: ${scenario} stops without a claim receipt`, async t => { + const {store, contender, requests} = await setup(t); + const initial = await store.loadAuthority(); + assert.equal(initial.status, "loaded"); + if (initial.status !== "loaded") throw new Error("missing fixture"); + let commits = 0; + const wrapped = beforeCommit(store, async () => { + const current = await contender.loadAuthority(); + if (current.status !== "loaded") throw new Error("missing fixture"); + assert.equal((await contender.commitAuthority({operation_id: `peer-${++commits}`, + expected_provider_revision: current.provider_revision, next_projection: current.head, + events: [], receipts: []})).status, "applied"); + }); + const result = await claim(wrapped, {...requests[0], ...(scenario === "pinned" + ? {expected_provider_revision: initial.provider_revision} : {})}, + async () => scenario !== "source-change" || commits === 0); + assert.equal(commits, scenario === "exhausted" ? 3 : 1); + assert.equal(result.status, scenario === "source-change" ? "failed" : "conflict", JSON.stringify(result)); + assert.equal((await store.readReceipt(requests[0].operation_id)).status, "missing"); + }); + } +}