diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index dfc4ce7bdd..738c97b9e6 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -235,7 +235,7 @@ These priorities do not change live Goal quota or authorize experiments/cloud re ### R5: TS Convergence and Local Persistence -Existing-lease L3 checkpoint: renew/transfer/release share one TS transaction and provider handle; real File/SQLite/PostgreSQL and immutable-legacy comparison cover handover, cleanup and historical replay. [Boundary and remaining callers](../../reference/canonical-lease-renew.md); R5, D2/D3 and new-Goal default qualification remain open. +L3 checkpoint: standalone acquisition/takeover, atomic claim admission and maintenance share typed lease facts/rules and provider opening. Exact acquisition retry verifies current execution proof; real CLI completion can recover missing Markdown display. Full-state scope conflicts, process interruption and File/SQLite/PostgreSQL read-only rehearsal are covered. [Remaining executor and integration boundaries](../../reference/canonical-lease-renew.md); R5, D2/D3 and default qualification remain open. - **Owner:** TS T0–T4 and shared-authority D1–D3; retain their numbering and gates. - **Selection:** prioritize an entire hot-path transaction or recovery lifecycle used by R1–R4. Record before/after callers, owners, crossings, actual deletions and performance. Stop adding per-field Python→TS RPCs; do not rebuild the merged Todo update. diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md index 6d6e350ae3..c9901d0398 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md @@ -235,7 +235,7 @@ R2 的一条依赖必须通过真实 LoopX Agent 间的请求/产物交接完成 ### R5:TS 收敛与本地持久化 -既有 lease 的 L3 检查点:renew/transfer/release 共用 TS 事务和 provider handle;真实 File/SQLite/PostgreSQL 与不可变 legacy 对照覆盖转交、清理和历史 replay。[交付边界与剩余 caller](../../reference/canonical-lease-renew.md);R5、D2/D3 和新 Goal 默认化资格仍未完成。 +L3 检查点:独立领取/接管、原子 claim 准入与维护共用 typed lease facts/rules 和 provider opening;原领取重试校验当前执行 proof,真实 CLI 完成可恢复缺失 Markdown 展示。覆盖完整 scope 冲突、进程中断及 File/SQLite/PostgreSQL 只读演练。[剩余 executor 与集成边界](../../reference/canonical-lease-renew.md);R5、D2/D3 和默认化资格仍未完成。 - **Owner:** TS RFC T0–T4、shared-authority D1–D3;保留两套编号及原门禁。 - **选择规则:** 优先迁移 R1–R4 热路径的一笔完整事务或恢复生命周期,附前后 caller/owner/crossing 表、实际删除和性能证据。不要继续按单字段增加 Python→TS RPC;不要重建已合入的 Todo update。 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 172744f591..91db2e1ee2 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3039,7 +3039,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 | Inventory actual CLI/Turn/Chat callers; close remaining effect-owned user decisions, delegated owner actions and Monitor lifecycle transitions with validated actor/grant facts. | 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 | Existing-lease renew/transfer/release now share one TS transaction, record materializer and provider opening handle; CLI readback and service-factory PostgreSQL are exercised. | Native/imported scale fixtures, real process loss, stale proof, no-op receipts and historical replay pass. [Operation and four-arm rehearsal](../../reference/canonical-lease-renew.md) distinguish preserved legacy outcomes from improved replay. Acquire/reclaim and executor fence adoption remain explicit caller 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. 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. | | B / L4: Leased Monitor poll and settlement | Compose observation, generation and independent successors with the current lease fence. Reuse the existing quota settlement protocol and exact business receipt. | L2/L3; real polling failure, duplicate/no-change observations, crash between business and quota settlement, and competing writers. Do not pretend separate authorities share a database transaction. | | 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. | diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index e4c0ab3d1b..d86fc08a7c 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -2409,7 +2409,7 @@ canonical renew 候选,#4328 是 SQLite D2 首批测量/恢复候选;它 | --- | --- | --- | | A/L1:Monitor 配置(本切片) | 现有 `todo update` 配置进入 TS planner/CAS/receipt,删除 Python 重复 intent 字段表;区分配置与观察 hash、时间、代数。 | 普通 CLI/API、清除/省略、active lease proof、no-op/replay、展示失败恢复、完整 fixture 和真实 provider。不宣称完成委托 Chat 或 leased polling。 | | A/L2:公共 mutation admission 闭合 | 盘点 CLI/Turn/Chat 实际 caller;以可信 actor/grant 事实闭合剩余 effect-owned 用户决策、委托 owner 动作和 Monitor lifecycle。 | 复用已合并 T1 owner,不开通通用 raw patch;验证权限拒绝和 caller 响应,删除替代的 Python admission,列全未支持命令。 | -| A/L3:canonical lease 生命周期 | 既有 lease 的 renew/transfer/release 已共用 TS 整笔事务、record materializer 与 provider opening handle;覆盖 CLI 回读和 service-factory PostgreSQL。 | native/imported 规模 fixture、真实进程中断、旧 proof、no-op receipt 与历史 replay 已验证。[操作与四臂演练](../../reference/canonical-lease-renew.md) 区分保持的旧结果和改进的 replay。Acquire/reclaim 与 executor fence 接入仍是明确 caller 工作;D1–D3/default hold 保留。 | +| A/L3:canonical lease 生命周期 | 独立 acquire/接管、原子 claim 的 lease 准入及维护复用 TS facts/decision/materializer 与同一 provider opening fence。Acquire 成功必须校验当前执行 proof;canonical 完成可恢复缺失展示。 | 已覆盖完整 head scope 冲突、归档/失效 holder、创建 CAS 原样重试、旧执行、进程中断、真实 CLI 与四臂演练。[操作及剩余 caller](../../reference/canonical-lease-renew.md)。跨外部 effect 的 executor 持锁 fence 仍为明确工作;保留 D1–D3/default hold。 | | B/L4:leased Monitor poll 与 settlement | 组合观察、变化代数、独立 successor 和现有 lease fence;复用 quota settlement 与精确业务回执。 | L2/L3;真实 polling 失败、重复/无变化、业务提交到 quota settlement 间崩溃和并发。不能假装不同 authority 共享一个数据库事务。 | | B/L5:consumer 与展示闭合 | 核对 #4316,审计 Turn/quota/Dashboard/Chat 的来源,复用 projection outbox 完成 D1 新鲜度和恢复。 | 验证 CLI、Lark/Chat、打包 frontend 的受影响交互;缺失/陈旧展示、权威空状态、pending 投影及超过 UI 上限的数据。逐个删除晋升后的 legacy fallback。 | | A–C/L6:本地持久化资格 | 延续 contributor 认领的 #4224/#4328,在选定 SQLite profile 上补齐第 7.2 节 ledger,复用 File/NoKV 对照。 | capacity、真实进程/crash/restore/upgrade、历史 receipt/scan、consumer lag、支持的 runtime/OS,以及另行授权的 >=10 天合成 soak。缺项继续 hold。 | diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index fa9b056d56..6fc16bca00 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -109,19 +109,23 @@ native creation, archival, receipt replay, and store reopen are tested without Markdown metadata. Python only adapts the typed read result to the compatibility summary. This is a contract checkpoint, not a completed CLI lifecycle cutover. -### Existing-lease transaction closure (2026-09-17) - -Renew, transfer and release now share `coordination/task_lease_lifecycle.ts` and -the existing typed lifecycle decision/record materializer. The Python adapter -routes promoted commands once; the local opening handle owns provider identity, -including service-factory PostgreSQL. Renew-only constructors and duplicated -record materialization are retired; legacy storage remains for its real callers. -Transfer preserves Todo claims/scopes; release accepts expired or deregistered -owners only with matching proof. No-op cleanup seals a receipt; historical -replay cannot reacquire execution authority. Archived renew/transfer and unsafe -generation increments fail closed. See [operation, compatibility and four-arm -rehearsal](../../reference/canonical-lease-renew.md). This closes existing-lease -L3 mutations, not acquisition/reclaim, executor fences, D2/D3 or new-Goal defaults. +### Lease acquisition and lifecycle convergence (2026-09-18) + +Standalone acquire/takeover and maintenance now share the local provider/source +fence. `task_lease_acquire_decision.ts` owns acquire admission and materialization; +legacy acquire and canonical atomic Todo claim reuse it. `task_lease_state.ts` +provides full canonical facts, including archived-holder exclusion from scope +conflicts. Python sends registration facts through one native request, without +reconstructing the canonical Todo/lease head. Generation exhaustion fails closed. + +An acquire receipt alone is not current execution authority: exact create-CAS +retry recovers the original decision and verifies the current owner/key/epoch; +renewal returns current proof while expiry/release/transfer cannot revive it. +Canonical completion can rebuild missing Markdown display through the existing +outbox. Real CLI, scale/native/imported fixtures, process loss and four-arm +read-only rehearsal cover the boundary. See [operation and compatibility](../../reference/canonical-lease-renew.md). +Executor-held external-effect locks, remaining L2/L4/L5 consumers, D2/D3 and +new-Goal defaults remain separate; this is not full L3 or T4 retirement. ### Local provider opening boundary (2026-09-13) 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 9898a6f3e3..c6295b7dcf 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 @@ -89,16 +89,20 @@ coordination 路径使用同一份语言中立的 `coordination_state_contract_v 仅将 typed read result 适配为兼容 summary。这是 contract 检查点,不是已经完成的 CLI lifecycle cutover。 -### 既有 lease 整笔事务闭合(2026-09-17) - -renew、transfer、release 共用 `coordination/task_lease_lifecycle.ts` 与既有 typed -lifecycle decision/record materializer。Python 对 promoted 命令只路由一次;local -opening handle 持有 provider 身份,含 service factory 的 PostgreSQL。独立的 renew -构造器与重复记录修改规则已移除,旧存储仅为实际 caller 保留。转交保留 Todo claim -和 scopes;到期或 owner 注销后,释放仍须匹配当前 proof。no-op 清理封存 receipt, -历史 replay 不重新授予执行权;归档后的续租/转交与不安全 generation 递增会拒绝。 -见[操作、兼容与四臂演练](../../reference/canonical-lease-renew.md)。本次闭合 L3 -既有 lease 修改,不包含 acquire/reclaim、executor fence、D2/D3 或新 Goal 默认化。 +### Lease 领取与生命周期收敛(2026-09-18) + +独立 acquire/接管和维护共用 local provider/source fence。`task_lease_acquire_decision.ts` +拥有领取准入和 materializer,legacy acquire 与 canonical 原子 Todo claim 复用; +`task_lease_state.ts` 解释完整 canonical facts,归档 holder 不再阻塞 scope。Python +只通过一次 native 请求传注册事实,不重建 canonical Todo/lease head;generation +耗尽明确拒绝。 + +Acquire receipt 本身不证明当前执行权:创建 CAS 的原样重试恢复原决定,再检查 +当前 owner/key/epoch;续约后返回当前 proof,过期/释放/转交不会复活旧执行。 +Canonical 完成可经既有 outbox 重建缺失的 Markdown 展示。真实 CLI、规模及 +native/imported fixture、进程中断和只读四臂演练覆盖此边界。见[操作与兼容](../../reference/canonical-lease-renew.md)。 +跨外部 effect 的 executor 持锁、剩余 L2/L4/L5 caller、D2/D3 和新 Goal 默认化仍 +独立验收;本批不是全部 L3 或 T4 retirement。 ### Local provider opening 边界(2026-09-13) diff --git a/docs/reference/canonical-lease-renew.md b/docs/reference/canonical-lease-renew.md index 57c2e14332..6325ecb0dc 100644 --- a/docs/reference/canonical-lease-renew.md +++ b/docs/reference/canonical-lease-renew.md @@ -1,10 +1,29 @@ -# Canonical lease renewal, transfer and release +# Canonical lease acquisition and lifecycle -An already promoted local File or SQLite Goal runs `task-lease renew`, -`transfer` and `release` through one TypeScript coordination transaction owner. -Previously only renewal used that provider path; transfer/release hit the legacy -writer fence. Unpromoted Goals retain the legacy transaction path. Selecting a -provider does not promote a Goal or grant execution authority. +An already promoted File or SQLite Goal runs `task-lease acquire`, `renew`, +`transfer` and `release` through TS-owned provider transactions. Fresh execution +and expired/ineffective-holder takeover now use the complete canonical head; +previously standalone acquire still entered the fenced legacy file writer. +Unpromoted Goals retain their legacy wire and transaction. A provider selector +does not promote a Goal or grant execution authority. + +For a new execution on an open, eligible Todo with no prior lease: + +```bash +loopx --registry registry.json task-lease acquire \ + --goal-id example-goal --todo-id todo_work --owner agent-a \ + --idempotency-key execution-a --expected-version 0 \ + --ttl-seconds 600 --write-scope 'src/**' +``` + +After expiry/release, use a new execution key and the current version from +inspect, rather than version 0. Active effective holders and overlapping scopes +on other effective leases reject acquisition. Archived, excluded, unregistered +or claim-conflicting holders do not block another eligible execution. The scan +uses all retained records, not the bounded operator display. Atomic +`todo claim + lease` shares the same facts, decision and record materializer. +Standalone acquire uses requested scopes; atomic claim retains its existing +Todo-required scope intent. Neither operation expands a permission grant. ## Operate the current lease @@ -32,6 +51,7 @@ versions, not these example numbers. | Operation | State change | Admission | | --- | --- | --- | +| Acquire / takeover | Version +1 and epoch +1; new execution identity and expiry | Open active Todo, registered eligible actor, no effective conflicting holder or overlapping execution; optional version CAS | | Renew | Version +1; owner/key/epoch/scopes unchanged; expiry is runtime clock + TTL | Active lease, current proof, registered eligible owner, active open Todo | | Transfer | Version +1 and epoch +1; replace owner/key; retain scopes; set expiry | Active lease, current proof, registered sender and eligible receiver; new execution key | | Release | Retain version/epoch; persist released status and timestamps | Current owner/key/version proof; expiry, removed registration and closed/archived Todo do not prevent cleanup | @@ -39,27 +59,39 @@ versions, not these example numbers. A missing lease at expected version 0 or an already released matching lease returns a durable `no_change` receipt. Its storage cursor can advance while lease state, timestamps and domain events remain unchanged. A wrong version or -wrong proof never becomes successful cleanup. Renew/transfer reject archived +wrong proof never becomes successful cleanup. Acquire/renew/transfer reject archived Todos and safe-integer generation exhaustion before writing. The generation check also protects the legacy path; release at that generation remains legal. ## Commit, retry and readback One provider CAS commits the lease projection, event and original receipt. -Operation identity includes operation, Goal, Todo, owner, execution key and -expected version. The immutable request digest additionally binds TTL and the +Maintenance operation identity includes operation, Goal, Todo, owner, execution +key and expected version. The immutable request digest additionally binds TTL and the transfer receiver/key. A retry with changed intent is rejected. The shipped renewal receipt schema, identity and digest encoding remain compatible. -`status=replayed` and `idempotent=true` return historical results even after a +For maintenance, `status=replayed` and `idempotent=true` return historical results even after a later renewal, transfer, release or expiry. They do not grant present execution rights or renew again. Freeze the original request after a lost/ambiguous response; recover its receipt, then inspect current state before new work. -The canonical-only lifecycle request is closed and versioned. The prior +Acquisition identity binds Goal/Todo/owner/execution key; the request digest +also binds original expected version, TTL and scope set. An exact retry of +`--expected-version 0` recovers its receipt instead of failing against the +version it created. A changed request under that identity is rejected. +Acquisition **success requires current proof**: an active, still-eligible lease +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. + +The canonical-only acquire and lifecycle requests are closed and versioned. The prior renew-only wire remains accepted for renewal only. An older runtime rejects the new schema entirely. Missing/invalid fences, changed registration facts before -renew/transfer, unavailable providers and CAS conflicts fail closed. They never +acquire/renew/transfer, unavailable providers and CAS conflicts fail closed. They never fall back to a lease file or stale/malformed Markdown. Release admission uses the existing proof without a registration snapshot. The CLI still resolves its runtime root from the registry or explicit override. @@ -76,22 +108,26 @@ matching store incarnation. There is no credential-bearing CLI option, implicit service activation or fallback. The real PostgreSQL rehearsal exercises this public native lifecycle route as well as the storage contract. -This closes existing-lease mutation under shared-authority L3 / roadmap R5. -Fresh acquisition/reclaim, executor holder/terminal fences and automatic -cross-agent result return retain their own callers and qualification. A lease +This closes standalone acquisition/takeover and existing-lease mutation under +shared-authority L3 / roadmap R5. Real CLI validation carries newly acquired +proof into canonical Todo completion and released readback. Complete/supersede +can use canonical state with a missing Markdown display, then rebuild it through +the existing projection outbox; legacy source requirements remain unchanged. +Executor holder/terminal locks across external effects and automatic cross-agent +result return retain their own callers and qualification. A lease transfer alone does not prove a completed collaboration journey. No storage format, default profile, active-Goal migration, D2 soak or D3 promotion changes. Rollback retains canonical state, receipts and the writer fence. Older code may -reject transfer/release or the new wire; plan for current lease expiry and +reject acquire/transfer/release or the new wire; plan for current lease expiry and restore compatible code. Do not remove the fence or revive stale lease files. ## Validation -The shared production-scale fixture adds an execution-handover dimension while +The shared production-scale fixture covers fresh execution, takeover and handover while retaining its mixed status, decision and historical-lease population. Every AuthorityStore conformance arm covers native/imported records, negative -admission, stale senders, response loss, CAS competition, no-op sealing and +admission, a live scope holder beyond display limits, stale senders, response loss, CAS competition, no-op sealing and historical replay. Real File/SQLite CLI and killed-process tests cover the host boundary; PostgreSQL uses an isolated real server. NoKV coverage uses its existing test transport and is not service qualification. @@ -107,47 +143,50 @@ uv run --extra test python examples/control_plane/authority-lease-lifecycle-rehe ``` Use the source-checkout Python environment and a qualified SQLite Node runtime. -The runner adds one synthetic Todo/lease only to disposable copies, compares +The runner adds two synthetic Todos and one initial lease only to disposable copies, compares all operation results and non-target records, and verifies that the live source -is unchanged. It reports the legacy historical-replay rejection as an explicit -semantic improvement, not as normalized parity. Optional `--private-diagnostics` +is unchanged. It separately reports the legacy create-CAS retry mismatch and maintenance +historical-replay rejection as semantic improvements, not normalized parity. Optional `--private-diagnostics` keeps raw failures in an owner-only file that must not be published. This does not replace [D2 capacity and continuity qualification](sqlite-authority-store.md). ## 中文操作与语义 -已 promoted 的 File/SQLite Goal,其 renew、transfer、release 现在共用一笔 TS -canonical transaction。此前只有 renew 能走 provider,另两项会被旧 writer fence -挡住。未 promoted 的默认路径保持;选择 provider 本身不构成 promotion。 - -按上面的命令先 inspect,再使用实际 owner/key/version。续租只增加 version;转交 -增加 version 和 epoch,必须换 execution key,并检查双方注册及接收方资格;两者 -保留 write scopes。转交不改变 Todo claim,不覆盖排除规则。释放保留版本与 epoch, -记录 released 状态和时间;只凭当前 owner/key/version 清理,允许已到期、已注销 -owner 和已关闭/归档 Todo。版本不匹配仍拒绝。 - -缺失 lease 的 version 0 清理及已释放 lease 的匹配清理,会封存 `no_change` receipt: -存储 cursor 可以前进,但 lease、时间和 domain event 不变。归档 Todo 禁止再次 -续租/转交;version 或 epoch 无法安全递增时,两种存储路径都拒绝,仍允许释放。 - -lease/event/receipt 原子提交。重试身份绑定 operation、Goal、Todo、owner、execution -key 和 expected version;不可变 intent 另绑定 TTL 及接收者/key。改 intent 重试会 -被拒绝。既有 renew receipt 和 digest 保持兼容。历史 replay 可以跨后续转交、续租、 -释放和到期,返回原结果而不改变当前状态,也不重新授予执行权。响应不确定时冻结 -原请求恢复,再 inspect 当前状态。 - -新 wire 只接受既有 lease 的三类操作,旧 renew wire 仍只接受 renew;旧 runtime -会拒绝新 schema。fence、provider、提交前注册源变化或 CAS 失败都不回退。返回值 -包含 provider/revision/cursor 及版本冲突细节,不创建第二份 lease/shadow 状态, -不为 lease-only 修改重写 Markdown。释放 admission 不需要注册快照;CLI 仍从 registry -或显式 override 解析 runtime root。 - -PostgreSQL 复用同一个 opening handle,但必须由 service owner 提供已有的 scoped -factory 并核对 incarnation;没有新增凭据参数或自动启用服务。真实四臂演练验证 -相同公共 native 入口;NoKV 仅经过已有测试 transport,不代表其服务已合格。 - -这是 L3/R5 的既有 lease 修改闭环,fresh acquire/reclaim、executor fence 及自动 -结果返回仍须各自验收。没有改存储格式、默认 provider、活动 Goal、D2 soak 或 D3 -晋升。回滚保留 canonical state/receipt/fence,考虑旧版本拒绝操作与 lease 到期, -恢复兼容代码;不能删 fence 或复活旧 lease 文件。上面的只读快照演练只修改隔离 -副本,输出摘要,核对源与无关记录不变,不能替代完整持久性资格。 +已 promoted 的 File/SQLite Goal,其 acquire、renew、transfer、release 现在使用 +同一 provider opening/source fence 与 TS 规则;新领取和失效持有者接管不再进入 +旧文件 writer。旧 wire 与未 promoted 路径保留,选择 provider 不构成 promotion。 + +新执行按上面的 acquire 命令领取;没有旧 lease 时 expected-version 为 0,否则 +用 inspect 的当前版本及新的 execution key。领取/接管同时增加 version 和 epoch, +续约只增 version,转交同时增二者并换 key;释放保留 generation。转交不改变 Todo +claim、不覆盖 exclusion 或扩大 scope;释放只凭匹配 proof,允许到期或注销 owner +清理。所有需递增的入口都拒绝安全整数耗尽,仍允许释放。 + +完整 canonical Todo/lease 集合决定 scope 冲突,不能只看 UI 页面。归档、排除、 +注销或与当前 claim 冲突的 holder 不阻挡新的合格执行。独立 acquire 和原子的 +Todo claim + lease 共用状态解释、准入和 materializer;前者使用请求 scopes,后者 +仍使用 Todo required scopes。scope 是执行冲突声明,不扩张权限。 + +一笔 CAS 保存 lease/event/原 receipt。维护操作的身份与既有 renew digest 保持 +兼容,历史 replay 可跨后续修改,但不授予当前执行权。Acquire 的身份绑定 +Goal/Todo/owner/key,digest 另绑定原 expected version、TTL、scope set;创建时 +expected-version 0 的原样重试可恢复回执,改变参数会拒绝。 + +Acquire 的成功还必须核对当前有效 owner/key/epoch 和资格。同一执行续约后, +重试返回 `lease` 中的当前版本/到期时间,以及 `original_receipt` 中不可变的原始 +决定;`current_provider_revision/current_cursor` 标识当前读回。已转交、到期或释放 +的旧执行不能凭 receipt 复活。Todo claim 和维护 receipt 仍是历史语义,不能将其 +当成新的 acquire 响应。 + +canonical acquire 与 lifecycle 各有封闭 wire,旧 renew wire 只接受 renew;旧 +runtime 不识别新 acquire schema。fence、provider、注册源变化和 CAS 错误不回退 +旧文件。lease-only 命令不生成第二份 shadow 或重写 Markdown。真实 CLI 验证包含 +新领取→续约→释放→新执行→完成;complete/supersede 可在 Markdown 展示丢失时 +读取 canonical state,完成后经原 projection outbox 重建,旧路径仍要求源文件。 + +PostgreSQL 使用已有 service-owned scoped factory 和 incarnation 检查,无新凭据 +参数或自动启用。四臂演练使用相同公共 native 入口;NoKV 仍只经过测试 transport。 +L3/R5 的独立领取/接管及维护由此可用,跨外部 effect 的 executor holder/terminal +锁和自动结果返回仍需各自验收。D2 soak、D3、默认 profile、活动 Goal 迁移和跨主机 +部署没有改变。回滚保留 canonical state/receipt/fence 并恢复兼容代码,不得复活旧 +lease 文件。演练仅修改隔离副本,核对源与无关记录不变,不能代替完整持久性资格。 diff --git a/examples/control_plane/authority-lease-lifecycle-rehearsal.py b/examples/control_plane/authority-lease-lifecycle-rehearsal.py index 4374fcefbb..0d6f5d34a0 100644 --- a/examples/control_plane/authority-lease-lifecycle-rehearsal.py +++ b/examples/control_plane/authority-lease-lifecycle-rehearsal.py @@ -35,6 +35,8 @@ let raw=''; for await (const chunk of process.stdin) raw+=chunk; const input=JSON.parse(raw); const moduleAt=(root,path)=>import(pathToFileURL(join(root,'loopx/control_plane',path)).href); +const {executeTaskLeaseAcquire: currentAcquire}=await moduleAt(input.repo,'work_items/task_lease_acquire.ts'); +const {executeTaskLeaseAcquire: legacyAcquire}=await moduleAt(input.baseline_repo,'work_items/task_lease_acquire.ts'); const {executeTaskLeaseLifecycle: current}=await moduleAt(input.repo,'work_items/task_lease_lifecycle.ts'); const {executeTaskLeaseLifecycle: legacy}=await moduleAt(input.baseline_repo,'work_items/task_lease_lifecycle.ts'); const {FileAuthorityStore}=await moduleAt(input.repo,'coordination/file_authority_store.ts'); @@ -45,12 +47,14 @@ const {canonicalAuthorityBytes,canonicalAuthoritySha256}=await moduleAt(input.repo,'coordination/authority_store_codec.ts'); const {engageLegacyCoordinationWriterFence}=await moduleAt(input.repo,'coordination/legacy_writer_fence.ts'); const digest=value=>canonicalAuthoritySha256(value); -const goal=input.goal_id, target='todo_lifecycle_rehearsal'; +const goal=input.goal_id, target='todo_lifecycle_rehearsal', acquisition='todo_acquire_rehearsal'; +assert(!input.projection.todos.some(t=>t.todo_id===acquisition)); assert(!input.projection.todos.some(t=>t.todo_id===target)); const initial=structuredClone(input.projection); initial.todos.push({schema_version:'todo_item_v0',todo_id:target,role:'agent',status:'open',done:false, text:'Isolated lease lifecycle rehearsal',archive_state:'active',source_section:'Agent Todo', claimed_by:null,excluded_agents:[],task_class:'advancement_task'}); +initial.todos.push({...initial.todos.at(-1),todo_id:acquisition,text:'Isolated acquisition and takeover rehearsal'}); initial.todos.sort((a,b)=>a.todo_idb.todo_id?1:0); initial.todo_read_model=coordinationTodoReadModel(initial.todos,initial.todo_read_model.schema_version); initial.handoff_mode='hard_lease'; @@ -100,6 +104,28 @@ state:'engaged',goal_id:goal,fence_id:'rehearsal',source_version:'state:1',source_projection_sha256:digest(initial), expected_shadow_provider_revision:h.provider_revision}})).status,'applied'); } + const acquireBase={schema_version:store?'loopx_canonical_task_lease_acquire_request_v0':'loopx_task_lease_acquire_native_v0', + runtime_root:runtime,goal_id:goal,todo_id:acquisition,owner:'agent-a',idempotency_key:'acquire-a', + expected_version:0,ttl_seconds:600,write_scopes:['isolated-acquisition/**'],authority}; + const invokeAcquire=(request,now)=>(store?currentAcquire:legacyAcquire)(request,{...dependencies,now:()=>new Date(now)}); + const acquired=await invokeAcquire(acquireBase,'2026-09-13T10:05:00Z'); + assert.equal(acquired.ok,true,`${arm} acquire: ${acquired.error_code}`); + assert.deepEqual(acquired.lease,{schema_version:'task_lease_v0',goal_id:goal,todo_id:acquisition, + owner:'agent-a',idempotency_key:'acquire-a',version:1,lease_epoch:1,status:'active', + write_scopes:['isolated-acquisition/**'],acquire_ttl_seconds:600,acquired_at:'2026-09-13T10:05:00Z', + updated_at:'2026-09-13T10:05:00Z',expires_at:'2026-09-13T10:15:00Z'}); + const acquireReplay=await invokeAcquire(acquireBase,'2026-09-13T10:06:00Z'); + if(store) {assert.equal(acquireReplay.status,'replayed');assert.deepEqual(acquireReplay.lease,acquired.lease);} + else assert.equal(acquireReplay.error_code,'version_mismatch'); + const takeover=await invokeAcquire({...acquireBase,owner:'agent-b',idempotency_key:'acquire-b',expected_version:1},'2026-09-14T10:05:00Z'); + assert.equal(takeover.ok,true,`${arm} takeover: ${takeover.error_code}`); + assert.deepEqual(takeover.lease,{...acquired.lease,owner:'agent-b',idempotency_key:'acquire-b',version:2,lease_epoch:2, + acquired_at:'2026-09-14T10:05:00Z',updated_at:'2026-09-14T10:05:00Z',expires_at:'2026-09-14T10:15:00Z'}); + if(store) { + const staleAcquire=await invokeAcquire(acquireBase,'2026-09-14T10:06:00Z'); + assert.equal(staleAcquire.error_code,'idempotency_key_reuse'); + assert.equal((await import('node:fs')).existsSync(join(leaseDir,`${acquisition}.json`)),false); + } const base={schema_version:store?'loopx_canonical_task_lease_lifecycle_request_v0':'loopx_task_lease_lifecycle_native_v0', runtime_root:runtime,goal_id:goal,todo_id:target,owner:'agent-a',idempotency_key:'rehearsal-a', expected_version:3,ttl_seconds:600,authority,current_time:'2026-09-13T10:05:00Z'}; @@ -125,24 +151,24 @@ if (store) { assert.equal(replay.status,'replayed');assert.deepEqual(replay.original_receipt,results[0].original_receipt); assert.deepEqual(replay.lease,results[0].lease); - const final=await store.loadAuthority();assert.equal(final.status,'loaded');assert.equal(final.cursor,'4'); + const final=await store.loadAuthority();assert.equal(final.status,'loaded');assert.equal(final.cursor,'6'); assert.deepEqual(final.head.todos,initial.todos); - assert.deepEqual(final.head.leases.filter(l=>l.todo_id!==target),initial.leases.filter(l=>l.todo_id!==target)); + assert.deepEqual(final.head.leases.filter(l=>l.todo_id!==target&&l.todo_id!==acquisition),initial.leases.filter(l=>l.todo_id!==target)); assert.deepEqual(await readFile(join(leaseDir,`${target}.json`)),legacyBefore); finalHeads[arm]=final.head; - report[arm]={passed:true,commits:4,historical_replay:'original_receipt'}; + report[arm]={passed:true,commits:6,historical_replay:'original_receipt',acquire_retry:'current_proof',retired_acquire_rejected:true}; } else { assert.equal(replay.error_code,'lifecycle_receipt_state_mismatch'); report[arm]={passed:true,historical_replay:'baseline_rejected_after_later_mutation'}; } - observations[arm]=results.map(r=>({lease:r.lease,handoff_mode:r.handoff_mode})); + observations[arm]=[acquired,takeover,...results].map(r=>({lease:r.lease})); } for (const arm of ['file','sqlite','postgresql']) assert.deepEqual(observations[arm],observations.legacy); assert.deepEqual(finalHeads.file,finalHeads.sqlite); assert.deepEqual(finalHeads.file,finalHeads.postgresql); process.stdout.write(JSON.stringify({schema_version:'authority_lease_lifecycle_rehearsal_v0', initial_todos:initial.todos.length,initial_leases:initial.leases.length,fixture_sha256:digest(initial), observation_sha256:digest(observations.file),provider_head_sha256:digest(finalHeads.file),arms:report, - intentional_delta:'canonical_historical_replay_survives_later_mutations',non_target_records_unchanged:true})); + intentional_delta:'canonical_receipt_recovery_and_current_acquire_proof;_explicit_create_CAS_retry_no_longer_mismatches',non_target_records_unchanged:true})); } finally { for (const table of ['authority_receipts','authority_events','authority_commits','authority_heads']) await pool.query(`DELETE FROM loopx_control_plane.${table} WHERE tenant_id=$1`,[tenant]); diff --git a/examples/shared-goal-authority-e2e/correctness.md b/examples/shared-goal-authority-e2e/correctness.md index ffd642c65b..5e835e9cd7 100644 --- a/examples/shared-goal-authority-e2e/correctness.md +++ b/examples/shared-goal-authority-e2e/correctness.md @@ -177,8 +177,10 @@ and `tests/control_plane/test_shadow_fence_caller_parity_e2e.py`; the `baseline` entries of that fixture document earlier revisions and are never executed. -Promoted CLI renew, transfer and release rows now assert provider commits and -independent lease readback; the direct legacy-wire rows remain fenced. The CLI +Promoted CLI acquire, renew, transfer and release rows assert provider commits and +independent lease readback; fresh acquisition C is separate from keyless A and +existing-lease B so one newly working caller does not invalidate another row's intended negative precondition. Acquisition replay verifies current proof. +The direct legacy-wire rows remain fenced. The CLI scenario carries the validated current owner/key/version after renewal or transfer, rejects the superseded completion proof, and previews with the current proof before release. Release follows the active-lease/quiescence checks and diff --git a/examples/shared-goal-authority-e2e/mutants.py b/examples/shared-goal-authority-e2e/mutants.py index 4ed37338b7..3830ef3f81 100644 --- a/examples/shared-goal-authority-e2e/mutants.py +++ b/examples/shared-goal-authority-e2e/mutants.py @@ -328,11 +328,13 @@ def apply(source: str) -> str: ), "tests/control_plane/test_shadow_drain_e2e.py::test_public_mutation_has_no_second_snapshot_mirror")) # Fenced sibling-caller parity: each mutant is caught by exactly one parity row. +# Promoted CLI acquisition now succeeds through canonical authority. Probe the +# retained native legacy entrypoint for its fence diagnostic/envelope contract. CASES.append(Case("native_fence_remediation_truncated", ( (COORDINATION + "legacy_writer_fence.ts", replacement( 'export const LEGACY_WRITER_FENCED_REMEDIATION =\n "legacy coordination writer is fenced; use the promoted canonical authority ({authority_mode}) for goal {goal_id}; fence {fence_id}; the primary record was not changed";', 'export const LEGACY_WRITER_FENCED_REMEDIATION =\n "legacy coordination writer is fenced";')), -), "tests/control_plane/test_shadow_fence_caller_parity_e2e.py::test_fence_caller_parity[cli-task_lease_acquire-engaged]")) +), "tests/control_plane_ts/legacy_writer_fence_caller_parity.test.ts", pattern="^fence parity: ts-acquire-engaged$")) CASES.append(Case("python_fence_remediation_truncated", ( (COORDINATION + "legacy_writer_fence.py", replacement( 'LEGACY_WRITER_FENCED_REMEDIATION = (\n "legacy coordination writer is fenced; use the promoted canonical authority "\n "({authority_mode}) for goal {goal_id}; fence {fence_id}; "\n "the primary record was not changed"\n)', @@ -342,7 +344,7 @@ def apply(source: str) -> str: (COORDINATION + "legacy_writer_fence.ts", replacement( " this.payload = { write_check: writeCheck };", " this.payload = writeCheck;")), -), "tests/control_plane/test_shadow_fence_caller_parity_e2e.py::test_fence_caller_parity[cli-task_lease_acquire-engaged]")) +), "tests/control_plane_ts/legacy_writer_fence_caller_parity.test.ts", pattern="^fence parity: ts-acquire-engaged$")) CASES.append(Case("fence_acquire_receipt_fabricated", ( ("loopx/control_plane/work_items/task_lease_acquire.ts", replacement( ' return { code: error.code, message: error.message, payload: error.payload, stage: "validation", kind: "permission_denied" };', diff --git a/loopx/control_plane/coordination/coordination_state_contract.generated.ts b/loopx/control_plane/coordination/coordination_state_contract.generated.ts index 5874791688..b3c3f8348c 100644 --- a/loopx/control_plane/coordination/coordination_state_contract.generated.ts +++ b/loopx/control_plane/coordination/coordination_state_contract.generated.ts @@ -78,6 +78,7 @@ export const DELIVERY_WORKSPACE_SNAPSHOT_REQUEST_SCHEMA = "loopx_delivery_worksp export const DELIVERY_WORKSPACE_SNAPSHOT_RESULT_SCHEMA = "loopx_delivery_workspace_result_v0"; export const TASK_LEASE_ACQUIRE_REQUEST_SCHEMA = "loopx_task_lease_acquire_native_v0"; +export const TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA = "loopx_canonical_task_lease_acquire_request_v0"; export const TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA = "loopx_task_lease_lifecycle_native_v0"; export const TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA = "loopx_canonical_task_lease_renew_request_v0"; export const TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA = "loopx_canonical_task_lease_lifecycle_request_v0"; @@ -295,6 +296,7 @@ export const COORDINATION_STATE_CONTRACT = deepFreeze({ }, "task_lease_protocol": { "acquire_request_schema": TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, + "canonical_acquire_request_schema": TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA, "lifecycle_request_schema": TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, "canonical_renew_request_schema": TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA, "canonical_lifecycle_request_schema": TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA diff --git a/loopx/control_plane/coordination/coordination_state_contract_generated.py b/loopx/control_plane/coordination/coordination_state_contract_generated.py index a63534552b..b3ead85018 100644 --- a/loopx/control_plane/coordination/coordination_state_contract_generated.py +++ b/loopx/control_plane/coordination/coordination_state_contract_generated.py @@ -164,6 +164,7 @@ def _freeze(value: Any) -> Any: 'request_schema': 'loopx_delivery_workspace_request_v0', 'result_schema': 'loopx_delivery_workspace_result_v0'}, 'task_lease_protocol': {'acquire_request_schema': 'loopx_task_lease_acquire_native_v0', + 'canonical_acquire_request_schema': 'loopx_canonical_task_lease_acquire_request_v0', 'lifecycle_request_schema': 'loopx_task_lease_lifecycle_native_v0', 'canonical_renew_request_schema': 'loopx_canonical_task_lease_renew_request_v0', 'canonical_lifecycle_request_schema': 'loopx_canonical_task_lease_lifecycle_request_v0'}, @@ -262,6 +263,7 @@ def _freeze(value: Any) -> Any: DELIVERY_WORKSPACE_SNAPSHOT_RESULT_SCHEMA: Final[str] = 'loopx_delivery_workspace_result_v0' TASK_LEASE_ACQUIRE_REQUEST_SCHEMA: Final[str] = 'loopx_task_lease_acquire_native_v0' +TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA: Final[str] = 'loopx_canonical_task_lease_acquire_request_v0' TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA: Final[str] = 'loopx_task_lease_lifecycle_native_v0' TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA: Final[str] = 'loopx_canonical_task_lease_renew_request_v0' TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA: Final[str] = 'loopx_canonical_task_lease_lifecycle_request_v0' diff --git a/loopx/control_plane/coordination/coordination_state_contract_v0.json b/loopx/control_plane/coordination/coordination_state_contract_v0.json index f1d1cc7688..c80f748a3c 100644 --- a/loopx/control_plane/coordination/coordination_state_contract_v0.json +++ b/loopx/control_plane/coordination/coordination_state_contract_v0.json @@ -174,6 +174,7 @@ }, "task_lease_protocol": { "acquire_request_schema": "loopx_task_lease_acquire_native_v0", + "canonical_acquire_request_schema": "loopx_canonical_task_lease_acquire_request_v0", "lifecycle_request_schema": "loopx_task_lease_lifecycle_native_v0", "canonical_renew_request_schema": "loopx_canonical_task_lease_renew_request_v0", "canonical_lifecycle_request_schema": "loopx_canonical_task_lease_lifecycle_request_v0" diff --git a/loopx/control_plane/coordination/task_lease_acquire.ts b/loopx/control_plane/coordination/task_lease_acquire.ts new file mode 100644 index 0000000000..5c14a89e54 --- /dev/null +++ b/loopx/control_plane/coordination/task_lease_acquire.ts @@ -0,0 +1,128 @@ +/** A lease acquisition is one full-head admission/CAS with a retained receipt. + * Unlike historical maintenance replay, success here must supply current proof. */ +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 {decideTaskLeaseAcquire, materializeTaskLeaseAcquire} from "../work_items/task_lease_acquire_decision.ts"; +import {leaseOwnerRejection} from "../work_items/task_lease_eligibility.ts"; +import {normalizeGoalId, normalizeTodoId, normalizeOwner, normalizeIdempotencyKey, + normalizeWriteScopes, normalizeTtl, leaseEpoch, leaseVersion, leaseIsActive, + TaskLeaseAcquireError} from "../work_items/task_lease_acquire.ts"; + +export interface CanonicalTaskLeaseAcquireInput { + goal_id: string; todo_id: string; owner: string; idempotency_key: string; + expected_version: number | null; ttl_seconds: number | null; + write_scopes: readonly string[]; registered_agents: readonly string[]; now: Date; +} + +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 = {}) => + ({schema_version: schema, status: "failed", changed: false, reason_code: code, reason, failure_stage: "validation", ...detail}); + let input: CanonicalTaskLeaseAcquireInput & {ttl_seconds: number}; + try { + input = {...raw, goal_id: normalizeGoalId(raw.goal_id), todo_id: normalizeTodoId(raw.todo_id), + owner: normalizeOwner(raw.owner), idempotency_key: normalizeIdempotencyKey(raw.idempotency_key), + write_scopes: normalizeWriteScopes(raw.write_scopes), ttl_seconds: normalizeTtl(raw.ttl_seconds), + registered_agents: raw.registered_agents.map(normalizeOwner)}; + if (input.expected_version !== null && (!Number.isSafeInteger(input.expected_version) || input.expected_version < 0)) { + return failed("invalid_expected_version", "expected lease version must be a non-negative safe integer"); + } + if (!(input.now instanceof Date) || !Number.isFinite(input.now.valueOf())) return failed("invalid_clock", "lease acquire requires a valid clock"); + } catch (error) { + return failed(error instanceof TaskLeaseAcquireError ? error.code : "invalid_canonical_acquire_request", + error instanceof Error ? error.message : "invalid canonical lease acquire"); + } + const identityFields = {goal_id: input.goal_id, todo_id: input.todo_id, owner: input.owner, idempotency_key: input.idempotency_key}; + const identity = {schema_version: "loopx_canonical_task_lease_acquire_receipt_v0", + operation_id: `lease-acquire:${canonicalAuthoritySha256(identityFields)}`, goal_id: input.goal_id, + request_sha256: canonicalAuthoritySha256({...identityFields, expected_version: input.expected_version, + ttl_seconds: input.ttl_seconds, write_scopes: [...input.write_scopes].sort()})}; + const receipt = new CoordinationCommandReceipt({result_schema: schema, identity, failure: failed, + decode(original) { + const payload = commandReceiptResult(original); + const lease = canonicalTaskLease(canonicalAuthorityObject(payload.fields.lease, "acquire receipt lease"), input.goal_id, input.todo_id); + const scopes = [...(lease.write_scopes ?? []) as string[]].sort(); + if (JSON.stringify(scopes) !== JSON.stringify([...input.write_scopes].sort()) || + (lease.acquire_ttl_seconds != null && lease.acquire_ttl_seconds !== input.ttl_seconds) || + (payload.changed && input.expected_version !== null && leaseVersion(lease) !== input.expected_version + 1)) { + throw new AuthorityStoreProtocolError("acquire receipt does not match its original parameters"); + } + if (lease.owner !== input.owner || lease.idempotency_key !== input.idempotency_key || lease.status !== "active" || + leaseVersion(lease) < 1 || leaseEpoch(lease) < 1 || !leaseIsActive(lease, new Date(String(lease.acquired_at)))) { + throw new AuthorityStoreProtocolError("acquire receipt does not match its execution identity"); + } + 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); + } + // 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}; + }; + + let committed = false; + try { + const replay = await receipt.read(store); + if (replay !== null) return await currentProof(replay); + const {head, facts, mode} = await readFacts(); + 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, + ...facts, command: input}); + if (decision.outcome === "rejected" || decision.outcome === "conflict") { + return failed(decision.code, `canonical task lease acquire rejected: ${decision.code}`, { + handoff_mode: mode, expected_version: input.expected_version, actual_version: leaseVersion(facts.current), + ...(facts.todo ? {todo_status: facts.todo.status, claimed_by: facts.todo.claimed_by, excluded_agents: [...facts.todo.excluded_agents]} : {}), + ...(decision.conflict_indexes.length ? {conflicts: decision.conflict_indexes.map(i => facts.other_leases[i])} : {})}); + } + const changed = decision.outcome === "apply"; + const lease = changed ? materializeTaskLeaseAcquire(input, input, decision, input.now) : facts.current; + if (!lease) throw new AuthorityStoreProtocolError("accepted acquire lacks a lease"); + const commit = changed ? prepareCoordinationProjectionCommit({goal_id: input.goal_id, + operation_id: identity.operation_id, expected_provider_revision: head.provider_revision, projection: head.head, + mutations: [{kind: "lease_upsert", lease}]}) : {operation_id: identity.operation_id, + expected_provider_revision: head.provider_revision, next_projection: head.head, events: [], receipts: []}; + commit.receipts = [{...identity, result: {changed, lease, handoff_mode: mode}}]; + await beforeCommit?.(lease); + committed = true; + const result = await receipt.commit(store, commit); + return await currentProof(result); + } catch (error) { + if (committed) return {...failed("canonical_acquire_recovery_required", "acquire may be durable; retry the same request"), + status: "ambiguous", failure_stage: "durable_writeback", recovery: {operation_id: identity.operation_id, retry_with_same_operation_id: true}}; + return failed(error instanceof TaskLeaseAcquireError ? error.code : "invalid_canonical_acquire_state", + error instanceof Error ? error.message : "canonical acquisition state is invalid"); + } +} diff --git a/loopx/control_plane/coordination/task_lease_lifecycle.ts b/loopx/control_plane/coordination/task_lease_lifecycle.ts index 4218ae4a46..2db75321d2 100644 --- a/loopx/control_plane/coordination/task_lease_lifecycle.ts +++ b/loopx/control_plane/coordination/task_lease_lifecycle.ts @@ -1,3 +1,4 @@ +import {canonicalTaskLease, canonicalLeaseTodoFact} from "./task_lease_state.ts"; /** Lease mutations share one canonical revision, decision and durable receipt. * Providers own persistence only; replay never grants current execution rights. */ import type {JsonObject} from "../effect_program.ts"; @@ -35,24 +36,6 @@ export interface CanonicalTaskLeaseLifecycleInput { now: Date; } -function leaseRecord(value: JsonObject, input: CanonicalTaskLeaseLifecycleInput): LeaseRecord { - if ((value.schema_version !== undefined && value.schema_version !== "task_lease_v0") || - (value.goal_id !== undefined && value.goal_id !== input.goal_id) || value.todo_id !== input.todo_id || - (value.status !== "active" && value.status !== "released")) { - throw new AuthorityStoreProtocolError("canonical lease identity or schema is invalid"); - } - if (typeof value.owner !== "string" || typeof value.idempotency_key !== "string" || - normalizeOwner(value.owner) !== value.owner || normalizeIdempotencyKey(value.idempotency_key) !== value.idempotency_key) { - throw new AuthorityStoreProtocolError("canonical lease owner and execution key must be normalized strings"); - } - leaseVersion(value); leaseEpoch(value); - if (value.write_scopes !== undefined && (!Array.isArray(value.write_scopes) || - value.write_scopes.some(scope => typeof scope !== "string"))) { - throw new AuthorityStoreProtocolError("canonical lease write_scopes must be strings"); - } - return value; -} - export async function executeCanonicalTaskLeaseLifecycle(store: AuthorityStore, raw: CanonicalTaskLeaseLifecycleInput, beforeCommit?: (lease: JsonObject | null) => Promise): Promise { const contract = CONTRACTS[raw.operation]; @@ -96,7 +79,7 @@ export async function executeCanonicalTaskLeaseLifecycle(store: AuthorityStore, } let lease: LeaseRecord; try { - lease = leaseRecord(canonicalAuthorityObject(fields.lease, "lease receipt record"), input); + lease = canonicalTaskLease(canonicalAuthorityObject(fields.lease, "lease receipt record"), input.goal_id, input.todo_id); leaseIsActive(lease, new Date(0)); // Validate time syntax, not present-day authority. } catch (error) { throw new AuthorityStoreProtocolError(error instanceof Error ? error.message : "invalid lease receipt record"); @@ -120,13 +103,12 @@ export async function executeCanonicalTaskLeaseLifecycle(store: AuthorityStore, const index = indexCoordinationProjection(head.head, input.goal_id); validateCoordinationTodoReadModel(head.head, input.goal_id); const todo = index.todos.get(input.todo_id), rawLease = index.leases.get(input.todo_id); - const lease = rawLease ? leaseRecord(rawLease, input) : null; + const lease = rawLease ? canonicalTaskLease(rawLease, input.goal_id, input.todo_id) : null; const excluded = todo?.excluded_agents ?? []; if (!Array.isArray(excluded) || excluded.some(value => typeof value !== "string")) return failed("invalid_coordination_projection", "Todo exclusions must be strings"); const mode = requireStringLiteral(head.head.handoff_mode ?? "legacy", HANDOFF_MODES, "canonical handoff_mode"); const decision = decideTaskLeaseLifecycle({handoff_mode: mode, registered_agents: input.registered_agents, - todo: todo && todo.archive_state === "active" ? {todo_id: input.todo_id, status: String(todo.status), - claimed_by: normalizeAgent(todo.claimed_by), excluded_agents: excluded as string[]} : null, + todo: canonicalLeaseTodoFact(todo), lease: lease ? {present: true, active: leaseIsActive(lease, input.now), status: String(lease.status), owner: normalizeOwner(lease.owner), idempotency_key: normalizeIdempotencyKey(lease.idempotency_key), version: leaseVersion(lease), lease_epoch: leaseEpoch(lease), write_scopes: (lease.write_scopes ?? []) as string[], acquire_ttl_seconds: null} : null, diff --git a/loopx/control_plane/coordination/task_lease_state.ts b/loopx/control_plane/coordination/task_lease_state.ts new file mode 100644 index 0000000000..ca1225db36 --- /dev/null +++ b/loopx/control_plane/coordination/task_lease_state.ts @@ -0,0 +1,58 @@ +/** Full canonical facts shared by lease acquisition, maintenance and Todo claim. */ +import type {JsonObject} from "../effect_program.ts"; +import {AuthorityStoreProtocolError} from "./authority_store_codec.ts"; +import {indexCoordinationProjection} from "./coordination_projection.ts"; +import {normalizeTodoAgent} from "./todo_agents.ts"; +import {leaseOwnerRejection} from "../work_items/task_lease_eligibility.ts"; +import {leaseVersion, leaseEpoch, leaseInteger, leaseIsActive, normalizeOwner, + normalizeIdempotencyKey, type LeaseRecord, type TodoFact} from "../work_items/task_lease_acquire.ts"; +import type {AcquireDecisionInput} from "../work_items/task_lease_acquire_decision.ts"; + +export function canonicalTaskLease(value: JsonObject, goalId: string, todoId: string): LeaseRecord { + if ((value.schema_version !== undefined && value.schema_version !== "task_lease_v0") || + (value.goal_id !== undefined && value.goal_id !== goalId) || value.todo_id !== todoId || + (value.status !== "active" && value.status !== "released")) { + throw new AuthorityStoreProtocolError("canonical lease identity or schema is invalid"); + } + if (typeof value.owner !== "string" || typeof value.idempotency_key !== "string" || + normalizeOwner(value.owner) !== value.owner || normalizeIdempotencyKey(value.idempotency_key) !== value.idempotency_key) { + throw new AuthorityStoreProtocolError("canonical lease owner and execution key must be normalized strings"); + } + leaseVersion(value); leaseEpoch(value); + if (value.write_scopes !== undefined && (!Array.isArray(value.write_scopes) || + value.write_scopes.some(scope => typeof scope !== "string"))) { + throw new AuthorityStoreProtocolError("canonical lease write_scopes must be strings"); + } + return value; +} + +export function canonicalLeaseTodoFact(todo: JsonObject | undefined): TodoFact | null { + if (!todo || todo.archive_state !== "active") return null; + const excluded = todo.excluded_agents ?? []; + if (!Array.isArray(excluded)) throw new AuthorityStoreProtocolError("Todo exclusions must be an array"); + return {todo_id: String(todo.todo_id), status: String(todo.status), + claimed_by: todo.claimed_by == null ? null : normalizeTodoAgent(todo.claimed_by, "todo.claimed_by"), + excluded_agents: excluded.map(value => normalizeTodoAgent(value, "todo.excluded_agents"))}; +} + +/** Every retained lease is paired with its current Todo, never a display page. */ +export function canonicalTaskLeaseAcquireFacts(index: ReturnType, + goalId: string, todoId: string, registered: readonly string[], now: Date): + Pick & {current: LeaseRecord | null} { + const raw = index.leases.get(todoId); + const current = raw ? canonicalTaskLease(raw, goalId, todoId) : null; + const todo = canonicalLeaseTodoFact(index.todos.get(todoId)); + const lease = current === null ? null : {present: true, active: leaseIsActive(current, now), + status: String(current.status), owner: String(current.owner), idempotency_key: String(current.idempotency_key), + version: leaseVersion(current), lease_epoch: leaseEpoch(current), + write_scopes: (current.write_scopes ?? []) as string[], acquire_ttl_seconds: leaseInteger(current, "acquire_ttl_seconds")}; + const other_leases = [...index.leases].flatMap(([id, rawLease]) => { + if (id === todoId) return []; + const candidate = canonicalTaskLease(rawLease, goalId, id); + const active = leaseIsActive(candidate, now); + return [{todo_id: id, active, + effective: active && leaseOwnerRejection(canonicalLeaseTodoFact(index.todos.get(id)), String(candidate.owner), registered) === null, + write_scopes: (candidate.write_scopes ?? []) as string[]}]; + }); + return {todo, lease, other_leases, current}; +} diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index 552d93dcb2..482adedb17 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -1,4 +1,5 @@ -import {leaseOwnerRejection as ownerRejection} from "../work_items/task_lease_eligibility.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"; import type { AuthorityStore, AuthorityStoreCommit } from "./authority_store.ts"; import { @@ -17,18 +18,12 @@ import { type CoordinationProjectionMutation, } from "./coordination_projection.ts"; import { - evaluateTaskLeaseAcquireDecision, - leaseEpoch, leaseInteger, leaseIsActive, normalizeAgent, normalizeIdempotencyKey, normalizeTtl, normalizeWriteScopes, - TASK_LEASE_SCHEMA_VERSION, - utcIsoformat, - type LeaseRecord, - type TodoFact, } from "../work_items/task_lease_acquire.ts"; export const COORDINATION_TODO_CLAIM_RESULT_SCHEMA = @@ -303,7 +298,7 @@ function normalizeLeaseRequest(value: unknown): CoordinationTodoClaimLeaseReques }; } -function todoLeaseFact(todo: JsonObject): TodoFact { +function claimWriteScopes(todo: JsonObject): string[] { const requiredWriteScopes = todo.required_write_scopes ?? []; if (!Array.isArray(requiredWriteScopes) || requiredWriteScopes.some((scope) => typeof scope !== "string")) { @@ -317,23 +312,7 @@ function todoLeaseFact(todo: JsonObject): TodoFact { "todo.required_write_scopes contains an invalid or duplicate scope", ); } - return { - todo_id: String(todo.todo_id), - status: typeof todo.status === "string" ? todo.status : "", - claimed_by: normalizeAgent(todo.claimed_by), - excluded_agents: normalizeExcludedAgents(todo.excluded_agents), - role: typeof todo.role === "string" ? todo.role : undefined, - task_class: typeof todo.task_class === "string" ? todo.task_class : null, - bound_agent: normalizeAgent(todo.bound_agent), - blocks_agent: normalizeAgent(todo.blocks_agent), - }; -} - -function leaseDecisionInteger(value: unknown, label: string): number { - if (!Number.isSafeInteger(value) || Number(value) < 0) { - throw new AuthorityStoreProtocolError(`${label} must be a non-negative safe integer`); - } - return Number(value); + return writeScopes; } function activeLeaseForOwner( @@ -510,50 +489,14 @@ export async function executeCoordinationTodoClaim( try { const currentLease = projection.leases.get(input.todo_id); if (handoffMode === "hard_lease" && leaseRequest !== null) { - const todoFact = todoLeaseFact(todo); - const currentActive = currentLease !== undefined && leaseIsActive(currentLease, input.now); - const otherLeases = projection.lease_todo_ids.flatMap((todoId) => { - if (todoId === input.todo_id) return []; - const candidate = projection.leases.get(todoId)!; - const active = leaseIsActive(candidate, input.now); - const otherTodo = projection.todos.get(todoId); - return [{ - todo_id: todoId, - active, - effective: active && otherTodo !== undefined && ownerRejection( - todoLeaseFact(otherTodo), - normalizeAgent(candidate.owner), - input.registered_agents, - ) === null, - write_scopes: normalizeWriteScopes(candidate.write_scopes), - }]; - }); - const writeScopes = normalizeWriteScopes(todo.required_write_scopes ?? []); - const decision = evaluateTaskLeaseAcquireDecision({ - handoff_mode: handoffMode, - registered_agents: [...input.registered_agents], - todo: todoFact, - lease: currentLease === undefined ? null : { - present: true, - active: currentActive, - status: typeof currentLease.status === "string" ? currentLease.status : null, - owner: normalizeAgent(currentLease.owner), - idempotency_key: typeof currentLease.idempotency_key === "string" - ? currentLease.idempotency_key : null, - version: leaseInteger(currentLease, "version") ?? 0, - lease_epoch: leaseEpoch(currentLease), - write_scopes: normalizeWriteScopes(currentLease.write_scopes), - acquire_ttl_seconds: leaseInteger(currentLease, "acquire_ttl_seconds"), - }, - other_leases: otherLeases, - command: { - owner: authority.owner, - idempotency_key: leaseRequest.idempotency_key, - ttl_seconds: leaseRequest.ttl_seconds, - write_scopes: writeScopes, - expected_version: leaseRequest.expected_version, - }, - }); + // Validate required scopes before planning; a caller cannot omit a required conflict. + const writeScopes = claimWriteScopes(todo); + const facts = canonicalTaskLeaseAcquireFacts(projection, input.goal_id, input.todo_id, input.registered_agents, input.now); + const decision = evaluateTaskLeaseAcquireDecision({handoff_mode: handoffMode, + registered_agents: [...input.registered_agents], ...facts, + command: {owner: authority.owner, idempotency_key: leaseRequest.idempotency_key, + ttl_seconds: leaseRequest.ttl_seconds, write_scopes: writeScopes, + expected_version: leaseRequest.expected_version}}); if (decision.outcome === "no_change") { if (currentLease === undefined) { throw new AuthorityStoreProtocolError( @@ -566,27 +509,9 @@ export async function executeCoordinationTodoClaim( if (decision.next_lease === null) { throw new AuthorityStoreProtocolError("lease acquire apply is missing next_lease"); } - const acquiredAt = utcIsoformat(input.now); - lease = { - schema_version: TASK_LEASE_SCHEMA_VERSION, - goal_id: input.goal_id, - todo_id: input.todo_id, - owner: authority.owner, - idempotency_key: leaseRequest.idempotency_key, - write_scopes: writeScopes, - acquire_ttl_seconds: leaseRequest.ttl_seconds, - version: leaseDecisionInteger(decision.next_lease.version, "next_lease.version"), - lease_epoch: leaseDecisionInteger( - decision.next_lease.lease_epoch, - "next_lease.lease_epoch", - ), - acquired_at: acquiredAt, - updated_at: acquiredAt, - expires_at: utcIsoformat( - new Date(input.now.valueOf() + leaseRequest.ttl_seconds * 1_000), - ), - status: "active", - } satisfies LeaseRecord; + lease = materializeTaskLeaseAcquire(input, {owner: authority.owner, + idempotency_key: leaseRequest.idempotency_key, ttl_seconds: leaseRequest.ttl_seconds, + write_scopes: writeScopes, expected_version: leaseRequest.expected_version}, decision, input.now); leaseChanged = true; } else { return failure( diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index bf94e3139a..218ac783af 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -119,8 +119,8 @@ import { validateDeliveryClaim } from "./work_items/delivery_outcome.ts"; import { evaluateTaskLeaseAcquireDecision, evaluateTaskLeaseWriteScopesOverlap, - executeTaskLeaseAcquire, -} from "./work_items/task_lease_acquire.ts"; +} from "./work_items/task_lease_acquire_decision.ts"; +import {executeTaskLeaseAcquire} from "./work_items/task_lease_acquire.ts"; import { executeTaskLeaseLifecycle } from "./work_items/task_lease_lifecycle.ts"; import { commitLocalAuthorityShadowEntry, diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index ce941765f4..0c4337d53d 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -100,6 +100,9 @@ def _route_terminal_call(command: str, call: Mapping[str, Any]) -> dict[str, Any goal_id=goal_id, project=call.get("project"), state_file=call.get("state_file"), + # Completion consumes canonical state and can rebuild its display just + # like archive. The legacy caller below still requires its source file. + require_existing=False, ) complete = command == "complete" return terminal_canonical_todo_if_promoted( diff --git a/loopx/control_plane/work_items/canonical_task_lease_lifecycle.ts b/loopx/control_plane/work_items/canonical_task_lease_lifecycle.ts index 6b1518b4dc..6af3b89f4b 100644 --- a/loopx/control_plane/work_items/canonical_task_lease_lifecycle.ts +++ b/loopx/control_plane/work_items/canonical_task_lease_lifecycle.ts @@ -1,3 +1,4 @@ +import type {AuthorityStore} from "../coordination/authority_store.ts"; import type {JsonObject} from "../effect_program.ts"; import {canonicalAuthoritySha256} from "../coordination/authority_store_codec.ts"; import {openLocalAuthorityStoreHandle, localAuthorityOpenFailure, @@ -10,7 +11,7 @@ import {executeCanonicalTaskLeaseLifecycle} from "../coordination/task_lease_lif import {revalidateAuthoritySources, TaskLeaseAcquireError, type AuthorityFacts} from "./task_lease_acquire.ts"; import type {TaskLeaseLifecycleDecisionOperation} from "./task_lease_lifecycle_decision.ts"; -interface LocalLeaseRequest { +export interface LocalLeaseRequest { operation: TaskLeaseLifecycleDecisionOperation; runtime_root: string; goal_id: string; todo_id: string; owner: string | null; idempotency_key: string | null; @@ -19,11 +20,20 @@ interface LocalLeaseRequest { authority: AuthorityFacts | null; } -/** One fenced local opening boundary. A selected service provider must supply - * its service-owned factory; request data can never carry credentials/clients. */ -export async function mutateCanonicalTaskLease(request: LocalLeaseRequest, - dependencies: {now: () => Date; beforeWrite?: (lease: JsonObject) => void | Promise; - authorityProvider?: LocalAuthorityProviderDependencies}): Promise { +export interface LocalTaskLeaseDependencies { + now: () => Date; + beforeWrite?: (lease: JsonObject) => void | Promise; + authorityProvider?: LocalAuthorityProviderDependencies; +} + +/** The same local promotion/source fence surrounds acquire and maintenance. */ +export async function withCanonicalTaskLeaseAuthority( + request: {runtime_root: string; goal_id: string; operation: string; + owner: string | null; idempotency_key: string | null; authority: AuthorityFacts | null}, + dependencies: LocalTaskLeaseDependencies, + execute: (store: AuthorityStore, guards: { + revalidate: () => Promise; beforeCommit: (lease: JsonObject | null) => Promise; + }) => Promise): Promise { const root = request.runtime_root, goalId = request.goal_id; const initialFence = await loadLegacyCoordinationWriterFence(root, goalId); const evidence: JsonObject = {source_authority: null, decision_read_from_provider: false, legacy_fallback_used: false}; @@ -46,19 +56,17 @@ export async function mutateCanonicalTaskLease(request: LocalLeaseRequest, return rejected("authority_required", "canonical lease mutation needs registered actor context"); } evidence.decision_read_from_provider = true; - const result = await executeCanonicalTaskLeaseLifecycle(store, {operation: request.operation, - goal_id: goalId, todo_id: request.todo_id, owner: request.owner, idempotency_key: request.idempotency_key, - expected_version: request.expected_version, ttl_seconds: request.ttl_seconds, - new_owner: request.new_owner, new_idempotency_key: request.new_idempotency_key, - registered_agents: request.authority?.registered_agents ?? [], now: dependencies.now()}, async lease => { - if (lease) await dependencies.beforeWrite?.(lease); + const revalidate = async () => { await verifyFence(); - // Release only relinquishes the existing proof. Removed actors and a - // missing registry must still be able to clean up their own lease. + // Release relinquishes an existing proof even after actor deregistration. if (request.operation !== "release" && request.authority) { await revalidateAuthoritySources(request.authority.source_receipts); } - }); + }; + const result = await execute(store, {revalidate, beforeCommit: async lease => { + if (lease) await dependencies.beforeWrite?.(lease); + await revalidate(); + }}); return {...result, ...evidence}; }); } catch (error) { @@ -67,3 +75,13 @@ export async function mutateCanonicalTaskLease(request: LocalLeaseRequest, return {...rejected("canonical_lease_route_failed", error instanceof Error ? error.message : "canonical lease route failed"), ...localAuthorityOpenFailure(error)}; } } + +export async function mutateCanonicalTaskLease(request: LocalLeaseRequest, + dependencies: LocalTaskLeaseDependencies): Promise { + return await withCanonicalTaskLeaseAuthority(request, dependencies, (store, guards) => + executeCanonicalTaskLeaseLifecycle(store, {operation: request.operation, + goal_id: request.goal_id, todo_id: request.todo_id, owner: request.owner!, idempotency_key: request.idempotency_key!, + expected_version: request.expected_version, ttl_seconds: request.ttl_seconds, + new_owner: request.new_owner, new_idempotency_key: request.new_idempotency_key, + registered_agents: request.authority?.registered_agents ?? [], now: dependencies.now()}, guards.beforeCommit)); +} diff --git a/loopx/control_plane/work_items/task_lease_acquire.ts b/loopx/control_plane/work_items/task_lease_acquire.ts index f7f5684b41..6205d5b923 100644 --- a/loopx/control_plane/work_items/task_lease_acquire.ts +++ b/loopx/control_plane/work_items/task_lease_acquire.ts @@ -1,3 +1,6 @@ +import {executeCanonicalTaskLeaseAcquire} from "../coordination/task_lease_acquire.ts"; +import {withCanonicalTaskLeaseAuthority} from "./canonical_task_lease_lifecycle.ts"; +import {evaluateTaskLeaseAcquireDecision, materializeTaskLeaseAcquire, type AcquireDecisionOtherLease} from "./task_lease_acquire_decision.ts"; import {leaseOwnerRejection as ownerRejection} from "./task_lease_eligibility.ts"; import { ShadowManagementError, requireShadowPrimaryWriteAllowed } from "../coordination/shadow_management.ts"; import { parseIsoTimestamp } from "../runtime_timestamp.ts"; @@ -25,7 +28,7 @@ import { decodeLocalAuthorityShadowBinding, type LocalAuthorityShadowBinding, } from "../coordination/local_authority_shadow_outbox.ts"; -import { TASK_LEASE_ACQUIRE_REQUEST_SCHEMA } from "../coordination/coordination_state_contract.generated.ts"; +import { TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA } from "../coordination/coordination_state_contract.generated.ts"; export const TASK_LEASE_ACQUIRE_REQUEST_SCHEMA_VERSION = TASK_LEASE_ACQUIRE_REQUEST_SCHEMA; @@ -74,6 +77,7 @@ export interface AuthorityFacts { } interface AcquireRequest { + canonical: boolean; runtime_root: string; goal_id: string; owner: string; @@ -102,48 +106,6 @@ export interface LeaseRecord extends JsonObject { expires_at?: unknown; } -interface AcquireDecisionLease { - present: boolean; - active: boolean; - status: string | null; - owner: string | null; - idempotency_key: string | null; - version: number; - lease_epoch: number; - write_scopes: readonly string[]; - acquire_ttl_seconds: number | null; -} - -interface AcquireDecisionOtherLease { - todo_id: string; - active: boolean; - effective: boolean; - write_scopes: readonly string[]; -} - -interface AcquireDecisionInput { - handoff_mode: string; - registered_agents: readonly string[]; - todo: TodoFact | null; - lease: AcquireDecisionLease | null; - other_leases: readonly AcquireDecisionOtherLease[]; - command: { - owner: string; - idempotency_key: string; - ttl_seconds: number; - write_scopes: readonly string[]; - expected_version: number | null; - }; -} - -interface AcquireDecision extends JsonObject { - outcome: "apply" | "no_change" | "conflict" | "rejected"; - code: string; - idempotent: boolean; - next_lease: JsonObject | null; - conflict_indexes: number[]; -} - interface TaskLeaseFailure { code: string; message: string; @@ -158,6 +120,7 @@ interface ExecutionContext { } export interface TaskLeaseAcquireDependencies { + authorityProvider?: import("../coordination/local_authority_provider.ts").LocalAuthorityProviderDependencies; now?: () => Date; beforeWrite?: (lease: JsonObject) => void | Promise; } @@ -435,14 +398,23 @@ function normalizeHandoffMode(value: unknown): string { function decodeRequest(value: unknown): AcquireRequest { const request = requireJsonObject(value, "task lease acquire request"); - if (request.schema_version !== TASK_LEASE_ACQUIRE_REQUEST_SCHEMA_VERSION) { + const canonical = request.schema_version === TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA; + if (request.schema_version !== TASK_LEASE_ACQUIRE_REQUEST_SCHEMA_VERSION && !canonical) { throw new EffectRuntimeRequestError("Task-lease acquire request schema mismatch"); } // Decode the authority envelope first. Besides keeping the boundary // fail-closed, this preserves the public error ordering used by the native // CLI: a missing authority projection is reported before unrelated fields. - const authority = decodeTaskLeaseAuthority(request.authority); + if (canonical) { + const allowed = new Set(["schema_version", "runtime_root", "goal_id", "todo_id", "owner", "idempotency_key", + "ttl_seconds", "write_scopes", "expected_version", "authority"]); + const unsupported = Object.keys(request).find(key => !allowed.has(key)); + if (unsupported) throw new TaskLeaseAcquireError(`canonical acquire does not accept ${unsupported}`, "invalid_canonical_acquire_request"); + } + const rawAuthority = requireJsonObject(request.authority, "acquire authority"); + const authority = decodeTaskLeaseAuthority(canonical ? {...rawAuthority, handoff_mode: "legacy", todos: [], todo_projection_error: null} : rawAuthority); return { + canonical, runtime_root: stringValue(request.runtime_root, "runtime_root"), goal_id: normalizeGoalId(request.goal_id), owner: normalizeOwner(request.owner), @@ -631,321 +603,6 @@ function ownerFailure( }); } -function classMatch( - pattern: string, - start: number, - value: string, -): { end: number; matches: boolean } | null { - let end = start + 1; - if (pattern[end] === "!") end += 1; - if (pattern[end] === "]") end += 1; - end = pattern.indexOf("]", end); - if (end < 0) return null; - let body = pattern.slice(start + 1, end); - const negated = body.startsWith("!"); - if (negated) body = body.slice(1); - let matches = false; - for (let index = 0; index < body.length; index += 1) { - if (index + 2 < body.length && body[index + 1] === "-") { - if (body[index] <= value && value <= body[index + 2]) matches = true; - index += 2; - } else if (body[index] === value) { - matches = true; - } - } - return { end, matches: negated ? !matches : matches }; -} - -function fnmatchcase(value: string, pattern: string): boolean { - const memo = new Map(); - const match = (valueIndex: number, patternIndex: number): boolean => { - const key = `${valueIndex}:${patternIndex}`; - const cached = memo.get(key); - if (cached !== undefined) return cached; - let result: boolean; - if (patternIndex === pattern.length) { - result = valueIndex === value.length; - } else if (pattern[patternIndex] === "*") { - result = match(valueIndex, patternIndex + 1) || - (valueIndex < value.length && match(valueIndex + 1, patternIndex)); - } else if (valueIndex === value.length) { - result = false; - } else if (pattern[patternIndex] === "?") { - result = match(valueIndex + 1, patternIndex + 1); - } else if (pattern[patternIndex] === "[") { - const characterClass = classMatch(pattern, patternIndex, value[valueIndex]); - result = characterClass === null - ? value[valueIndex] === "[" && match(valueIndex + 1, patternIndex + 1) - : characterClass.matches && match(valueIndex + 1, characterClass.end + 1); - } else { - result = value[valueIndex] === pattern[patternIndex] && - match(valueIndex + 1, patternIndex + 1); - } - memo.set(key, result); - return result; - }; - return match(0, 0); -} - -function scopeLiteralPrefix(scope: string): string { - const indexes = ["*", "?", "["] - .map((token) => scope.indexOf(token)) - .filter((index) => index >= 0); - return indexes.length > 0 ? scope.slice(0, Math.min(...indexes)) : scope; -} - -function scopePairOverlaps(left: string, right: string): boolean { - if (left === right) return true; - if (["*", "**", "./"].includes(left) || ["*", "**", "./"].includes(right)) { - return true; - } - const leftGlob = ["*", "?", "["].some((token) => left.includes(token)); - const rightGlob = ["*", "?", "["].some((token) => right.includes(token)); - if (leftGlob && !rightGlob) { - const prefix = scopeLiteralPrefix(left); - return fnmatchcase(right, left) || - (prefix.endsWith("/") && right.replace(/\/$/u, "") === prefix.replace(/\/$/u, "")); - } - if (rightGlob && !leftGlob) { - const prefix = scopeLiteralPrefix(right); - return fnmatchcase(left, right) || - (prefix.endsWith("/") && left.replace(/\/$/u, "") === prefix.replace(/\/$/u, "")); - } - if (leftGlob && rightGlob) { - const leftPrefix = scopeLiteralPrefix(left); - const rightPrefix = scopeLiteralPrefix(right); - return !leftPrefix || !rightPrefix || leftPrefix.startsWith(rightPrefix) || - rightPrefix.startsWith(leftPrefix); - } - const leftRoot = left.replace(/\/$/u, ""); - const rightRoot = right.replace(/\/$/u, ""); - return (left.endsWith("/") && right.startsWith(`${leftRoot}/`)) || - (right.endsWith("/") && left.startsWith(`${rightRoot}/`)); -} - -function writeScopesOverlap(left: readonly string[], right: readonly string[]): boolean { - if (left.length === 0 || right.length === 0) return false; - return left.some((a) => right.some((b) => scopePairOverlaps(a, b))); -} - -function decisionStringArray(value: unknown, label: string): string[] { - if (!Array.isArray(value) || value.some((item) => typeof item !== "string")) { - throw new EffectRuntimeRequestError(`${label} must be an array of strings`); - } - return [...value] as string[]; -} - -export function evaluateTaskLeaseWriteScopesOverlap(value: unknown): JsonObject { - const input = requireJsonObject(value, "task lease write-scope overlap"); - return { - overlap: writeScopesOverlap( - decisionStringArray(input.left, "left"), - decisionStringArray(input.right, "right"), - ), - }; -} - -function decisionBoolean(value: unknown, label: string): boolean { - if (typeof value !== "boolean") { - throw new EffectRuntimeRequestError(`${label} must be a boolean`); - } - return value; -} - -function decisionInteger(value: unknown, label: string): number { - if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { - throw new EffectRuntimeRequestError(`${label} must be a non-negative safe integer`); - } - return value; -} - -function decisionNullableString(value: unknown, label: string): string | null { - if (value === null || value === undefined) return null; - if (typeof value !== "string") { - throw new EffectRuntimeRequestError(`${label} must be a string or null`); - } - return value; -} - -function decodeDecisionTodo(value: unknown): TodoFact | null { - if (value === null || value === undefined) return null; - const todo = requireJsonObject(value, "task lease acquire decision todo"); - return { - todo_id: stringValue(todo.todo_id, "todo.todo_id"), - status: stringValue(todo.status, "todo.status"), - claimed_by: decisionNullableString(todo.claimed_by, "todo.claimed_by"), - excluded_agents: decisionStringArray(todo.excluded_agents, "todo.excluded_agents"), - }; -} - -function decodeDecisionLease(value: unknown): AcquireDecisionLease | null { - if (value === null || value === undefined) return null; - const lease = requireJsonObject(value, "task lease acquire decision lease"); - // Accept old callers without trusting their derived eligibility hint. - if (lease.effective !== undefined) decisionBoolean(lease.effective, "lease.effective"); - return { - present: decisionBoolean(lease.present, "lease.present"), - active: decisionBoolean(lease.active, "lease.active"), - status: decisionNullableString(lease.status, "lease.status"), - owner: decisionNullableString(lease.owner, "lease.owner"), - idempotency_key: decisionNullableString( - lease.idempotency_key, - "lease.idempotency_key", - ), - version: decisionInteger(lease.version, "lease.version"), - lease_epoch: decisionInteger(lease.lease_epoch, "lease.lease_epoch"), - write_scopes: decisionStringArray(lease.write_scopes, "lease.write_scopes"), - acquire_ttl_seconds: optionalInteger( - lease.acquire_ttl_seconds, - "lease.acquire_ttl_seconds", - ), - }; -} - -function decodeAcquireDecisionInput(value: unknown): AcquireDecisionInput { - const input = requireJsonObject(value, "task lease acquire decision"); - const command = requireJsonObject(input.command, "task lease acquire decision command"); - const rawOtherLeases = input.other_leases; - if (!Array.isArray(rawOtherLeases)) { - throw new EffectRuntimeRequestError("other_leases must be an array"); - } - const otherLeases = rawOtherLeases.map((raw, index) => { - const lease = requireJsonObject(raw, `other_leases[${index}]`); - return { - todo_id: stringValue(lease.todo_id, `other_leases[${index}].todo_id`), - active: decisionBoolean(lease.active, `other_leases[${index}].active`), - effective: decisionBoolean(lease.effective, `other_leases[${index}].effective`), - write_scopes: decisionStringArray( - lease.write_scopes, - `other_leases[${index}].write_scopes`, - ), - }; - }); - return { - handoff_mode: stringValue(input.handoff_mode, "handoff_mode"), - registered_agents: decisionStringArray( - input.registered_agents, - "registered_agents", - ), - todo: decodeDecisionTodo(input.todo), - lease: decodeDecisionLease(input.lease), - other_leases: otherLeases, - command: { - owner: stringValue(command.owner, "command.owner"), - idempotency_key: stringValue( - command.idempotency_key, - "command.idempotency_key", - ), - ttl_seconds: decisionInteger(command.ttl_seconds, "command.ttl_seconds"), - write_scopes: decisionStringArray( - command.write_scopes, - "command.write_scopes", - ), - expected_version: optionalInteger( - command.expected_version, - "command.expected_version", - ), - }, - }; -} - -function acquireDecisionResult( - outcome: AcquireDecision["outcome"], - code: string, - options: { - idempotent?: boolean; - nextLease?: JsonObject | null; - conflictIndexes?: number[]; - } = {}, -): AcquireDecision { - return { - outcome, - code, - idempotent: options.idempotent ?? false, - next_lease: options.nextLease ?? null, - conflict_indexes: options.conflictIndexes ?? [], - }; -} - -/** - * Canonical pure decision for both local file acquire and shared coordination. - * Locking, source revalidation, persistence, provider CAS, and receipts stay in - * their respective execution layers. - */ -export function evaluateTaskLeaseAcquireDecision(value: unknown): AcquireDecision { - const input = decodeAcquireDecisionInput(value); - const { command, lease } = input; - if (lease !== null && lease.active && (!lease.present || lease.status === "released")) { - return acquireDecisionResult("rejected", "invalid_lease_snapshot"); - } - if (input.handoff_mode === "soft_claim") { - return acquireDecisionResult("rejected", "handoff_mode_forbids_lease"); - } - const rejection = ownerRejection( - input.todo ?? undefined, - command.owner || null, - input.registered_agents, - ); - if (rejection !== null) { - return acquireDecisionResult("rejected", rejection); - } - const actualVersion = lease !== null && lease.present ? lease.version : 0; - if ( - command.expected_version !== null && command.expected_version !== actualVersion - ) { - return acquireDecisionResult("conflict", "version_mismatch"); - } - // The old wire effective hint is not authority over the supplied owner facts. - if (lease !== null && lease.present && lease.active && - ownerRejection(input.todo, lease.owner, input.registered_agents) === null) { - if ( - lease.owner === command.owner && - lease.idempotency_key === command.idempotency_key - ) { - const scopesMatch = equalScopeSets(lease.write_scopes, command.write_scopes); - const ttlMatches = lease.acquire_ttl_seconds === null || - lease.acquire_ttl_seconds === command.ttl_seconds; - if (!scopesMatch || !ttlMatches) { - return acquireDecisionResult("rejected", "idempotency_key_reuse"); - } - return acquireDecisionResult("no_change", "lease_acquire_replay", { - idempotent: true, - }); - } - return acquireDecisionResult("conflict", "todo_lease_conflict"); - } - if ( - lease !== null && lease.present && - lease.idempotency_key === command.idempotency_key - ) { - return acquireDecisionResult("rejected", "idempotency_key_reuse"); - } - const conflictIndexes = input.other_leases.flatMap((other, index) => - other.active && other.effective && - writeScopesOverlap(command.write_scopes, other.write_scopes) - ? [index] - : [] - ); - if (conflictIndexes.length > 0) { - return acquireDecisionResult("conflict", "write_scope_conflict", { - conflictIndexes, - }); - } - return acquireDecisionResult("apply", "lease_acquire", { - nextLease: { - present: true, - active: true, - status: "active", - owner: command.owner, - idempotency_key: command.idempotency_key, - version: actualVersion + 1, - lease_epoch: (lease?.lease_epoch ?? 0) + 1, - write_scopes: [...command.write_scopes], - acquire_ttl_seconds: command.ttl_seconds, - }, - }); -} - async function currentSourceReceipt(receipt: SourceReceipt): Promise { try { const bytes = await readFile(receipt.path); @@ -1108,12 +765,6 @@ function transitionError( ); } -function equalScopeSets(left: readonly string[], right: readonly string[]): boolean { - const a = [...new Set(left)].sort((first, second) => first.localeCompare(second)); - const b = [...new Set(right)].sort((first, second) => first.localeCompare(second)); - return a.length === b.length && a.every((value, index) => value === b[index]); -} - function successEnvelope( request: AcquireRequest, lease: LeaseRecord, @@ -1163,6 +814,7 @@ const VALIDATION_FAILURE_CODES = new Set([ ...PERMISSION_DENIED_CODES, "todo_lease_conflict", "write_scope_conflict", + "lease_generation_exhausted", "authority_source_changed", "corrupt_lease", ]); @@ -1318,27 +970,7 @@ async function commitAcquire( ); } - const acquiredAt = utcIsoformat(at); - const lease: LeaseRecord = { - schema_version: TASK_LEASE_SCHEMA_VERSION, - goal_id: request.goal_id, - todo_id: request.todo_id, - owner: request.owner, - idempotency_key: request.idempotency_key, - write_scopes: [...request.write_scopes], - acquire_ttl_seconds: request.ttl_seconds, - version: decisionInteger(decision.next_lease.version, "next_lease.version"), - lease_epoch: decisionInteger( - decision.next_lease.lease_epoch, - "next_lease.lease_epoch", - ), - acquired_at: acquiredAt, - updated_at: acquiredAt, - expires_at: utcIsoformat( - new Date(at.valueOf() + request.ttl_seconds * 1_000), - ), - status: "active", - }; + const lease = materializeTaskLeaseAcquire(request, request, decision, at); await dependencies.beforeWrite?.(lease); await revalidateAuthoritySources(request.authority.source_receipts); const captureRequired = request.runtime_shadow !== null || @@ -1392,6 +1024,28 @@ export async function executeTaskLeaseAcquire( } try { + if (request.canonical) { + const result = await withCanonicalTaskLeaseAuthority({...request, operation: "acquire"}, { + ...dependencies, now: dependencies.now ?? (() => new Date()), + }, async (store, guards) => { + await guards.revalidate(); + return await executeCanonicalTaskLeaseAcquire(store, {...request, + registered_agents: request.authority.registered_agents, now: (dependencies.now ?? (() => new Date()))()}, guards.beforeCommit); + }); + if (["applied", "no_change", "replayed", "recovered"].includes(String(result.status))) { + const idempotent = result.status !== "applied"; + const envelope = successEnvelope(request, result.lease as LeaseRecord, "", acquireEffectId(request), idempotent); + delete envelope.lease_path; + return {...result, ...envelope, acquired: !idempotent && result.changed === true}; + } + const envelope = failureEnvelope({code: String(result.reason_code ?? result.conflict_kind ?? "canonical_acquire_failed"), + message: String(result.reason ?? "canonical acquisition failed; inspect before retrying"), + kind: result.reason_code === "legacy_writer_fence_read_failed" ? "permission_denied" : undefined, + payload: result, stage: result.failure_stage === "durable_writeback" ? "durable_writeback" : "validation"}, context); + delete envelope.lease_path; + // Keep the public acquire schema even when the inner command has diagnostics. + return {...envelope, schema_version: TASK_LEASE_SCHEMA_VERSION}; + } return await withFileMutationLock( legacyCoordinationLeaseLockPath(request.runtime_root, request.goal_id), () => withFileMutationLock(taskLeaseLockPath(request), async () => { diff --git a/loopx/control_plane/work_items/task_lease_acquire_adapter.py b/loopx/control_plane/work_items/task_lease_acquire_adapter.py index 0b37f53c33..0de1a0d388 100644 --- a/loopx/control_plane/work_items/task_lease_acquire_adapter.py +++ b/loopx/control_plane/work_items/task_lease_acquire_adapter.py @@ -14,6 +14,7 @@ from ..coordination.coordination_state_contract_generated import ( LOCAL_AUTHORITY_SHADOW_BINDING_SCHEMA, TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, + TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA, TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, ) @@ -410,14 +411,17 @@ def execute_native_task_lease_acquire( from ..effect_runtime import effect_runtime_result + from ..coordination.local_authority import local_authority_is_promoted + + canonical = local_authority_is_promoted(runtime_root=runtime_root, goal_id=str(goal_id)) for attempt in range(TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS): - authority = task_lease_acquire_authority_facts( + authority = _canonical_lease_authority_facts(registry_path, str(goal_id)) if canonical else task_lease_acquire_authority_facts( registry_path=registry_path, goal_id=str(goal_id or ""), todo_id=str(todo_id or ""), ) request = { - "schema_version": TASK_LEASE_ACQUIRE_NATIVE_SCHEMA_VERSION, + "schema_version": TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA if canonical else TASK_LEASE_ACQUIRE_NATIVE_SCHEMA_VERSION, "runtime_root": str(runtime_root), "goal_id": goal_id, "todo_id": todo_id, @@ -430,7 +434,7 @@ def execute_native_task_lease_acquire( } registry = load_registry(registry_path) goal = _registry_goal(registry, str(goal_id)) - if resolve_coordination_runtime_shadow_config(goal).enabled: + if not canonical and resolve_coordination_runtime_shadow_config(goal).enabled: request["runtime_shadow"] = { "schema_version": LOCAL_AUTHORITY_SHADOW_BINDING_SCHEMA, "provider": "file_v0", @@ -443,6 +447,14 @@ def execute_native_task_lease_acquire( and attempt + 1 < TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS ): continue + if canonical: + if payload.get("ok") is True and ( + payload.get("source_authority") not in ("file_v0", "sqlite_v0") + or payload.get("decision_read_from_provider") is not True + or payload.get("legacy_fallback_used") is not False + ): + raise RuntimeError("canonical acquire omitted valid provider evidence") + return payload return _finalize_native_acquire_result( payload, authority=authority, diff --git a/loopx/control_plane/work_items/task_lease_acquire_decision.ts b/loopx/control_plane/work_items/task_lease_acquire_decision.ts new file mode 100644 index 0000000000..0fbd3c248c --- /dev/null +++ b/loopx/control_plane/work_items/task_lease_acquire_decision.ts @@ -0,0 +1,407 @@ +/** Shared acquire/reclaim admission. IO and durable receipts belong to callers. */ +import {leaseOwnerRejection as ownerRejection} from "./task_lease_eligibility.ts"; +import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import type {JsonObject} from "../effect_program.ts"; +import type {TodoFact, LeaseRecord} from "./task_lease_acquire.ts"; + +export interface AcquireDecisionLease { + present: boolean; + active: boolean; + status: string | null; + owner: string | null; + idempotency_key: string | null; + version: number; + lease_epoch: number; + write_scopes: readonly string[]; + acquire_ttl_seconds: number | null; +} + +export interface AcquireDecisionOtherLease { + todo_id: string; + active: boolean; + effective: boolean; + write_scopes: readonly string[]; +} + +export interface AcquireDecisionInput { + handoff_mode: string; + registered_agents: readonly string[]; + todo: TodoFact | null; + lease: AcquireDecisionLease | null; + other_leases: readonly AcquireDecisionOtherLease[]; + command: { + owner: string; + idempotency_key: string; + ttl_seconds: number; + write_scopes: readonly string[]; + expected_version: number | null; + }; +} + +export interface AcquireDecision extends JsonObject { + outcome: "apply" | "no_change" | "conflict" | "rejected"; + code: string; + idempotent: boolean; + next_lease: JsonObject | null; + conflict_indexes: number[]; +} + +function stringValue(value: unknown, label: string): string { + if (typeof value !== "string") { + throw new EffectRuntimeRequestError(`${label} must be a string`); + } + return value; +} + +function optionalInteger(value: unknown, label: string): number | null { + if (value === null || value === undefined) return null; + if (typeof value !== "number" || !Number.isInteger(value)) { + throw new EffectRuntimeRequestError(`${label} must be an integer or null`); + } + return value; +} + +function classMatch( + pattern: string, + start: number, + value: string, +): { end: number; matches: boolean } | null { + let end = start + 1; + if (pattern[end] === "!") end += 1; + if (pattern[end] === "]") end += 1; + end = pattern.indexOf("]", end); + if (end < 0) return null; + let body = pattern.slice(start + 1, end); + const negated = body.startsWith("!"); + if (negated) body = body.slice(1); + let matches = false; + for (let index = 0; index < body.length; index += 1) { + if (index + 2 < body.length && body[index + 1] === "-") { + if (body[index] <= value && value <= body[index + 2]) matches = true; + index += 2; + } else if (body[index] === value) { + matches = true; + } + } + return { end, matches: negated ? !matches : matches }; +} + +function fnmatchcase(value: string, pattern: string): boolean { + const memo = new Map(); + const match = (valueIndex: number, patternIndex: number): boolean => { + const key = `${valueIndex}:${patternIndex}`; + const cached = memo.get(key); + if (cached !== undefined) return cached; + let result: boolean; + if (patternIndex === pattern.length) { + result = valueIndex === value.length; + } else if (pattern[patternIndex] === "*") { + result = match(valueIndex, patternIndex + 1) || + (valueIndex < value.length && match(valueIndex + 1, patternIndex)); + } else if (valueIndex === value.length) { + result = false; + } else if (pattern[patternIndex] === "?") { + result = match(valueIndex + 1, patternIndex + 1); + } else if (pattern[patternIndex] === "[") { + const characterClass = classMatch(pattern, patternIndex, value[valueIndex]); + result = characterClass === null + ? value[valueIndex] === "[" && match(valueIndex + 1, patternIndex + 1) + : characterClass.matches && match(valueIndex + 1, characterClass.end + 1); + } else { + result = value[valueIndex] === pattern[patternIndex] && + match(valueIndex + 1, patternIndex + 1); + } + memo.set(key, result); + return result; + }; + return match(0, 0); +} + +function scopeLiteralPrefix(scope: string): string { + const indexes = ["*", "?", "["] + .map((token) => scope.indexOf(token)) + .filter((index) => index >= 0); + return indexes.length > 0 ? scope.slice(0, Math.min(...indexes)) : scope; +} + +function scopePairOverlaps(left: string, right: string): boolean { + if (left === right) return true; + if (["*", "**", "./"].includes(left) || ["*", "**", "./"].includes(right)) { + return true; + } + const leftGlob = ["*", "?", "["].some((token) => left.includes(token)); + const rightGlob = ["*", "?", "["].some((token) => right.includes(token)); + if (leftGlob && !rightGlob) { + const prefix = scopeLiteralPrefix(left); + return fnmatchcase(right, left) || + (prefix.endsWith("/") && right.replace(/\/$/u, "") === prefix.replace(/\/$/u, "")); + } + if (rightGlob && !leftGlob) { + const prefix = scopeLiteralPrefix(right); + return fnmatchcase(left, right) || + (prefix.endsWith("/") && left.replace(/\/$/u, "") === prefix.replace(/\/$/u, "")); + } + if (leftGlob && rightGlob) { + const leftPrefix = scopeLiteralPrefix(left); + const rightPrefix = scopeLiteralPrefix(right); + return !leftPrefix || !rightPrefix || leftPrefix.startsWith(rightPrefix) || + rightPrefix.startsWith(leftPrefix); + } + const leftRoot = left.replace(/\/$/u, ""); + const rightRoot = right.replace(/\/$/u, ""); + return (left.endsWith("/") && right.startsWith(`${leftRoot}/`)) || + (right.endsWith("/") && left.startsWith(`${rightRoot}/`)); +} + +function writeScopesOverlap(left: readonly string[], right: readonly string[]): boolean { + if (left.length === 0 || right.length === 0) return false; + return left.some((a) => right.some((b) => scopePairOverlaps(a, b))); +} + +function decisionStringArray(value: unknown, label: string): string[] { + if (!Array.isArray(value) || value.some((item) => typeof item !== "string")) { + throw new EffectRuntimeRequestError(`${label} must be an array of strings`); + } + return [...value] as string[]; +} + +export function evaluateTaskLeaseWriteScopesOverlap(value: unknown): JsonObject { + const input = requireJsonObject(value, "task lease write-scope overlap"); + return { + overlap: writeScopesOverlap( + decisionStringArray(input.left, "left"), + decisionStringArray(input.right, "right"), + ), + }; +} + +function decisionBoolean(value: unknown, label: string): boolean { + if (typeof value !== "boolean") { + throw new EffectRuntimeRequestError(`${label} must be a boolean`); + } + return value; +} + +function decisionInteger(value: unknown, label: string): number { + if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { + throw new EffectRuntimeRequestError(`${label} must be a non-negative safe integer`); + } + return value; +} + +function decisionNullableString(value: unknown, label: string): string | null { + if (value === null || value === undefined) return null; + if (typeof value !== "string") { + throw new EffectRuntimeRequestError(`${label} must be a string or null`); + } + return value; +} + +function decodeDecisionTodo(value: unknown): TodoFact | null { + if (value === null || value === undefined) return null; + const todo = requireJsonObject(value, "task lease acquire decision todo"); + return { + todo_id: stringValue(todo.todo_id, "todo.todo_id"), + status: stringValue(todo.status, "todo.status"), + claimed_by: decisionNullableString(todo.claimed_by, "todo.claimed_by"), + excluded_agents: decisionStringArray(todo.excluded_agents, "todo.excluded_agents"), + }; +} + +function decodeDecisionLease(value: unknown): AcquireDecisionLease | null { + if (value === null || value === undefined) return null; + const lease = requireJsonObject(value, "task lease acquire decision lease"); + // Accept old callers without trusting their derived eligibility hint. + if (lease.effective !== undefined) decisionBoolean(lease.effective, "lease.effective"); + return { + present: decisionBoolean(lease.present, "lease.present"), + active: decisionBoolean(lease.active, "lease.active"), + status: decisionNullableString(lease.status, "lease.status"), + owner: decisionNullableString(lease.owner, "lease.owner"), + idempotency_key: decisionNullableString( + lease.idempotency_key, + "lease.idempotency_key", + ), + version: decisionInteger(lease.version, "lease.version"), + lease_epoch: decisionInteger(lease.lease_epoch, "lease.lease_epoch"), + write_scopes: decisionStringArray(lease.write_scopes, "lease.write_scopes"), + acquire_ttl_seconds: optionalInteger( + lease.acquire_ttl_seconds, + "lease.acquire_ttl_seconds", + ), + }; +} + +function decodeAcquireDecisionInput(value: unknown): AcquireDecisionInput { + const input = requireJsonObject(value, "task lease acquire decision"); + const command = requireJsonObject(input.command, "task lease acquire decision command"); + const rawOtherLeases = input.other_leases; + if (!Array.isArray(rawOtherLeases)) { + throw new EffectRuntimeRequestError("other_leases must be an array"); + } + const otherLeases = rawOtherLeases.map((raw, index) => { + const lease = requireJsonObject(raw, `other_leases[${index}]`); + return { + todo_id: stringValue(lease.todo_id, `other_leases[${index}].todo_id`), + active: decisionBoolean(lease.active, `other_leases[${index}].active`), + effective: decisionBoolean(lease.effective, `other_leases[${index}].effective`), + write_scopes: decisionStringArray( + lease.write_scopes, + `other_leases[${index}].write_scopes`, + ), + }; + }); + return { + handoff_mode: stringValue(input.handoff_mode, "handoff_mode"), + registered_agents: decisionStringArray( + input.registered_agents, + "registered_agents", + ), + todo: decodeDecisionTodo(input.todo), + lease: decodeDecisionLease(input.lease), + other_leases: otherLeases, + command: { + owner: stringValue(command.owner, "command.owner"), + idempotency_key: stringValue( + command.idempotency_key, + "command.idempotency_key", + ), + ttl_seconds: decisionInteger(command.ttl_seconds, "command.ttl_seconds"), + write_scopes: decisionStringArray( + command.write_scopes, + "command.write_scopes", + ), + expected_version: optionalInteger( + command.expected_version, + "command.expected_version", + ), + }, + }; +} + +function acquireDecisionResult( + outcome: AcquireDecision["outcome"], + code: string, + options: { + idempotent?: boolean; + nextLease?: JsonObject | null; + conflictIndexes?: number[]; + } = {}, +): AcquireDecision { + return { + outcome, + code, + idempotent: options.idempotent ?? false, + next_lease: options.nextLease ?? null, + conflict_indexes: options.conflictIndexes ?? [], + }; +} + +/** + * Canonical pure decision for both local file acquire and shared coordination. + * Locking, source revalidation, persistence, provider CAS, and receipts stay in + * their respective execution layers. + */ +export function evaluateTaskLeaseAcquireDecision(value: unknown): AcquireDecision { + return decideTaskLeaseAcquire(decodeAcquireDecisionInput(value)); +} + +export function decideTaskLeaseAcquire(input: AcquireDecisionInput): AcquireDecision { + const { command, lease } = input; + if (lease !== null && lease.active && (!lease.present || lease.status === "released")) { + return acquireDecisionResult("rejected", "invalid_lease_snapshot"); + } + if (input.handoff_mode === "soft_claim") { + return acquireDecisionResult("rejected", "handoff_mode_forbids_lease"); + } + const rejection = ownerRejection( + input.todo ?? undefined, + command.owner || null, + input.registered_agents, + ); + if (rejection !== null) { + return acquireDecisionResult("rejected", rejection); + } + const actualVersion = lease !== null && lease.present ? lease.version : 0; + if ( + command.expected_version !== null && command.expected_version !== actualVersion + ) { + return acquireDecisionResult("conflict", "version_mismatch"); + } + // The old wire effective hint is not authority over the supplied owner facts. + if (lease !== null && lease.present && lease.active && + ownerRejection(input.todo, lease.owner, input.registered_agents) === null) { + if ( + lease.owner === command.owner && + lease.idempotency_key === command.idempotency_key + ) { + const scopesMatch = equalScopeSets(lease.write_scopes, command.write_scopes); + const ttlMatches = lease.acquire_ttl_seconds === null || + lease.acquire_ttl_seconds === command.ttl_seconds; + if (!scopesMatch || !ttlMatches) { + return acquireDecisionResult("rejected", "idempotency_key_reuse"); + } + return acquireDecisionResult("no_change", "lease_acquire_replay", { + idempotent: true, + }); + } + return acquireDecisionResult("conflict", "todo_lease_conflict"); + } + if ( + lease !== null && lease.present && + lease.idempotency_key === command.idempotency_key + ) { + return acquireDecisionResult("rejected", "idempotency_key_reuse"); + } + const conflictIndexes = input.other_leases.flatMap((other, index) => + other.active && other.effective && + writeScopesOverlap(command.write_scopes, other.write_scopes) + ? [index] + : [] + ); + if (conflictIndexes.length > 0) { + return acquireDecisionResult("conflict", "write_scope_conflict", { + conflictIndexes, + }); + } + if (actualVersion >= Number.MAX_SAFE_INTEGER || (lease?.lease_epoch ?? 0) >= Number.MAX_SAFE_INTEGER) { + return acquireDecisionResult("rejected", "lease_generation_exhausted"); + } + return acquireDecisionResult("apply", "lease_acquire", { + nextLease: { + present: true, + active: true, + status: "active", + owner: command.owner, + idempotency_key: command.idempotency_key, + version: actualVersion + 1, + lease_epoch: (lease?.lease_epoch ?? 0) + 1, + write_scopes: [...command.write_scopes], + acquire_ttl_seconds: command.ttl_seconds, + }, + }); +} + +function equalScopeSets(left: readonly string[], right: readonly string[]): boolean { + const a = [...new Set(left)].sort((first, second) => first.localeCompare(second)); + const b = [...new Set(right)].sort((first, second) => first.localeCompare(second)); + return a.length === b.length && a.every((value, index) => value === b[index]); +} + +/** Materialize the one admitted generation for standalone and atomic Todo claim. */ +export function materializeTaskLeaseAcquire(identity: {goal_id: string; todo_id: string}, + command: AcquireDecisionInput["command"], decision: AcquireDecision, now: Date): LeaseRecord { + if (decision.outcome !== "apply" || decision.next_lease === null) { + throw new EffectRuntimeRequestError("lease materialization requires an apply decision"); + } + const at = now.toISOString().replace(/\.\d{3}Z$/u, "Z"); + return {schema_version: "task_lease_v0", goal_id: identity.goal_id, todo_id: identity.todo_id, + owner: command.owner, idempotency_key: command.idempotency_key, + write_scopes: [...command.write_scopes], acquire_ttl_seconds: command.ttl_seconds, + version: decisionInteger(decision.next_lease.version, "next_lease.version"), + lease_epoch: decisionInteger(decision.next_lease.lease_epoch, "next_lease.lease_epoch"), + acquired_at: at, updated_at: at, + expires_at: new Date(now.valueOf() + command.ttl_seconds * 1000).toISOString().replace(/\.\d{3}Z$/u, "Z"), + status: "active"}; +} diff --git a/scripts/generate_coordination_state_contract.py b/scripts/generate_coordination_state_contract.py index ae761b31a6..e30a0dc15d 100644 --- a/scripts/generate_coordination_state_contract.py +++ b/scripts/generate_coordination_state_contract.py @@ -107,6 +107,7 @@ ) TASK_LEASE_PROTOCOL_KEYS = ( "acquire_request_schema", + "canonical_acquire_request_schema", "lifecycle_request_schema", "canonical_renew_request_schema", "canonical_lifecycle_request_schema", diff --git a/tests/control_plane/test_canonical_lease_acquire.py b/tests/control_plane/test_canonical_lease_acquire.py new file mode 100644 index 0000000000..9fe9b3ee68 --- /dev/null +++ b/tests/control_plane/test_canonical_lease_acquire.py @@ -0,0 +1,84 @@ +"""Acquire must grant current execution proof from the selected authority.""" +import json +import os +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"]) +def test_public_acquire_renew_complete_and_retired_retry(tmp_path, monkeypatch, provider): + isolate_sqlite_runtime(tmp_path, monkeypatch) + runtime, state, registry = tmp_path / "runtime", tmp_path / "state.md", tmp_path / "registry.json" + goal, target = "lease-admission", "todo_admission" + state.write_text("# Synthetic admission\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": "Acquire and finish canonical work", "archive_state": "active", "source_section": "Agent Todo", + "index": 1, "task_class": "advancement_task", "claimed_by": "agent-a", + }]) + initialize_canonical_authority(runtime, goal, projection, state_path=state, provider=provider) + state.unlink() + cli_repo = Path(os.environ.get("LOOPX_LEASE_REPLAY_REPO", REPO)) + + def cli(command, *args, expected_exit=0): + child = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), "--format", "json", + *command, "--goal-id", goal, "--todo-id", target, *args], cwd=cli_repo, + capture_output=True, text=True, timeout=60, check=False) + assert child.returncode == expected_exit, child.stdout + child.stderr + return json.loads(child.stdout) + + proof = ["--owner", "agent-a", "--idempotency-key", "admission-a"] + acquire = [*proof, "--expected-version", "0", "--ttl-seconds", "600", "--write-scope", "src/**"] + try: + first = cli(["task-lease", "acquire"], *acquire) + assert first["ok"] and first["acquired"] + assert first["source_authority"] == provider + "_v0" + assert first["lease"]["version"] == first["lease"]["lease_epoch"] == 1 + assert first["lease"]["write_scopes"] == ["src/**"] + replay = cli(["task-lease", "acquire"], *acquire) + assert replay["idempotent"] and not replay["acquired"] + assert replay["lease"] == first["lease"] + assert replay["original_receipt"] == first["original_receipt"] + renewed = cli(["task-lease", "renew"], *proof, "--expected-version", "1", "--ttl-seconds", "600") + current = cli(["task-lease", "acquire"], *acquire) + assert current["lease"] == renewed["lease"] + assert current["original_receipt"] == first["original_receipt"] + retired = cli(["task-lease", "release"], *proof, "--expected-version", "2") + assert retired["released"] + rejected = cli(["task-lease", "acquire"], *acquire, expected_exit=1) + assert rejected["error_code"] == "idempotency_key_reuse" + inspected = cli(["task-lease", "inspect"]) + assert not inspected["active"] and inspected["lease"] == retired["lease"] + assert not (runtime / "goals" / goal / "task-leases" / f"{target}.json").exists() + assert not state.exists() + next_proof = ["--owner", "agent-a", "--idempotency-key", "admission-next", "--expected-version", "2"] + next_execution = cli(["task-lease", "acquire"], *next_proof, "--ttl-seconds", "600") + assert next_execution["lease"]["version"] == 3 and next_execution["lease"]["lease_epoch"] == 2 + complete_args = ["--agent-id", "agent-a", "--task-lease-idempotency-key", "admission-next", + "--task-lease-expected-version", "3", "--evidence", "validation://lease-admission", "--no-follow-up"] + preview = cli(["todo", "complete"], *complete_args, "--dry-run") + assert preview["ok"] + assert not state.exists() + completed = cli(["todo", "complete"], *complete_args) + assert completed["ok"] + assert state.exists() + readback = cli(["todo", "list"]) + assert readback["todos"][0]["status"] == "done" + terminal = cli(["task-lease", "inspect"]) + assert terminal["lease"]["status"] == "released" and not terminal["active"] + assert (cli(["task-lease", "acquire"], *next_proof, "--ttl-seconds", "600", expected_exit=1))["error_code"] == "todo_not_open" + 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=cli_repo, capture_output=True, text=True, timeout=30, check=True) diff --git a/tests/control_plane/test_shadow_fence_caller_parity_e2e.py b/tests/control_plane/test_shadow_fence_caller_parity_e2e.py index 5e489c1db5..ae1d489bd3 100644 --- a/tests/control_plane/test_shadow_fence_caller_parity_e2e.py +++ b/tests/control_plane/test_shadow_fence_caller_parity_e2e.py @@ -27,7 +27,7 @@ "leases=json.loads(sys.argv[2])); " "value['handoff_mode']='hard_lease'; print(json.dumps(value))" ) -PLACEHOLDERS = ("runtime_root", "todo_a", "todo_b", "todo_gate") +PLACEHOLDERS = ("runtime_root", "todo_a", "todo_b", "todo_c", "todo_gate") class Workspace: @@ -60,7 +60,7 @@ def fence_path(self) -> Path: def normalize(self, value: object) -> object: text = json.dumps(value) text = text.replace(json.dumps(str(self.w.root))[1:-1], "{runtime_root}") - for key in ("todo_gate", "todo_b", "todo_a"): + for key in ("todo_gate", "todo_c", "todo_b", "todo_a"): if key in self.ids: text = text.replace(self.ids[key], "{" + key + "}") return json.loads(text) @@ -100,6 +100,7 @@ def build_seeded(path: Path, mode: str, name: str, *, gate: bool) -> Workspace: assert ws.call("handoff-mode", "set", "--mode", "hard_lease")["ok"] is True ws.ids["todo_a"] = ws.w.add("Parity target A") ws.ids["todo_b"] = ws.w.add("Parity leased B") + ws.ids["todo_c"] = ws.w.add("Parity fresh acquisition C") if gate: gate_add = ws.call("todo", "add", "--role", "user", "--task-class", "user_gate", "--global-gate", "--text", "Approve the parity plan") @@ -111,7 +112,7 @@ def build_seeded(path: Path, mode: str, name: str, *, gate: bool) -> Workspace: lease = ws.control["cli-task_lease_acquire-absent"]["envelope"] assert lease.get("acquired") is True, lease ws.lease_version = int(lease["lease"]["version"]) - seed_and_fence(ws, ("todo_a", "todo_b", "todo_gate") if gate else ("todo_a", "todo_b")) + seed_and_fence(ws, ("todo_a", "todo_b", "todo_c", "todo_gate") if gate else ("todo_a", "todo_b", "todo_c")) return ws @@ -146,7 +147,7 @@ def build_w5(path: Path) -> Workspace: def lease_args(ws: Workspace, verb: str) -> tuple[str, ...]: todo_b, version = ws.ids.get("todo_b", ""), str(ws.lease_version) if verb == "acquire": - return ("task-lease", "acquire", "--todo-id", ws.ids.get("todo_a", ""), "--owner", "agent-a", + return ("task-lease", "acquire", "--todo-id", ws.ids.get("todo_c", ws.ids.get("todo_a", "")), "--owner", "agent-a", "--idempotency-key", "parity-acquire", "--ttl-seconds", "3600") common = ("task-lease", verb, "--todo-id", todo_b, "--owner", ws.lease_owner, "--idempotency-key", ws.lease_key, "--expected-version", version) @@ -277,6 +278,25 @@ def assert_canonical_lease_transition(ws: Workspace, caller: str, before: dict, assert native(ws.w, "read", {}) == canonical_after +def assert_canonical_acquisition(ws: Workspace, before: dict, observed: dict) -> None: + """New execution C must not change the existing B proof or keyless A tests.""" + after = native(ws.w, "read", {}) + old, current = before["head"], after["head"] + assert int(current["cursor"]) == int(old["cursor"]) + 1 + assert current["head"]["todos"] == old["head"]["todos"] + target = ws.ids["todo_c"] + assert [r for r in current["head"]["leases"] if r["todo_id"] != target] == old["head"]["leases"] + lease = next(r for r in current["head"]["leases"] if r["todo_id"] == target) + assert lease["owner"] == "agent-a" and lease["idempotency_key"] == "parity-acquire" + assert lease["version"] == lease["lease_epoch"] == 1 and lease["status"] == "active" + assert observed["envelope"]["lease"] == ws.normalize(lease) + assert "lease_path" not in observed["envelope"] + replay = ws.call(*lease_args(ws, "acquire")) + assert replay["ok"] and replay["idempotent"] and not replay["acquired"] + assert replay["lease"] == lease + assert native(ws.w, "read", {}) == after + + @pytest.mark.parametrize("row", load_rows(), ids=[row["id"] for row in load_rows()]) def test_fence_caller_parity(workspaces: Callable[[str], Workspace], row: dict) -> None: if row["id"] == "fixture-missing": @@ -284,7 +304,7 @@ def test_fence_caller_parity(workspaces: Callable[[str], Workspace], row: dict) ws = workspaces(row["workspace"]) canonical_before = native(ws.w, "read", {}) if row["caller"] in { "task_lease_renew", "task_lease_transfer", "task_lease_release", - } else None + } or (row["caller"] == "task_lease_acquire" and row["fence_state"] == "engaged") else None observed = observe_row(ws, row) assert observed["exit"] == row["exit"], observed if row.get("match") == "subset": @@ -298,7 +318,10 @@ def test_fence_caller_parity(workspaces: Callable[[str], Workspace], row: dict) assert observed["envelope"].get("provider_revision"), observed if canonical_before is not None: # Public canonical callers change the provider; legacy wire rows stay fenced. - assert_canonical_lease_transition(ws, row["caller"], canonical_before, observed) + if row["caller"] == "task_lease_acquire": + assert_canonical_acquisition(ws, canonical_before, observed) + else: + assert_canonical_lease_transition(ws, row["caller"], canonical_before, observed) assert observed["effect"] == row["effect"], observed assert observed["outbox_added"] == row.get("outbox_added", []), observed diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 728b00d729..ef1cb020b6 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,3 +1,4 @@ +import {registerLeaseAcquisitionConformance} from "./lease_acquisition_conformance.ts"; import {registerLeaseLifecycleConformance} from "./lease_lifecycle_conformance.ts"; import {registerMonitorConfigurationConformance} from "./monitor_configuration_conformance.ts"; import {registerAuthorityScanConformance} from "./authority_scan_conformance.ts"; @@ -211,6 +212,7 @@ export function registerAuthorityStoreConformance( factory: AuthorityStoreConformanceFactory, ): void { registerLeaseLifecycleConformance(providerName, factory); + registerLeaseAcquisitionConformance(providerName, factory); registerAuthorityScanConformance(providerName, factory); registerOwnershipObservationConformance(providerName, factory); registerNativePlanningUpdateConformance(providerName, factory); diff --git a/tests/control_plane_ts/canonical_task_lease_renew.test.ts b/tests/control_plane_ts/canonical_task_lease_renew.test.ts index d67ff0744b..b6a78fed20 100644 --- a/tests/control_plane_ts/canonical_task_lease_renew.test.ts +++ b/tests/control_plane_ts/canonical_task_lease_renew.test.ts @@ -1,3 +1,4 @@ +import {executeTaskLeaseAcquire} from "../../loopx/control_plane/work_items/task_lease_acquire.ts"; import assert from "node:assert/strict"; import {createHash} from "node:crypto"; import {mkdtemp, writeFile, rm, access} from "node:fs/promises"; @@ -18,7 +19,7 @@ import {executeCanonicalTaskLeaseLifecycle} from "../../loopx/control_plane/coor import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; import {taskLeaseOperationIdentity, taskLeaseOperationRequestDigest} from "../../loopx/control_plane/work_items/task_lease_operation_identity.ts"; import {legacyCoordinationWriterFencePath} from "../../loopx/control_plane/coordination/legacy_writer_fence.ts"; -import {TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/coordination_state_contract.generated.ts"; +import {TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA, TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_LIFECYCLE_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/coordination_state_contract.generated.ts"; import {shadowManagementStatePath} from "../../loopx/control_plane/coordination/shadow_management.ts"; import {atomicWriteJson} from "../../loopx/control_plane/effect_runtime_io.ts"; @@ -73,6 +74,57 @@ for (const provider of ["file", "sqlite"] as const) { const skip = provider === "sqlite" ? sqliteSkipReason() : undefined; const providerTest = (name: string, body: (t: TestContext) => Promise) => test(name, {skip}, body); + function acquireRequest(request: Awaited>["request"]) { + return {schema_version: TASK_LEASE_CANONICAL_ACQUIRE_REQUEST_SCHEMA, runtime_root: request.runtime_root, + goal_id: request.goal_id, todo_id: request.todo_id, owner: "agent-a", idempotency_key: "acquire-new", + expected_version: 1, ttl_seconds: 600, write_scopes: ["src/**"], authority: request.authority}; + } + const acquireNow = new Date("2026-09-13T10:11:00Z"); + providerTest(`${provider} native acquisition ignores stale Todo/mode and preserves the legacy fence`, async t => { + const {store, request, state} = await fixture(t, provider); + const command = acquireRequest(request); + const before = await store.loadAuthority(); + const legacy = await executeTaskLeaseAcquire({...command, schema_version: TASK_LEASE_ACQUIRE_REQUEST_SCHEMA}, {now: () => acquireNow}); + assert.equal(legacy.error_code, "legacy_coordination_writer_fenced"); + assert.deepEqual(await store.loadAuthority(), before); + const result = await executeTaskLeaseAcquire({...command, authority: {...command.authority, + handoff_mode: "soft_claim", todos: []}}, {now: () => acquireNow}); + assert.equal(result.ok, true, JSON.stringify(result)); assert.equal(result.source_authority, provider + "_v0"); + assert.equal(result.legacy_fallback_used, false); assert.equal(result.lease_path, undefined); + assert.equal((result.lease as Record).version, 2); + assert.equal((result.lease as Record).lease_epoch, 8); + await access(state); + }); + for (const fault of ["source", "fence", "field"] as const) { + providerTest(`${provider} acquisition fails closed on ${fault} without legacy or canonical writes`, async t => { + const {store, request, state, runtime, goal} = await fixture(t, provider); + const command = acquireRequest(request), before = await store.loadAuthority(); + if (fault === "fence") await rm(legacyCoordinationWriterFencePath(runtime, goal)); + const result = await executeTaskLeaseAcquire(fault === "field" ? {...command, lock_token: "wrong-domain"} : command, + {now: () => acquireNow, beforeWrite: fault === "source" ? async () => {await writeFile(state, "changed source");} : undefined}); + assert.equal(result.error_code, {source: "authority_source_changed", fence: "canonical_acquire_fence_missing", + field: "invalid_canonical_acquire_request"}[fault]); + assert.deepEqual(await store.loadAuthority(), before); + }); + } + for (const boundary of ["before", "after"] as const) { + providerTest(`${provider} real process death ${boundary} acquire commit preserves one new generation`, async t => { + const {root, store, request} = await fixture(t, provider); + const command = acquireRequest(request), config = join(root, "acquire-child.json"); + await writeFile(config, JSON.stringify({...command, operation: "acquire", now: acquireNow.toISOString()})); + const before = await store.loadAuthority(); + const child = spawnSync(process.execPath, ["--no-warnings", "--experimental-sqlite", "--experimental-strip-types", CHILD, config, boundary, "600"], {encoding: "utf8", timeout: 30000}); + assert.equal(child.signal, "SIGKILL", child.stderr); + if (boundary === "before") assert.deepEqual(await store.loadAuthority(), before); + const result = await executeTaskLeaseAcquire(command, {now: () => acquireNow}); + assert.equal(result.ok, true, JSON.stringify(result)); assert.equal(result.status, boundary === "before" ? "applied" : "replayed"); + const final = await store.loadAuthority(); if (final.status !== "loaded") throw new Error("missing head"); + assert.equal(final.cursor, "2"); assert.equal((result.lease as Record).version, 2); + assert.equal((result.lease as Record).lease_epoch, 8); + assert.equal((await executeTaskLeaseAcquire(command, {now: () => acquireNow})).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), final); + }); + } providerTest(`${provider} legacy wire requests remain fenced`, async t => { const {store, request} = await fixture(t, provider); const before = await store.loadAuthority(); const result = await executeTaskLeaseLifecycle({...request, schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION}, {now: () => NOW}); diff --git a/tests/control_plane_ts/canonical_task_lease_renew_process.ts b/tests/control_plane_ts/canonical_task_lease_renew_process.ts index f14df6fc2f..4afd0c6b82 100644 --- a/tests/control_plane_ts/canonical_task_lease_renew_process.ts +++ b/tests/control_plane_ts/canonical_task_lease_renew_process.ts @@ -1,4 +1,5 @@ -/** Disposable child for renewal race and lost-response integration tests. */ +import {executeCanonicalTaskLeaseAcquire} from "../../loopx/control_plane/coordination/task_lease_acquire.ts"; +/** Disposable child for acquisition/lifecycle crash and race qualification. */ import {readFileSync} from "node:fs"; import {openLocalAuthorityStore} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; import {executeCanonicalTaskLeaseLifecycle} from "../../loopx/control_plane/coordination/task_lease_lifecycle.ts"; @@ -24,7 +25,8 @@ const measured: AuthorityStore = { return result; }, }; -const result = await executeCanonicalTaskLeaseLifecycle(measured, {...request, +const execute = request.operation === "acquire" ? executeCanonicalTaskLeaseAcquire : executeCanonicalTaskLeaseLifecycle; +const result = await execute(measured, {...request, registered_agents: ["agent-a", "agent-b"], now: new Date(request.now), ttl_seconds: request.operation === "release" ? null : Number(ttl)}); process.stdout.write(JSON.stringify(result)); if (process.connected) process.disconnect(); diff --git a/tests/control_plane_ts/lease_acquisition_conformance.ts b/tests/control_plane_ts/lease_acquisition_conformance.ts new file mode 100644 index 0000000000..b43f6caaa8 --- /dev/null +++ b/tests/control_plane_ts/lease_acquisition_conformance.ts @@ -0,0 +1,174 @@ +/** Independent execution-admission invariants run unchanged on every store. */ +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 {executeCanonicalTaskLeaseAcquire as acquire, type CanonicalTaskLeaseAcquireInput} from "../../loopx/control_plane/coordination/task_lease_acquire.ts"; +import {executeCanonicalTaskLeaseLifecycle as mutate} from "../../loopx/control_plane/coordination/task_lease_lifecycle.ts"; +import {executeCoordinationTodoClaim} from "../../loopx/control_plane/coordination/todo_claim.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; +import {productionScaleLeaseAcquisitionFixture} from "./production_scale_coordination_fixture.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; + +async function loaded(store: AuthorityStore) { + const result = await store.loadAuthority(); + assert.equal(result.status, "loaded"); + if (result.status !== "loaded") throw new Error("missing fixture authority"); + return result; +} + +export function registerLeaseAcquisitionConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + async function setup(t: test.TestContext, schema: "native" | "legacy" = "native", + change?: (projection: JsonObject, target: string, conflict: string) => void) { + const {store, contender} = await factory(t); + const fixture = productionScaleLeaseAcquisitionFixture("goal-a", schema); + const projection = fixture.projection as JsonObject; + change?.(projection, fixture.target, fixture.acquisition.conflict_todo_id); + const seed = authorityProjectionFixture("goal-a", projection.todos as JsonObject[], projection.leases as JsonObject[], schema, + {handoff_mode: projection.handoff_mode}); + assert.equal((await store.commitAuthority({expected_provider_revision: null, operation_id: "admission-seed", + next_projection: seed, events: [], receipts: []})).status, "applied"); + const request: CanonicalTaskLeaseAcquireInput = {goal_id: "goal-a", todo_id: fixture.target, owner: "agent-a", + idempotency_key: fixture.acquisition.execution_key, expected_version: 0, + ttl_seconds: fixture.acquisition.ttl_seconds, write_scopes: fixture.acquisition.write_scopes, + registered_agents: fixture.registered_agents, now: new Date(fixture.scenario.now)}; + return {store, contender, request, seed, fixture}; + } + + for (const schema of ["native", "legacy"] as const) { + test(`${provider} ${schema} acquire, current-proof replay, handover and new execution`, async t => { + const {store, contender, request, seed} = await setup(t, schema); + const first = await acquire(store, request); + assert.equal(first.status, "applied", JSON.stringify(first)); + assert.deepEqual(first.lease, {schema_version: "task_lease_v0", goal_id: "goal-a", todo_id: request.todo_id, + owner: "agent-a", idempotency_key: "acquire-a", version: 1, lease_epoch: 1, status: "active", + write_scopes: ["lease-admission/**"], acquire_ttl_seconds: 600, + acquired_at: "2026-09-13T10:05:00Z", updated_at: "2026-09-13T10:05:00Z", expires_at: "2026-09-13T10:15:00Z"}); + const after = await loaded(contender); + const replay = await acquire(contender, request); + assert.equal(replay.status, "replayed"); assert.deepEqual(replay.lease, first.lease); + assert.deepEqual(replay.original_receipt, first.original_receipt); + assert.deepEqual(await loaded(store), after); + const maintenance = {goal_id: request.goal_id, todo_id: request.todo_id, owner: request.owner, + idempotency_key: request.idempotency_key, expected_version: 1, ttl_seconds: 600, + registered_agents: request.registered_agents, now: new Date("2026-09-13T10:06:00Z")}; + const renewed = await mutate(store, {...maintenance, operation: "renew"}); + assert.equal(renewed.status, "applied"); + const current = await acquire(contender, request); + assert.equal(current.status, "replayed"); assert.deepEqual(current.lease, renewed.lease); + assert.deepEqual(current.original_receipt, first.original_receipt); + const transferred = await mutate(store, {...maintenance, operation: "transfer", expected_version: 2, + new_owner: "agent-b", new_idempotency_key: "receiver-b"}); + assert.equal(transferred.status, "applied"); + const beforeStale = await loaded(store); + assert.equal((await acquire(contender, request)).reason_code, "idempotency_key_reuse"); + assert.deepEqual(await loaded(store), beforeStale); + const released = await mutate(store, {...maintenance, operation: "release", owner: "agent-b", + idempotency_key: "receiver-b", expected_version: 3, ttl_seconds: null}); + assert.equal(released.status, "applied"); + const next = await acquire(contender, {...request, owner: "agent-b", idempotency_key: "acquire-b", expected_version: 3}); + assert.equal(next.status, "applied"); + assert.equal((next.lease as JsonObject).version, 4); assert.equal((next.lease as JsonObject).lease_epoch, 3); + const final = await loaded(store); + assert.equal(final.cursor, "6"); assert.deepEqual(final.head.todos, seed.todos); + assert.deepEqual((final.head.leases as JsonObject[]).filter(l => l.todo_id !== request.todo_id), seed.leases); + }); + } + + test(`${provider} lost acquire response recovers once and changed intent stays rejected`, async t => { + const {store, contender, request} = await setup(t); + let commits = 0; + const transport: AuthorityStore = {storeIdentity: () => store.storeIdentity(), loadAuthority: () => store.loadAuthority(), + scanCommitted: (...a) => store.scanCommitted(...a), readReceipt: id => store.readReceipt(id), + commitAuthority: async commit => {commits++; assert.equal((await store.commitAuthority(commit)).status, "applied"); throw new Error("lost response");}}; + const recovered = await acquire(transport, request); + assert.equal(recovered.status, "recovered"); assert.equal(commits, 1); + const after = await loaded(contender); + for (const changed of [{ttl_seconds: 900}, {write_scopes: ["different/**"]}, {expected_version: 1}]) { + assert.equal((await acquire(contender, {...request, ...changed})).reason_code, "coordination_operation_identity_mismatch"); + assert.deepEqual(await loaded(store), after); + } + }); + + test(`${provider} competing acquisitions cannot grant two overlapping executions`, async t => { + const {store, contender, request} = await setup(t); + const loser = await acquire(store, request, async () => { + const winner = await acquire(contender, {...request, owner: "agent-b", idempotency_key: "winner-b"}); + assert.equal(winner.status, "applied"); + }); + assert.equal(loser.status, "conflict"); + const current = await loaded(store); + assert.equal(current.cursor, "2"); + assert.equal((current.head.leases as JsonObject[]).find(l => l.todo_id === request.todo_id)!.owner, "agent-b"); + assert.equal((await acquire(store, {...request, expected_version: null})).reason_code, "todo_lease_conflict"); + assert.deepEqual(await loaded(store), current); + }); + + for (const dimension of ["live", "archived", "excluded", "claim", "deregistered", "expired"] as const) { + test(`${provider} complete scope scan adjudicates ${dimension} holder beyond display limits`, async t => { + const {store, request, fixture} = await setup(t, "native", (p, _target, conflict) => { + const todo = (p.todos as JsonObject[]).find(r => r.todo_id === conflict)!; + const lease = (p.leases as JsonObject[]).find(r => r.todo_id === conflict)!; + lease.write_scopes = ["lease-admission/shared/**"]; + if (dimension === "archived") todo.archive_state = "archive"; + if (dimension === "excluded") todo.excluded_agents = ["agent-b"]; + if (dimension === "claim") todo.claimed_by = "agent-a"; + if (dimension === "deregistered") lease.owner = "unregistered-agent"; + if (dimension === "expired") lease.expires_at = "2026-09-13T10:00:00Z"; + }); + const before = await loaded(store); + const result = await acquire(store, request); + if (dimension === "live") { + assert.equal(result.reason_code, "write_scope_conflict"); + assert.equal((result.conflicts as JsonObject[])[0]!.todo_id, fixture.acquisition.conflict_todo_id); + assert.deepEqual(await loaded(store), before); + } else assert.equal(result.status, "applied", JSON.stringify(result)); + }); + } + + for (const [dimension, code] of [["mode", "handoff_mode_forbids_lease"], ["archive", "todo_not_found"], + ["done", "todo_not_open"], ["claim", "owner_conflicts_with_claim"], ["excluded", "owner_excluded_from_todo"]] as const) { + test(`${provider} acquisition rejects ${dimension} with no receipt or mutation`, async t => { + const {store, request} = await setup(t, "native", (p, target) => { + const todo = (p.todos as JsonObject[]).find(r => r.todo_id === target)!; + if (dimension === "mode") p.handoff_mode = "soft_claim"; + if (dimension === "archive") todo.archive_state = "archive"; + if (dimension === "done") {todo.status = "done"; todo.done = true;} + if (dimension === "claim") todo.claimed_by = "agent-b"; + if (dimension === "excluded") todo.excluded_agents = ["agent-a"]; + }); + const before = await loaded(store); + assert.equal((await acquire(store, request)).reason_code, code); + assert.deepEqual(await loaded(store), before); + }); + } + + for (const reason of ["expired", "ineffective"] as const) { + test(`${provider} ${reason} execution is replaced with a new version and epoch`, async t => { + const {store, request} = await setup(t); + assert.equal((await acquire(store, request)).status, "applied"); + const takeover = {...request, owner: "agent-b", idempotency_key: "next-b", expected_version: 1, + now: new Date(reason === "expired" ? "2026-09-14T10:05:00Z" : "2026-09-13T10:06:00Z"), + registered_agents: reason === "ineffective" ? ["agent-b"] : request.registered_agents}; + const next = await acquire(store, takeover); + assert.equal(next.status, "applied"); + assert.equal((next.lease as JsonObject).version, 2); assert.equal((next.lease as JsonObject).lease_epoch, 2); + assert.equal((await acquire(store, {...request, now: takeover.now})).reason_code, "idempotency_key_reuse"); + }); + } + + test(`${provider} atomic Todo claim uses the same archived-holder conflict rule`, async t => { + const {store, request} = await setup(t, "native", (p, target, conflict) => { + const todo = (p.todos as JsonObject[]).find(r => r.todo_id === target)!; + todo.required_write_scopes = ["lease-admission/**"]; + (p.todos as JsonObject[]).find(r => r.todo_id === conflict)!.archive_state = "archive"; + (p.leases as JsonObject[]).find(r => r.todo_id === conflict)!.write_scopes = ["lease-admission/shared/**"]; + }); + const result = await executeCoordinationTodoClaim(store, {goal_id: request.goal_id, todo_id: request.todo_id, + claimed_by: request.owner, actor_agent_id: request.owner, expected_role: "agent", registered_agents: request.registered_agents, + operation_id: "claim-acquire", lease_request: {idempotency_key: request.idempotency_key, expected_version: 0, ttl_seconds: 600}, + dry_run: false, now: request.now}); + assert.equal(result.status, "applied", JSON.stringify(result)); + assert.equal((result.lease as JsonObject).version, 1); + }); +} diff --git a/tests/control_plane_ts/production_scale_coordination_fixture.ts b/tests/control_plane_ts/production_scale_coordination_fixture.ts index 9aaf903c3e..27eb179517 100644 --- a/tests/control_plane_ts/production_scale_coordination_fixture.ts +++ b/tests/control_plane_ts/production_scale_coordination_fixture.ts @@ -44,6 +44,8 @@ const envelope = JSON.parse(readFileSync(new URL( }; lease_lifecycle: {owner: string; receiver: string; execution_key: string; receiver_execution_key: string; version: number; lease_epoch: number; now: string}; + lease_acquisition: {execution_key: string; next_execution_key: string; write_scopes: string[]; ttl_seconds: number; + conflict_todo_id: string; conflict_write_scopes: string[]}; semantic_cases: Record>; presentation_cases: Record>; update_cases: Record>; @@ -487,3 +489,21 @@ export function productionScaleLeaseLifecycleFixture(goalId: string, {source_authority: "synthetic_production_scale_fixture", handoff_mode: "hard_lease"}); return {projection, target, scenario, registered_agents: fixture.registered_agents}; } + +/** Start without a target lease, with a live peer beyond the bounded display. */ +export function productionScaleLeaseAcquisitionFixture(goalId: string, + schema: AuthorityProjectionSchema = "native") { + const fixture = productionScaleLeaseLifecycleFixture(goalId, schema); + const scenario = envelope.lease_acquisition; + const todos = fixture.projection.todos as Record[]; + const leases = fixture.projection.leases as Record[]; + const target = todos.find(todo => todo.todo_id === fixture.target)!; + const oldLease = leases.find(lease => lease.todo_id === fixture.target)!; + const projection = authorityProjectionFixture(goalId, + [...todos, {...target, todo_id: scenario.conflict_todo_id, text: "Independent scope holder beyond display limits"}], + [...leases.filter(lease => lease.todo_id !== fixture.target), {...oldLease, + todo_id: scenario.conflict_todo_id, owner: "agent-b", idempotency_key: "scope-holder", + write_scopes: ["independent-work/**"]}], schema, + {source_authority: "synthetic_production_scale_fixture", handoff_mode: "hard_lease"}); + return {...fixture, projection, acquisition: scenario}; +} diff --git a/tests/control_plane_ts/task_lease_acquire.test.ts b/tests/control_plane_ts/task_lease_acquire.test.ts index 1ccc7fb586..bb94af08bb 100644 --- a/tests/control_plane_ts/task_lease_acquire.test.ts +++ b/tests/control_plane_ts/task_lease_acquire.test.ts @@ -685,3 +685,19 @@ for (const { expiresAt, active } of [ } }); } + +for (const field of ["version", "lease_epoch"] as const) { + test(`legacy acquisition rejects exhausted ${field} before durable writeback`, async t => { + const root = await workspace(t); + const command = await request(root); + assert.equal((await executeTaskLeaseAcquire(command, {now: () => FIXED_NOW})).ok, true); + const lease = {...await persistedLease(root), status: "released", [field]: Number.MAX_SAFE_INTEGER}; + await writeFile(leasePath(root), JSON.stringify(lease)); + const before = await readFile(leasePath(root), "utf8"); + const result = await executeTaskLeaseAcquire({...command, idempotency_key: "next-execution"}, {now: () => FIXED_NOW}); + assert.equal(result.error_code, "lease_generation_exhausted"); + assert.deepEqual(result.settlement, {effect_id: null, receipts: [], failure: { + step: "validation", kind: "writeback_rejected", code: "lease_generation_exhausted"}}); + assert.equal(await readFile(leasePath(root), "utf8"), before); + }); +} diff --git a/tests/control_plane_ts/task_lease_eligibility.test.ts b/tests/control_plane_ts/task_lease_eligibility.test.ts index ddf06216cc..ba27947a2e 100644 --- a/tests/control_plane_ts/task_lease_eligibility.test.ts +++ b/tests/control_plane_ts/task_lease_eligibility.test.ts @@ -1,6 +1,6 @@ import assert from "node:assert/strict"; import test from "node:test"; -import {evaluateTaskLeaseAcquireDecision} from "../../loopx/control_plane/work_items/task_lease_acquire.ts"; +import {evaluateTaskLeaseAcquireDecision} from "../../loopx/control_plane/work_items/task_lease_acquire_decision.ts"; import {decideTaskLeaseLifecycle} from "../../loopx/control_plane/work_items/task_lease_lifecycle_decision.ts"; import {evaluateTaskLeaseOwnerEligibility, leaseOwnerRejection, type LeaseEligibilityTodo} from "../../loopx/control_plane/work_items/task_lease_eligibility.ts"; @@ -90,3 +90,13 @@ test("eligibility decoder rejects malformed facts and does not infer authorisati assert.equal(evaluateTaskLeaseOwnerEligibility({todo: described, owner: "agent-a", registered_agents: ["agent-a"]}).code, "owner_conflicts_with_claim"); }); + +for (const generation of ["version", "lease_epoch"] as const) { + test(`new execution rejects exhausted ${generation} rather than duplicating a numeric fence`, () => { + const input = acquire(false); + input.lease = {...input.lease, active: false, [generation]: Number.MAX_SAFE_INTEGER}; + input.command.expected_version = input.lease.version; + const result = evaluateTaskLeaseAcquireDecision(input); + assert.equal(result.code, "lease_generation_exhausted"); assert.equal(result.outcome, "rejected"); + }); +} diff --git a/tests/fixtures/control_plane/coordination_production_scale_v0.json b/tests/fixtures/control_plane/coordination_production_scale_v0.json index 128f2e24e9..598630b2ef 100644 --- a/tests/fixtures/control_plane/coordination_production_scale_v0.json +++ b/tests/fixtures/control_plane/coordination_production_scale_v0.json @@ -44,6 +44,12 @@ "lease_epoch": 7, "now": "2026-09-13T10:05:00Z" }, + "lease_acquisition": { + "execution_key": "acquire-a", "next_execution_key": "acquire-b", + "write_scopes": ["lease-admission/**"], "ttl_seconds": 600, + "conflict_todo_id": "todo_zz_acquire_scope", + "conflict_write_scopes": ["lease-admission/shared/**"] + }, "semantic_cases": { "title_only_monitor": { "todo_id": "todo_fixture_title_monitor", diff --git a/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json b/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json index 16c18c89c8..6093a190c5 100644 --- a/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json +++ b/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json @@ -1078,35 +1078,22 @@ "baseline": null }, { - "id": "cli-task_lease_acquire-engaged", + "id": "cli-task_lease_acquire-canonical", "surface": "cli", "workspace": "w1", "caller": "task_lease_acquire", "fence_state": "engaged", - "exit": 1, + "exit": 0, "expect": { - "ok": false, + "ok": true, "schema_version": "task_lease_v0", "action": "acquire", - "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", - "error_code": "legacy_coordination_writer_fenced", - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "blocked", - "reason_code": "legacy_coordination_writer_fenced", - "authority_mode": "file_v0", - "fence_id": "caller-fixture" - }, - "settlement": { - "effect_id": null, - "receipts": [], - "failure": { - "step": "validation", - "kind": "permission_denied", - "code": "legacy_coordination_writer_fenced" - } - } + "legacy_fallback_used": false, + "acquired": true, + "idempotent": false, + "source_authority": "file_v0", + "decision_read_from_provider": true, + "status": "applied" }, "effect": { "added": [], @@ -1127,8 +1114,33 @@ "error": "native task-lease acquire result shape mismatch", "error_code": "RuntimeError", "schema_version": "task_lease_v0" + }, + "8330a974c": { + "ok": false, + "schema_version": "task_lease_v0", + "action": "acquire", + "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", + "error_code": "legacy_coordination_writer_fenced", + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "blocked", + "reason_code": "legacy_coordination_writer_fenced", + "authority_mode": "file_v0", + "fence_id": "caller-fixture" + }, + "settlement": { + "effect_id": null, + "receipts": [], + "failure": { + "step": "validation", + "kind": "permission_denied", + "code": "legacy_coordination_writer_fenced" + } + } } - } + }, + "match": "subset" }, { "id": "cli-task_lease_renew-canonical", @@ -1717,35 +1729,22 @@ "match": "subset" }, { - "id": "cli-task_lease_acquire-engaged-capture", + "id": "cli-task_lease_acquire-canonical-capture", "surface": "cli", "workspace": "w2", "caller": "task_lease_acquire", "fence_state": "engaged", - "exit": 1, + "exit": 0, "expect": { - "ok": false, + "ok": true, "schema_version": "task_lease_v0", "action": "acquire", - "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", - "error_code": "legacy_coordination_writer_fenced", - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "blocked", - "reason_code": "legacy_coordination_writer_fenced", - "authority_mode": "file_v0", - "fence_id": "caller-fixture" - }, - "settlement": { - "effect_id": null, - "receipts": [], - "failure": { - "step": "validation", - "kind": "permission_denied", - "code": "legacy_coordination_writer_fenced" - } - } + "legacy_fallback_used": false, + "acquired": true, + "idempotent": false, + "source_authority": "file_v0", + "decision_read_from_provider": true, + "status": "applied" }, "effect": { "added": [], @@ -1753,7 +1752,33 @@ "changed": [] }, "outbox_added": [], - "baseline": null + "baseline": { + "8330a974c": { + "ok": false, + "schema_version": "task_lease_v0", + "action": "acquire", + "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", + "error_code": "legacy_coordination_writer_fenced", + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "blocked", + "reason_code": "legacy_coordination_writer_fenced", + "authority_mode": "file_v0", + "fence_id": "caller-fixture" + }, + "settlement": { + "effect_id": null, + "receipts": [], + "failure": { + "step": "validation", + "kind": "permission_denied", + "code": "legacy_coordination_writer_fenced" + } + } + } + }, + "match": "subset" }, { "id": "cli-task_lease_renew-canonical-capture", @@ -1901,16 +1926,12 @@ "ok": false, "schema_version": "task_lease_v0", "action": "acquire", - "error": "legacy coordination writer fence must be engaged", + "legacy_fallback_used": false, "error_code": "legacy_writer_fence_read_failed", - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "failed", - "reason_code": "legacy_writer_fence_read_failed", - "reason": "legacy coordination writer fence must be engaged", - "authority_mode": "unknown_fail_closed" - }, + "error": "legacy coordination writer fence must be engaged", + "source_authority": null, + "decision_read_from_provider": false, + "status": "failed", "settlement": { "effect_id": null, "receipts": [], @@ -1927,7 +1948,33 @@ "changed": [] }, "outbox_added": [], - "baseline": null + "baseline": { + "8330a974c": { + "ok": false, + "schema_version": "task_lease_v0", + "action": "acquire", + "error": "legacy coordination writer fence must be engaged", + "error_code": "legacy_writer_fence_read_failed", + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "failed", + "reason_code": "legacy_writer_fence_read_failed", + "reason": "legacy coordination writer fence must be engaged", + "authority_mode": "unknown_fail_closed" + }, + "settlement": { + "effect_id": null, + "receipts": [], + "failure": { + "step": "validation", + "kind": "permission_denied", + "code": "legacy_writer_fence_read_failed" + } + } + } + }, + "match": "subset" }, { "id": "cli-todo_capture_followups-invalid", @@ -1973,16 +2020,12 @@ "ok": false, "schema_version": "task_lease_v0", "action": "acquire", - "error": "EISDIR: illegal operation on a directory, read", + "legacy_fallback_used": false, "error_code": "legacy_writer_fence_read_failed", - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "failed", - "reason_code": "legacy_writer_fence_read_failed", - "reason": "EISDIR: illegal operation on a directory, read", - "authority_mode": "unknown_fail_closed" - }, + "error": "EISDIR: illegal operation on a directory, read", + "source_authority": null, + "decision_read_from_provider": false, + "status": "failed", "settlement": { "effect_id": null, "receipts": [], @@ -1999,7 +2042,33 @@ "changed": [] }, "outbox_added": [], - "baseline": null + "baseline": { + "8330a974c": { + "ok": false, + "schema_version": "task_lease_v0", + "action": "acquire", + "error": "EISDIR: illegal operation on a directory, read", + "error_code": "legacy_writer_fence_read_failed", + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_a}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "failed", + "reason_code": "legacy_writer_fence_read_failed", + "reason": "EISDIR: illegal operation on a directory, read", + "authority_mode": "unknown_fail_closed" + }, + "settlement": { + "effect_id": null, + "receipts": [], + "failure": { + "step": "validation", + "kind": "permission_denied", + "code": "legacy_writer_fence_read_failed" + } + } + } + }, + "match": "subset" }, { "id": "cli-todo_capture_followups-unreadable", @@ -2042,15 +2111,14 @@ "fence_state": "engaged_without_store", "exit": 1, "expect": { + "ok": false, "schema_version": "task_lease_v0", - "status": "missing", + "action": "acquire", + "legacy_fallback_used": false, + "error_code": "canonical_lease_authority_unavailable", "source_authority": "file_v0", "decision_read_from_provider": true, - "legacy_fallback_used": false, - "ok": false, - "action": "acquire", - "error": "canonical Todo authority is unavailable", - "error_code": "local_authority_todo_list_unavailable" + "status": "missing" }, "effect": { "added": [], @@ -2059,8 +2127,20 @@ }, "outbox_added": [], "baseline": { - "note": "ee1b17217 printed schema_version loopx_local_coordination_todo_list_result_v0 because the exception payload was spread after the envelope keys" - } + "note": "ee1b17217 printed schema_version loopx_local_coordination_todo_list_result_v0 because the exception payload was spread after the envelope keys", + "8330a974c": { + "schema_version": "task_lease_v0", + "status": "missing", + "source_authority": "file_v0", + "decision_read_from_provider": true, + "legacy_fallback_used": false, + "ok": false, + "action": "acquire", + "error": "canonical Todo authority is unavailable", + "error_code": "local_authority_todo_list_unavailable" + } + }, + "match": "subset" } ] }