Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions docs/architecture/rfcs/typescript-control-plane-migration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 预算和展示窗口都不允许截断
Expand Down
17 changes: 17 additions & 0 deletions loopx/control_plane/coordination/todo_claim.ts
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,23 @@ export async function executeCoordinationTodoClaim(
store: AuthorityStore,
rawInput: CoordinationTodoClaimInput,
authoritySourcesCurrent: AuthoritySourceCheck = uncheckedAuthoritySource,
): Promise<CoordinationTodoClaimResult> {
// 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<CoordinationTodoClaimResult> {
let input: CoordinationTodoClaimInput;
try {
Expand Down
12 changes: 8 additions & 4 deletions tests/control_plane_ts/authority_store_conformance.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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"};
Expand Down Expand Up @@ -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();
Expand Down
120 changes: 120 additions & 0 deletions tests/control_plane_ts/claim_contention_conformance.ts
Original file line number Diff line number Diff line change
@@ -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<void>): 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<void>(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");
});
}
}
Loading