diff --git a/docs/architecture/rfcs/STATUS.md b/docs/architecture/rfcs/STATUS.md index 9e46348e84..c16dd68b3a 100644 --- a/docs/architecture/rfcs/STATUS.md +++ b/docs/architecture/rfcs/STATUS.md @@ -63,7 +63,7 @@ appendix may keep dated history, but no dated log heading may precede it. | [RFC: Shared Goal Alignment and Governed Amendment Protocol (v0)](shared-goal-alignment-and-governed-amendment-v0.md) | Accepted | none | [2 entries](ledger/shared-goal-alignment-and-governed-amendment-v0/) | | [RFC: LoopX Shared Control-Plane Authority and Pluggable State Providers (v0)](shared-goal-authority-state-provider-v0.md) | Accepted | none | [23 entries](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | Accepted | none | — | -| [RFC: TypeScript Control-Plane Migration Direction v0](typescript-control-plane-migration-v0.md) | Accepted | none | [13 entries](ledger/typescript-control-plane-migration-v0/) | +| [RFC: TypeScript Control-Plane Migration Direction v0](typescript-control-plane-migration-v0.md) | Accepted | none | [14 entries](ledger/typescript-control-plane-migration-v0/) | ## Superseded (0) diff --git a/docs/architecture/rfcs/STATUS.zh-CN.md b/docs/architecture/rfcs/STATUS.zh-CN.md index c4688126e2..fe9189d5b8 100644 --- a/docs/architecture/rfcs/STATUS.zh-CN.md +++ b/docs/architecture/rfcs/STATUS.zh-CN.md @@ -60,7 +60,7 @@ | [RFC:共享 Goal 对齐与受治理 Amendment 协议(v0)](shared-goal-alignment-and-governed-amendment-v0.zh-CN.md) | 已接受 | 无 | [2 条](ledger/shared-goal-alignment-and-governed-amendment-v0/) | | [RFC:LoopX 共享控制面权威与可插拔状态 Provider(v0)](shared-goal-authority-state-provider-v0.zh-CN.md) | 已接受 | 无 | [23 条](ledger/shared-goal-authority-state-provider-v0/) | | [RFC: Single-Owner Local Daemon (v0)](single-owner-local-daemon-v0.md) | 已接受 | none | — | -| [RFC:LoopX 控制面 TypeScript 渐进迁移方向 v0](typescript-control-plane-migration-v0.zh-CN.md) | 已接受 | 无 | [13 条](ledger/typescript-control-plane-migration-v0/) | +| [RFC:LoopX 控制面 TypeScript 渐进迁移方向 v0](typescript-control-plane-migration-v0.zh-CN.md) | 已接受 | 无 | [14 条](ledger/typescript-control-plane-migration-v0/) | ## 已被替代 (0) diff --git a/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.md b/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.md index 9d4d03200b..a3785900e9 100644 --- a/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.md +++ b/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.md @@ -206,6 +206,9 @@ profile's promotion status. Use existing domain tasks and roadmap checkpoints for execution; this plan does not require unrelated migrations before a bounded repair can ship. +The [quiet-to-due Monitor checkpoint](ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.md) +records bounded CLI replay/successor evidence and its explicit M2/M3 limits. + ## 12. Open decisions The first M2 implementation must choose its smallest sufficient exploration diff --git a/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.zh-CN.md b/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.zh-CN.md index 3bbfda824d..88d59bf2d7 100644 --- a/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.zh-CN.md +++ b/docs/architecture/rfcs/composable-state-machines-recovery-verification-v0.zh-CN.md @@ -167,6 +167,9 @@ provider 会使相应前提不成立;必须暴露这个事实,不能判定 执行沿用领域任务与总路线 checkpoint,不要求无关迁移先于有界修复完成。 +[Monitor 静默到到期 checkpoint](ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.zh-CN.md) +记录有界 CLI 重放/successor 证据,并明确 M2/M3 尚未覆盖的部分。 + ## 12. 未决事项 首个 M2 实现由测试与领域维护者选择足够小的探索方法及边界。优先复用 fixture diff --git a/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.md b/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.md new file mode 100644 index 0000000000..bad4638e04 --- /dev/null +++ b/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.md @@ -0,0 +1,125 @@ +# Quiet-to-due Monitor recovery qualification + +[中文镜像](2026-10-02-monitor-quiet-due-recovery.zh-CN.md) + +This checkpoint adds composition evidence for T2 and the +[recovery verification RFC](../../composable-state-machines-recovery-verification-v0.md). +Base `e38b057b3` already +contains the effect-identity repair from [#4335](https://github.com/loopx-project/loopx/pull/4335). +An older installed runtime can still exhibit that repaired defect. +The remaining repair handles an unpolled, bound Monitor that becomes blocked: +its original Turn now exposes conditional lifecycle recovery instead of ordinary +execution or conflicting replan selection. No new RPC or persisted schema is added. + +## Missing counterexample + +A quiet heartbeat automatically commits an observation without a Todo. If a +Monitor becomes due during that same Turn, its explicit binding and actual +observation must remain possible. The first observation cannot settle the later +Todo or occupy its effect identity. Testing a newly due Monitor after an ordinary +unbound guard misses this prerequisite: the quiet poll has already committed. + +```mermaid +flowchart LR + Q[Quiet guard: no Todo] --> O[Commit quiet observation] + O --> D[Monitor becomes due] + D --> B[Bind original Turn to Monitor] + B --> P[Commit exact poll once] + P --> L[Discard caller response] + L --> R[Read back and retry original poll] + R --> S[Turn settled; no quota debit] + S --> N[Material successor uses a new Turn] +``` + +## Executable boundary + +`tests/control_plane/test_monitor_quiet_due_recovery.py` enumerates twelve journeys: +legacy Markdown, canonical File and canonical SQLite, each with unchanged and +material observations, with and without a peer-scoped user gate and reminder. Each journey uses real CLI subprocesses, real TS effects +and isolated provider state. It performs two exact retries and one conflicting +result retry. The response is deliberately discarded after successful command +completion; this models acknowledgement loss, not a process crash inside commit. + +Independent assertions require two distinct polls (quiet and Todo-bound), +unchanged original observation, no refresh/spend records, no mutation on replay +or conflict, one material generation and successor only in the material case, +and settled readback that cannot execute further work in the original Turn. +The Monitor remains open; material successor selection uses a fresh Turn. + +Five additional journeys block an already-bound Monitor before its poll: legacy, +File and SQLite soft claims, plus File and SQLite hard leases. The +unmodified base incorrectly returns `normal_run` in this fixture; other frontier +states can attempt a conflicting replan binding. Head returns the existing +`unsettled_host_turn_recovery` mode with the original identity and no delivery +authority. Two readbacks preserve blocked state. Only after the fixture's blocker +is resolved does the test execute the projected restore command, re-enter the +original guard, poll and settle without spending quota. The projected restore +carries the independently verified reason and `--clear-resume-when`, as required +by the existing lifecycle owner. Hard-lease restoration grants no execution +authority: polling without a lease remains rejected, and a fresh lease is acquired +before the exact poll. Two additional File/SQLite cases execute the same projected +restore against an active holder and verify rejection with Todo, lease and Goal +state unchanged. + +Sensitivity was checked in a disposable checkout of the same base: restoring +Turn-only effect allocation makes the unchanged/legacy journey fail at its first +bound poll with `heartbeat_receipt_identity_conflict`. Unmodified base passes +all six journeys. This is a deliberate historical-rule mutation, not a claim +that the current base fails. The temporary mutant is not a shipped fixture. + +An additional counterexample supplies an incomplete Todo frontier with no aggregate +work lane and an unrelated peer gate. The exact settled Monitor still has to +return `heartbeat_settled_skip`. Base instead allows ordinary execution because +the adapter drops its typed phase when the aggregate lane is absent. The adapter +now preserves the bound Monitor projection independently of aggregate coverage; +TS still owns poll verification and phase classification. This regression is red +before the adapter change and green after it. + +A further composition case settles an advancement Turn, then moves its primary +Todo to `blocked` before an independent due Monitor is observed. +The emitted poll previously failed because admission accepted only open/done +primaries. The existing TS transaction owner now accepts this later hold only +when exact settlement readback is `settled`. Unsettled holds still fail closed; +the Monitor must independently pass its due, owner, gate and lease checks. +No held Todo is reopened, no new advancement is admitted and no extra debit is +created. Synthetic TS regressions cover the accepted hold and identity/gate +rejections (including deferred primaries); one CLI journey exercise in-flight writeback, spend, hold, poll and +exact replay. The accepted blocked-primary TS case fails before the change. + +## Ownership and limits + +| Boundary | Existing owner retained | +| --- | --- | +| Selection arbitration | `work_items/action_portfolio.ts` | +| Monitor transaction and immutable replay | `quota/monitor_poll_commit.ts` | +| Observation and successors in canonical authority | `coordination/todo_monitor_poll.ts` | +| Turn closeout readback | `quota/settlement_readback.ts`, `quota/settlement_phase.ts` | +| Bound Monitor lifecycle recovery | `quota/blocked_wait.ts` | +| Bound phase projection into an optional aggregate lane | `work_items/work_lane.py` | +| Legacy effect-id compatibility and transport | `quota/monitor_poll.py` | + +The related refactor shares the current-Turn recovery envelope between causal +waits and blocked Monitor recovery in `blocked_wait.ts`. Python passes the already +verified Monitor phase through the existing request and renders the typed repair; +it owns no second lifecycle rule. The added phase field is optional: earlier +requests retain their causal-wait behavior. Only active, correctly owned blocked +Monitors with `poll_due` qualify; missing, duplicate, foreign, archived and already +polled inputs do not enter this restoration route. Tests assert these exclusions. + +Default behavior changes for blocked, unpolled replay, for bound Monitor +projection when an incomplete frontier has no aggregate lane, and for independent +auxiliary observation after an exact settled advancement is held. The emitted +restore command is conditional on a verified resolved blocker; the projection +neither reopens the Todo nor establishes a poll/closeout receipt. Retain the +blocker when unresolved. The existing Todo writer still enforces mutation authority. +Runtime request count is unchanged; no Python rule deletion or performance gain +is claimed. Moving the retained Python effect-id compatibility resolver needs +separate pending-receipt/caller characterization and is deferred. + +This qualifies a bounded Monitor closeout and CLI successor-selection sequence, +not full M2/M3: lease transfer, mid-commit crashes, PostgreSQL, scheduler dispatch +and original-context App/Lark delivery are outside this test. Existing pending- +wait recovery tests separately cover original binding retention on File/SQLite. +No model, external provider or benchmark job runs. There is no frontend change. +Reverting restores the previous projection without rewriting persisted data, +but also restores the blocked-Monitor and missing-lane recovery gaps. diff --git a/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.zh-CN.md b/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.zh-CN.md new file mode 100644 index 0000000000..65b9315dff --- /dev/null +++ b/docs/architecture/rfcs/ledger/typescript-control-plane-migration-v0/2026-10-02-monitor-quiet-due-recovery.zh-CN.md @@ -0,0 +1,105 @@ +# Monitor 从静默到到期的恢复验证 + +[English mirror](2026-10-02-monitor-quiet-due-recovery.md) + +本 checkpoint 为 T2 和[恢复验证 RFC](../../composable-state-machines-recovery-verification-v0.zh-CN.md) +补充组合证据。 +基线 `e38b057b3` 已包含 [#4335](https://github.com/loopx-project/loopx/pull/4335) +的 effect identity 修复;较旧的已安装 runtime 仍可能出现该缺陷。 +剩余修复处理已绑定但尚未 poll 的 Monitor 变为 blocked:原 Turn 现在暴露有条件的 +生命周期恢复,不再进入普通执行或冲突的 replan 选择。不新增 RPC 或持久化 schema。 + +## 缺失的反例 + +静默 heartbeat 会自动提交一条未绑定 Todo 的 observation。若 Monitor 在同一 Turn +内变为到期,仍必须能显式绑定并提交实际观察。前一条 observation 不能完成后续 +Todo 的结算,也不能占用它的 effect identity。仅从普通未绑定 guard 测试新到期 +Monitor 会遗漏这个前提:静默 poll 已经提交。 + +```mermaid +flowchart LR + Q[静默 guard:未绑定 Todo] --> O[提交静默观察] + O --> D[Monitor 到期] + D --> B[原 Turn 绑定 Monitor] + B --> P[准确 poll 只提交一次] + P --> L[丢弃调用方回执] + L --> R[回读并重试原 poll] + R --> S[Turn 已结算;不扣配额] + S --> N[有变化时 successor 使用新 Turn] +``` + +## 可执行边界 + +`tests/control_plane/test_monitor_quiet_due_recovery.py` 确定性枚举十二条旅程: +legacy Markdown、canonical File、canonical SQLite,分别观察无变化与有变化。 +每组分别覆盖有、无其他 peer 的 user gate 及普通用户提醒。 +各旅程使用真实 CLI 子进程、TS effects 和隔离 provider 状态,执行两次原请求重试 +和一次结果冲突重试。第一次命令成功后主动丢弃响应,模拟调用方确认丢失, +不等于模拟提交过程中的进程崩溃。 + +独立断言要求:静默和 Todo-bound 两条 poll 身份不同;原观察不变;没有 refresh/ +spend 记录;重放和冲突不改状态;仅有变化时生成一次 material generation 与一个 +successor;已结算回读不允许原 Turn 再执行工作。Monitor 仍为 open;有变化时, +后继选择通过新的 Turn 完成。 + +另五条旅程在绑定之后、poll 之前将 Monitor 标为 blocked:legacy、File、SQLite +的 soft claim,以及 File、SQLite 的 hard lease。未修改基线在这个 +fixture 中错误返回 `normal_run`;其他 frontier 状态还可能尝试冲突的 replan +绑定。修复后使用既有 `unsettled_host_turn_recovery`,保留原身份且不给交付权限。 +两次回读都保留 blocked 状态;仅在 fixture 的 blocker 已解除后,测试才执行 +投影的恢复命令、重入原 guard、poll 并无配额扣减地结束 Turn。恢复命令向现有 +生命周期 owner 提交独立核实的原因与 `--clear-resume-when`。hard lease 下恢复 +不授予执行权限:无 lease 的 poll 仍拒绝,取得新 lease 后才执行精确 poll。 +另外两条 File/SQLite 用例将相同投影命令用于存在活跃 holder 的状态,验证拒绝 +且 Todo、lease 与 Goal 状态均不变。 + +敏感性验证使用同一基线的临时 checkout:恢复仅按 Turn 分配 effect 的历史规则后, +无变化/legacy 旅程在首次绑定后的 poll 以 `heartbeat_receipt_identity_conflict` +失败。未修改基线的六条旅程全部通过。这是刻意注入历史规则的 mutation, +不声称当前基线仍有该缺陷。临时 mutant 不作为发布 fixture 保留。 + +另一个反例给出不完整的 Todo frontier:聚合 work lane 缺失,且存在其他 peer 的 +user gate。精确的已结算 Monitor 仍必须返回 `heartbeat_settled_skip`。基线 adapter +在聚合 lane 缺失时丢弃已验证 phase,错误恢复普通执行。修复后,绑定 Monitor 的 +投影不再依赖聚合覆盖完整性;poll 验证与 phase 判断仍归 TS。此反例在 adapter +修复前失败、修复后通过。 + +另一个组合场景先结算 advancement Turn,再将主 Todo 置为 `blocked`, +随后观察独立到期的 Monitor。此前命令已经投影,但准入只接受 open/done 主 Todo, +导致执行失败。现在由现有 TS 事务 owner 检查精确结算回读:仅 `settled` 才允许 +这些后续 hold。尚未结算的 hold 仍拒绝,Monitor 仍须独立通过到期、owner、gate 和 +lease 检查。不重开被 hold 的 Todo,不授予新交付,不增加扣费。合成 TS 回归覆盖 +该 hold 及身份/gate 拒绝(含 deferred 主 Todo);一条 CLI 旅程覆盖 in-flight 写回、扣费、hold、poll +和精确重放。blocked 主 Todo 的 TS 正例在修改前失败。 + +## 归属与限制 + +| 边界 | 保留的现有 owner | +| --- | --- | +| selection 仲裁 | `work_items/action_portfolio.ts` | +| Monitor 事务与不可变重放 | `quota/monitor_poll_commit.ts` | +| canonical authority 中的观察与 successor | `coordination/todo_monitor_poll.ts` | +| Turn 结算回读 | `quota/settlement_readback.ts`、`quota/settlement_phase.ts` | +| 绑定 Monitor 的生命周期恢复 | `quota/blocked_wait.ts` | +| 已绑定 phase 到可选聚合 lane 的投影 | `work_items/work_lane.py` | +| legacy effect-id 兼容与 transport | `quota/monitor_poll.py` | + +伴随重构在 `blocked_wait.ts` 内共享 causal wait 与 blocked Monitor 的 current-Turn +恢复 envelope。Python 通过现有请求传递已验证的 Monitor phase,并渲染 typed repair, +不另建生命周期判断。新增 phase 字段为可选,旧请求保留 causal-wait 行为。 +只有 active、owner 合法、状态 blocked 且 phase 为 `poll_due` 的 Monitor 进入此 +恢复路线;缺失、重复、其他 owner、已归档或已 poll 的输入均排除,测试覆盖这些边界。 + +默认行为改变上述 blocked、尚未 poll 的重入,不完整 frontier 缺少聚合 lane +时已绑定 Monitor 的投影,以及精确结算后的 advancement 被 hold 时的独立辅助观察。 +恢复命令以 blocker 已核实解除为 +前提;投影不会自行重开 Todo,也不构成 poll/closeout receipt。原因未解除时 +继续保留 blocked。既有 Todo writer 仍执行变更权限检查。runtime 请求数不变, +不计 Python 规则删除量或性能收益。保留的 Python effect-id 兼容 resolver 迁移 +需要另行刻画 pending receipt/caller,本次延后。 + +本次仅验证有界 Monitor 结算与 CLI successor selection,不完成整个 M2/M3: +lease transfer、提交中断、PostgreSQL、scheduler dispatch、App/Lark 原上下文投递 +均不在覆盖范围。现有 pending-wait 回归另行验证 File/SQLite 保留原绑定的恢复。 +不调用模型、外部 provider 或 benchmark Job,无前端变更。 +撤销修复可恢复旧投影且无需改写持久数据,但也会恢复 blocked Monitor 与聚合 lane 缺失的恢复缺口。 diff --git a/loopx/control_plane/quota/blocked_wait.ts b/loopx/control_plane/quota/blocked_wait.ts index eec64a79f2..15a8587998 100644 --- a/loopx/control_plane/quota/blocked_wait.ts +++ b/loopx/control_plane/quota/blocked_wait.ts @@ -137,9 +137,9 @@ export function prepareBlockedWait(value: unknown): JsonObject { export const RECEIPT_BOUND_WAIT_REQUEST_SCHEMA = "loopx_quota_receipt_bound_wait_request_v0"; -/** A pending dependency changes executable work, never the Turn's binding. - * This is a read-only recovery projection. Only refresh-state can validate and - * commit the existing blocked closeout; observing a wait does not settle it. */ +/** A pending dependency or blocked Monitor changes executable work, never the Turn's binding. + * This read-only projection cannot settle a Turn. Causal waits need verified + * refresh-state closeout; an unpolled blocked Monitor needs lifecycle repair. */ export function projectReceiptBoundWait(value: unknown): JsonObject { const request = requireJsonObject(value, "receipt-bound wait request"); if (request.schema_version !== RECEIPT_BOUND_WAIT_REQUEST_SCHEMA || @@ -150,11 +150,17 @@ export function projectReceiptBoundWait(value: unknown): JsonObject { const todos = request.todos.map(row => requireJsonObject(row, "receipt-bound Todo")); const matches = todos.filter(row => row.todo_id === request.todo_id); const todo = matches[0]; - if (matches.length !== 1 || !todo || todo.role !== "agent" || !["open", "deferred"].includes(String(todo.status)) || - todo.archive_state !== "active" || todo.task_class !== "advancement_task" || - todo.resume_ready !== false || (todo.claimed_by && todo.claimed_by !== request.agent_id)) { + if (matches.length !== 1 || !todo || todo.role !== "agent" || + todo.archive_state !== "active" || (todo.claimed_by && todo.claimed_by !== request.agent_id)) { return {status: "none"}; } + if (todo.task_class === "continuous_monitor" && todo.status === "blocked" && + request.monitor_phase === "poll_due") { + const reason = "The Monitor bound to this Turn is blocked. Inspect its blocker and existing effects; only after verifying the blocker is resolved, restore the Monitor and re-enter this same Turn. A blocked status is not a poll receipt and does not authorize replan or independent work."; + return boundRecovery(request, "lifecycle", reason); + } + if (!["open", "deferred"].includes(String(todo.status)) || + todo.task_class !== "advancement_task" || todo.resume_ready !== false) return {status: "none"}; // This projection repairs registered Todo dependencies. Other wait kinds // retain their existing route: in particular, a valid long timer must not // be rejected by the separate 1–30 minute blocked-retry writeback budget. @@ -164,16 +170,27 @@ export function projectReceiptBoundWait(value: unknown): JsonObject { const wait = prepareBlockedWait({...request, schema_version: BLOCKED_WAIT_REQUEST_SCHEMA, allow_turn_settlement_retry: false}); const reason = "The Todo bound to this Turn now waits on a dependency. Record its verified blocked closeout without spending quota; select independent work on the next host Turn."; + return boundRecovery(request, "blocked_writeback", reason, wait); +} + +/** Both routes preserve the existing binding and grant only control-plane recovery. */ +function boundRecovery( + request: JsonObject, repair: "lifecycle" | "blocked_writeback", reason: string, wait?: JsonObject, +): JsonObject { + const obligation = repair === "blocked_writeback" + ? "close_bound_wait_without_spend" : "repair_bound_monitor_lifecycle"; return { status: "recovery_required", recovery: {schema_version: "unsettled_host_turn_recovery_v0", scope: "current_turn", turn_instance_id: request.turn_instance_id, binding_kind: "todo", - binding_id: request.todo_id, repair: "blocked_writeback", wait}, - obligation: {lane: "control_plane_recovery", next_lane: "advancement_task", - obligation: "close_bound_wait_without_spend", contract: "repair_bound_turn_closeout", - contract_obligation: "close_bound_wait_without_spend", must_attempt_work: true, + binding_id: request.todo_id, repair, ...(wait ? {wait} : {})}, + obligation: {lane: "control_plane_recovery", + next_lane: repair === "blocked_writeback" ? "advancement_task" : "continuous_monitor", + obligation, contract: "repair_bound_turn_closeout", + contract_obligation: obligation, must_attempt_work: true, delivery_allowed: false, notify: "DONT_NOTIFY", spend_policy: "no spend for blocked closeout", - reason_code: "receipt_bound_pending_wait", reason, recommendation_reason: reason, + reason_code: repair === "blocked_writeback" ? "receipt_bound_pending_wait" : "unsettled_host_turn", + reason, recommendation_reason: reason, unsettled_reason: reason, recommended_action: reason}, }; } diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index 7d3d38b921..802bebc7fb 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -686,6 +686,7 @@ def build_live_quota_should_run_decision( goal_id=goal_id, agent_id=agent_id, todo_id=receipt_bound_todo_id, turn_instance_id=turn_instance_id, available_capabilities=available_capabilities, scheduler_execution_context=resolved_context, + monitor_phase=(receipt_bound_monitor_phase.value if receipt_bound_monitor_phase else None), ) if hook_dispatch["failures"]: payload["capability_hook_dispatch"] = { diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index ef6ad2f536..6ba82cbfff 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -576,7 +576,7 @@ async function auxiliaryMonitorAllowed( if (!request.runtime_root || !request.turn_instance_id || !decision.agent_id || observation.actor_agent_id !== decision.agent_id || todo === null || todo.todo_id !== settlementTodo || todo.task_class !== "advancement_task" || - !["open", "done"].includes(String(todo.status)) || + !["open", "done", "blocked"].includes(String(todo.status)) || (todo.claimed_by != null && todo.claimed_by !== decision.agent_id) || (todo.excluded_agents != null && (!Array.isArray(todo.excluded_agents) || todo.excluded_agents.some(value => typeof value !== "string") || @@ -589,6 +589,11 @@ async function auxiliaryMonitorAllowed( if (settlement?.found !== true || jsonObject(jsonObject(settlement.identity)?.result)?.failure !== null || progress?.schema_version !== "quota_settlement_progress_v0" || progress.state === "identity_required") throw conflict(); + // A later lifecycle hold cannot invalidate an already settled Turn's + // independent observation. An unsettled hold still requires recovery; + // neither case grants execution of the held advancement Todo. + if (todo.status === "blocked" && + progress.state !== "settled") throw conflict(); } const monitor = decision.registry_due_monitor; // Retain ordinary quota/due-work admission and capability/gate projections; diff --git a/loopx/control_plane/quota/unsettled_host_turn.py b/loopx/control_plane/quota/unsettled_host_turn.py index a86d060189..7d0b5cd94a 100644 --- a/loopx/control_plane/quota/unsettled_host_turn.py +++ b/loopx/control_plane/quota/unsettled_host_turn.py @@ -266,6 +266,7 @@ def apply_receipt_bound_wait_recovery( goal_id: str, agent_id: str, todo_id: str, turn_instance_id: str, available_capabilities: list[str] | None, scheduler_execution_context: Mapping[str, Any] | SchedulerExecutionContextResolution | None, + monitor_phase: str | None = None, ) -> bool: """Read full provider facts only when a replay lost its executable binding.""" from ...todos import list_goal_todos @@ -283,6 +284,7 @@ def apply_receipt_bound_wait_recovery( "schema_version": "loopx_quota_receipt_bound_wait_request_v0", "todos": items, "todo_id": todo_id, "agent_id": agent_id, "turn_instance_id": turn_instance_id, "observed_at": now_utc_iso(), + "monitor_phase": monitor_phase, }) if verdict.get("status") == "none": return False diff --git a/loopx/control_plane/work_items/unsettled_host_turn_contract.py b/loopx/control_plane/work_items/unsettled_host_turn_contract.py index 5367f5ac28..2d7825198d 100644 --- a/loopx/control_plane/work_items/unsettled_host_turn_contract.py +++ b/loopx/control_plane/work_items/unsettled_host_turn_contract.py @@ -26,6 +26,17 @@ def recovery_cli_actions( ) # The repair lane is a typed fact from the recovery transaction; this # renderer only turns it into operator commands. + if recovery.get("repair") == "lifecycle" and recovery.get("scope") == "current_turn": + bound_turn = shlex.quote(str(recovery["turn_instance_id"])) + return [ + (f"{command_prefix} todo list --goal-id {goal_id}{lifecycle_actor_args}" + f" --todo-id {shlex.quote(prior_todo_id)}"), + "Inspect the blocker and existing effects. Only if the blocker is verified resolved, replace with that verified reason and execute the restore below. Otherwise retain the blocked Todo and report the unresolved cause; do not reset receipts or choose another Turn to bypass it.", + (f"{command_prefix} todo update --goal-id {goal_id}{lifecycle_actor_args}" + f" --todo-id {shlex.quote(prior_todo_id)} --status open" + " --reason '' --clear-resume-when"), + f"{typed_quota_guard} --turn-instance-id {bound_turn}", + ] if recovery.get("repair") == "blocked_writeback": bound_turn = shlex.quote(str(recovery["turn_instance_id"])) return [ diff --git a/loopx/control_plane/work_items/work_lane.py b/loopx/control_plane/work_items/work_lane.py index b1647de469..ce063a1225 100644 --- a/loopx/control_plane/work_items/work_lane.py +++ b/loopx/control_plane/work_items/work_lane.py @@ -161,14 +161,17 @@ def preserve_heartbeat_receipt_bound_work_lane( ) -> dict[str, Any] | None: """Keep a committed same-turn Todo binding ahead of a newly due monitor.""" - if not isinstance(contract, dict): - return contract if not isinstance(selected_todo, dict): return contract todo_id = normalize_todo_id(selected_todo.get("todo_id")) if not todo_id or selected_todo.get("selection_binding") != "heartbeat_receipt": return contract if selected_todo.get("task_class") == TODO_TASK_CLASS_MONITOR: + # An incomplete aggregate can omit its work lane. The exact bound + # Monitor still carries the TS-verified phase; preserve its projection + # without deriving execution authority from aggregate completeness. + if not isinstance(contract, dict): + contract = {} raw_monitor_phase = selected_todo.get("receipt_bound_monitor_phase") try: if not isinstance(raw_monitor_phase, str): diff --git a/skills/loopx-self-repair/references/targeted-diagnostics.md b/skills/loopx-self-repair/references/targeted-diagnostics.md index bf973a617c..18d92ced60 100644 --- a/skills/loopx-self-repair/references/targeted-diagnostics.md +++ b/skills/loopx-self-repair/references/targeted-diagnostics.md @@ -7,6 +7,7 @@ does not authorize borrowing another session's identity or changing providers. | Missing fact | First useful surface | | --- | --- | | Why the current operation was rejected | Its existing JSON error, typed recovery and operation receipt; look up the error before opening other projections | +| Whether a receipt defect still needs a source fix | Resolve the failing executable and source revision, then run the same synthetic sequence on the intended current base. An installed release can predate an already merged repair; source tests do not prove installed behavior. | | Whether an ambiguous Todo create committed | `loopx --format json todo receipt --goal-id --operation-id `; inspect the receipt before any retry with the same identity | | Current lease ownership/version | `loopx --format json task-lease inspect --goal-id --todo-id ` | | One Todo's current state | `loopx --format json todo list --goal-id --todo-id `; use `--agent-id` when needed by the existing lane | @@ -20,6 +21,13 @@ options for all calls. `quota should-run` may establish host-Turn state; it is not a harmless latency probe. Never remove a Turn id, capability declaration, lease proof, revision or receipt check to make a command faster. +For an already repaired defect, keep the original receipt and binding intact. +Qualify the installed runtime after an authorized upgrade and add missing +composition coverage instead of duplicating the fix or marking the Todo blocked +merely to escape a receipt conflict. Upgrading and resetting task state are +different actions; neither a successful installation nor a fresh Turn proves +recovery of the original operation. + ## Read one response more than once, not one command When JSON may exceed the tool output budget, capture it once in an ignored diff --git a/tests/control_plane/test_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index 6940220b70..10ca8fdf4a 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -19,6 +19,69 @@ ) +def test_settled_advancement_hold_preserves_independent_monitor( + tmp_path: Path, +) -> None: + """A post-settlement hold does not strand a due independent observation.""" + status = "blocked" + project, runtime, registry = _write_fixture(tmp_path) + turn = "turn-settled-advancement-hold" + binding = ("--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--turn-instance-id", turn) + capabilities = ("--available-capability", "network", + "--available-capability", "external_evidence_poll") + guard = ("quota", "should-run", "--codex-app", *binding, + *capabilities, "--scan-path", str(project)) + rc, first = _run_cli(registry, runtime, *guard) + assert rc == 0, first + _append_newly_due_monitor(project, watch_only=True) + rc, admitted = _run_cli(registry, runtime, *guard) + assert rc == 0, admitted + poll_command = admitted["interaction_contract"]["cli_channel"][ + "auxiliary_monitor_poll"]["command"] + rc, refresh = _run_cli( + registry, runtime, "refresh-state", *binding, "--todo-id", TODO_ID, + "--classification", "fixture_inflight_validated", + "--delivery-batch-scale", "single_surface", + "--delivery-outcome", "outcome_progress", + "--delivery-boundary", "in_flight_continuation", + "--delivery-workspace-path", str(project), + "--no-global-sync", "--suppress-external-sinks", + ) + assert rc == 0, refresh + rc, spent = _run_cli(registry, runtime, *_projected_cli_args( + refresh["settlement_owed"]["command"], turn_instance_id=turn)) + assert rc == 0, spent + rc, held = _run_cli( + registry, runtime, "todo", "update", "--goal-id", GOAL_ID, + "--todo-id", TODO_ID, "--agent-id", AGENT_ID, "--status", status, + "--reason", "Await independent evidence after validated progress", + ) + assert rc == 0, held + rc, successor = _run_cli( + registry, runtime, "todo", "add", "--goal-id", GOAL_ID, + "--role", "agent", "--task-class", "advancement_task", + "--claimed-by", AGENT_ID, "--text", "Validate an independent artifact", + ) + assert rc == 0, successor + poll_args = tuple("observed-after-hold" if token == "${LOOPX_MONITOR_RESULT_HASH:?}" + else token for token in _projected_cli_args( + poll_command, turn_instance_id=turn)) + ("--scan-path", str(project)) + for expected_replay in (False, True): + rc, poll = _run_cli(registry, runtime, *poll_args) + assert rc == 0, (poll.get("error_code"), poll.get("reason"), poll.get("error")) + assert poll["replayed"] is expected_replay + assert poll["settlement_todo_id"] == TODO_ID + assert poll["turn_continuation"]["current_turn_settled"] is True + assert poll["turn_continuation"]["same_turn_independent_settlement_allowed"] is False + assert poll["settlement_resume"]["next_step"] is None + rc, listed = _run_cli(registry, runtime, "todo", "list", "--goal-id", GOAL_ID) + assert rc == 0, listed + assert next(t for t in listed["todos"] if t["todo_id"] == TODO_ID)["status"] == status + assert _classification_count(runtime, "quota_monitor_poll") == 1 + assert _spend_run_count(runtime) == 1 + + def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( tmp_path: Path, ) -> None: diff --git a/tests/control_plane/test_monitor_quiet_due_recovery.py b/tests/control_plane/test_monitor_quiet_due_recovery.py new file mode 100644 index 0000000000..47bfbcc192 --- /dev/null +++ b/tests/control_plane/test_monitor_quiet_due_recovery.py @@ -0,0 +1,403 @@ +"""Compose quiet observation, late binding and monitor closeout through the CLI. + +Independent oracle: the pre-due observation cannot settle a later Todo. Its +exact poll commits once, closes only that Turn, and never debits quota. Losing +the poll response permits readback/replay, not another observation or successor. +""" + +from __future__ import annotations + +import json +import shlex + +import pytest + +from canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, +) +from tests.control_plane.test_monitor_followthrough_contract import ( + AGENT_ID, + GOAL_ID, + _add_monitor, + _write_fixture, +) +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) +from loopx.control_plane.testing.canary_harness import run_json_cli, run_json_cli_result +from loopx.todos import add_goal_todo, list_goal_todos + + +@pytest.mark.parametrize("provider", ["legacy", "file", "sqlite"]) +@pytest.mark.parametrize("material", [False, True], ids=["unchanged", "material"]) +@pytest.mark.parametrize("peer_gate", [False, True], ids=["no-gate", "peer-gate"]) +def test_quiet_turn_binds_due_monitor_and_recovers_lost_response( + tmp_path, + monkeypatch, + provider, + material, + peer_gate, +): + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, state = _write_fixture(tmp_path) + add_goal_todo( + registry_path=registry, + goal_id=GOAL_ID, + role="agent", + text="Independent work awaiting a dependency", + task_class="advancement_task", + status="blocked", + claimed_by=AGENT_ID, + agent_id=AGENT_ID, + ) + monitor = _add_monitor( + registry, text="Observe a public target", target_key="public-target" + ) + if provider != "legacy": + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, + todos=list_goal_todos(registry_path=registry, goal_id=GOAL_ID)["todos"], + handoff_mode="soft_claim", + ) + initialize_canonical_authority( + runtime, GOAL_ID, projection, state_path=state, provider=provider + ) + + def call(*args): + return run_json_cli(*args, registry_path=registry, runtime_root=runtime) + + def polls(): + rows = [ + json.loads(line) + for line in (runtime / "goals" / GOAL_ID / "runs/index.jsonl") + .read_text() + .splitlines() + ] + assert not any( + row["classification"] in {"quota_slot_spent", "state_refreshed"} + for row in rows + ) + return [row for row in rows if row["classification"] == "quota_monitor_poll"] + + def todos(): + return call("todo", "list", "--goal-id", GOAL_ID)["todos"] + + scope = ( + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--runtime-profile", + "codex_app_ssh_goal", + ) + first = call("quota", "should-run", *scope, "--begin-turn") + assert first["effective_action"] == "monitor_quiet_skip", first + assert "settlement_identity" not in first["heartbeat_receipt"] + turn = first["heartbeat_receipt"]["turn_instance_id"] + guard = ("quota", "should-run", *scope, "--turn-instance-id", turn) + (quiet_poll,) = polls() + assert quiet_poll.get("todo_id") is None + assert call(*guard)["should_run"] is False + assert polls() == [quiet_poll] + + # A supported metadata update makes the synthetic monitor due; no clock + # sleep or rewrite of provider authority is needed after initialization. + call( + "todo", + "update", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--todo-id", + monitor["todo_id"], + "--next-due-at", + "2000-01-01T00:00:00Z", + "--watch-only", + ) + selected = call(*guard, "--todo-id", monitor["todo_id"]) + assert selected["heartbeat_receipt"]["status"] == "upgraded" + assert ( + selected["heartbeat_receipt"]["settlement_identity"]["todo_id"] + == monitor["todo_id"] + ) + selected = call(*guard) + assert ( + selected["agent_lane_next_action"]["receipt_bound_monitor_phase"] == "poll_due" + ) + assert selected["execution_obligation"]["must_attempt_work"] is True + assert polls() == [quiet_poll] + + if peer_gate: + call( + "todo", "add", "--goal-id", GOAL_ID, "--role", "user", + "--task-class", "user_gate", "--blocks-agent", "codex-main-control", + "--text", "Choose the peer's destination", + ) + call( + "todo", "add", "--goal-id", GOAL_ID, "--role", "user", + "--task-class", "user_action", "--bound-agent", AGENT_ID, + "--text", "Read the optional guide", + ) + + poll = ( + "quota", + "monitor-poll", + *scope, + "--turn-instance-id", + turn, + "--todo-id", + monitor["todo_id"], + "--target-key", + "public-target", + "--result-hash", + "observed-head", + "--execute", + *( + ( + "--material-change", + "--next-agent-todo", + "Validate the observed change", + "--next-action-kind", + "validate", + ) + if material + else () + ), + ) + # Deliberately discard the first process's response. Recovery uses only + # durable readback and an exact retry, as after a lost caller acknowledgement. + call(*poll) + committed = polls() + assert len(committed) == 2 + assert committed[0] == quiet_poll + assert committed[1]["todo_id"] == monitor["todo_id"] + assert ( + committed[1]["quota_monitor_poll_commit"]["effect_id"] + != quiet_poll["quota_monitor_poll_commit"]["effect_id"] + ) + committed_todos, committed_state = todos(), state.read_bytes() + + for _ in range(2): + readback = call(*guard) + assert readback["selected_todo"]["todo_id"] == monitor["todo_id"] + assert ( + readback["agent_lane_next_action"]["receipt_bound_monitor_phase"] + == "settled" + ) + assert readback["should_run"] is False + assert readback["execution_obligation"]["must_attempt_work"] is False + assert ( + readback["interaction_contract"]["agent_channel"]["delivery_allowed"] + is False + ) + replay = call(*poll) + assert replay["replayed"] is True and replay["appended"] is False + assert replay["turn_continuation"]["current_turn_settled"] is True + assert polls() == committed + assert todos() == committed_todos + assert state.read_bytes() == committed_state + + # A new result is a new intent, not an idempotent retry of this Turn. + changed_poll = tuple( + "different-head" if arg == "observed-head" else arg for arg in poll + ) + rc, conflict = run_json_cli_result( + *changed_poll, registry_path=registry, runtime_root=runtime + ) + assert ( + rc != 0 and conflict["error_code"] == "heartbeat_receipt_identity_conflict" + ), conflict + assert polls() == committed + assert todos() == committed_todos + assert state.read_bytes() == committed_state + current_monitor = next( + todo for todo in committed_todos if todo["todo_id"] == monitor["todo_id"] + ) + assert current_monitor["status"] == "open" # Turn closeout is not Todo completion. + assert int(current_monitor.get("material_change_generation") or 0) == int(material) + successors = [ + todo for todo in committed_todos if todo.get("action_kind") == "validate" + ] + assert len(successors) == int(material) + if material: + next_turn = call( + "quota", + "should-run", + *scope, + "--turn-instance-id", + "independent-successor-turn", + ) + assert next_turn["selected_todo"]["todo_id"] == successors[0]["todo_id"] + assert next_turn["execution_obligation"]["must_attempt_work"] is True + + +@pytest.mark.parametrize(("provider", "handoff_mode"), [ + ("legacy", "soft_claim"), ("file", "soft_claim"), ("sqlite", "soft_claim"), + ("file", "hard_lease"), ("sqlite", "hard_lease"), +]) +def test_blocked_bound_monitor_requires_verified_lifecycle_repair( + tmp_path, + monkeypatch, + provider, + handoff_mode, +): + """A blocker changes execution eligibility, never the committed binding.""" + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, state = _write_fixture(tmp_path) + monitor = _add_monitor( + registry, + text="Observe a public target", + target_key="public-target", + next_due_at="2000-01-01T00:00:00Z", + ) + if provider != "legacy": + initialize_canonical_authority( + runtime, + GOAL_ID, + build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, + todos=list_goal_todos(registry_path=registry, goal_id=GOAL_ID)["todos"], + handoff_mode=handoff_mode, + ), + state_path=state, + provider=provider, + ) + + def call(*args): + return run_json_cli(*args, registry_path=registry, runtime_root=runtime) + + scope = ( + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--runtime-profile", + "codex_app_ssh_goal", + "--turn-instance-id", + "blocked-monitor-turn", + ) + guard = ("quota", "should-run", *scope) + initial = call(*guard) + binding = initial["heartbeat_receipt"]["settlement_identity"] + assert binding["todo_id"] == monitor["todo_id"] + call( + "todo", + "update", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--todo-id", + monitor["todo_id"], + "--status", + "blocked", "--reason", "Temporary dependency unavailable", "--clear-resume-when", + ) + before = state.read_bytes() + for _ in range(2): + recovery = call(*guard) + assert recovery["effective_action"] == "unsettled_host_turn_recovery", recovery + assert recovery["heartbeat_receipt"]["settlement_identity"] == binding + assert recovery["heartbeat_receipt"]["status"] == "replayed" + assert recovery["unsettled_host_turn_recovery"]["repair"] == "lifecycle" + assert recovery["unsettled_host_turn_recovery"]["scope"] == "current_turn" + assert ( + recovery["interaction_contract"]["agent_channel"]["delivery_allowed"] + is False + ) + assert "replan_action_packet" not in recovery + assert state.read_bytes() == before + + # The fixture's temporary blocker is now independently resolved. Execute + # the offered conditional metadata repair, then re-enter the original guard. + actions = recovery["interaction_contract"]["cli_channel"]["next_cli_actions"] + restore = next(action for action in actions if " todo update " in action) + assert "--status open" in restore + command = [ + "Temporary dependency verified available" if token == "" else token + for token in shlex.split(restore) + ] + result = run_json_cli(*command[1:], registry_path=registry, runtime_root=runtime) + assert result["ok"] is True + resumed = call(*guard) + assert resumed["heartbeat_receipt"]["settlement_identity"] == binding + assert ( + resumed["agent_lane_next_action"]["receipt_bound_monitor_phase"] == "poll_due" + ) + lease_args = () + if handoff_mode == "hard_lease": + before_poll = state.read_bytes() + rc, denied = run_json_cli_result( + "quota", "monitor-poll", *scope, "--todo-id", monitor["todo_id"], + "--target-key", "public-target", "--result-hash", "verified-head", "--execute", + registry_path=registry, runtime_root=runtime, + ) + assert rc != 0 and denied["error_code"] == "monitor_poll_rejected", denied + assert state.read_bytes() == before_poll + call("task-lease", "acquire", "--goal-id", GOAL_ID, "--todo-id", monitor["todo_id"], + "--owner", AGENT_ID, "--idempotency-key", "restored-monitor", "--ttl-seconds", "900") + lease_args = ("--use-current-task-lease",) + poll = call( + "quota", + "monitor-poll", + *scope, + "--todo-id", + monitor["todo_id"], + "--target-key", + "public-target", + "--result-hash", + "verified-head", + "--execute", + *lease_args, + ) + assert poll["turn_continuation"]["current_turn_settled"] is True + assert call(*guard)["should_run"] is False + rows = [ + json.loads(line) + for line in (runtime / "goals" / GOAL_ID / "runs/index.jsonl") + .read_text() + .splitlines() + ] + assert sum(row["classification"] == "quota_monitor_poll" for row in rows) == 1 + assert not any( + row["classification"] in {"quota_slot_spent", "state_refreshed"} for row in rows + ) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_projected_restore_rejects_an_active_execution_lease(tmp_path, monkeypatch, provider): + """Lifecycle intent cannot bypass a live holder, even on a blocked record.""" + from loopx.control_plane.work_items.unsettled_host_turn_contract import recovery_cli_actions + from loopx.control_plane.coordination.local_authority import read_canonical_todos_if_promoted + + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, state = _write_fixture(tmp_path) + monitor = _add_monitor(registry, text="Observe a public target", target_key="public-target") + run_json_cli("todo", "update", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--todo-id", monitor["todo_id"], "--status", "blocked", registry_path=registry, + runtime_root=runtime) + lease = {"schema_version": "task_lease_v0", "goal_id": GOAL_ID, + "todo_id": monitor["todo_id"], "owner": AGENT_ID, "idempotency_key": "active-monitor", + "status": "active", "expires_at": "2099-01-01T00:00:00Z", "version": 1, + "lease_epoch": 1, "write_scopes": [], "acquire_ttl_seconds": 900} + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, todos=list_goal_todos(registry_path=registry, goal_id=GOAL_ID)["todos"], + handoff_mode="hard_lease", leases=[lease], + ) + initialize_canonical_authority(runtime, GOAL_ID, projection, state_path=state, provider=provider) + actions = recovery_cli_actions( + {"unsettled_host_turn_recovery": {"repair": "lifecycle", "scope": "current_turn", + "binding_id": monitor["todo_id"], "turn_instance_id": "held-monitor-turn"}}, + command_prefix="loopx", goal_id=GOAL_ID, lifecycle_actor_args=f" --agent-id {AGENT_ID}", + typed_quota_guard="loopx quota should-run", turn_instance_id="held-monitor-turn", + ) + restore = next(action for action in actions if " todo update " in action) + command = ["Dependency verified available" if token == "" else token + for token in shlex.split(restore)] + before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) + state_before = state.read_bytes() + rc, denied = run_json_cli_result(*command[1:], registry_path=registry, runtime_root=runtime) + assert rc != 0 and denied["error_code"] == "blocked_lifecycle_active_lease", denied + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) == before + assert state.read_bytes() == state_before diff --git a/tests/control_plane/test_settled_monitor_user_gate.py b/tests/control_plane/test_settled_monitor_user_gate.py index 0c50bcf626..708a579340 100644 --- a/tests/control_plane/test_settled_monitor_user_gate.py +++ b/tests/control_plane/test_settled_monitor_user_gate.py @@ -12,6 +12,36 @@ from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection from loopx.control_plane.testing.canary_harness import run_json_cli from loopx.todos import add_goal_todo, list_goal_todos +from loopx.control_plane.effect_program import ReceiptBoundMonitorPhase +from loopx.control_plane.quota.should_run import build_quota_should_run +from loopx.control_plane.testing.quota_fixtures import ( + quota_status_payload, quota_todo_item, quota_todo_summary, +) + + +def test_incomplete_frontier_cannot_erase_settled_monitor_receipt(): + monitor = quota_todo_item(todo_id="todo_monitor", title="Observe a public target", + task_class="continuous_monitor", claimed_by="agent-a", + next_due_at="2099-01-01T00:00:00Z", watch_only=True) + summary = quota_todo_summary([monitor], claim_scope_agent_id="agent-a") + # An incomplete legacy source cannot establish a monitor-only schedule. + # This changes aggregate coverage, not the exact persisted poll receipt. + summary["work_counts"]["complete"] = False + status = quota_status_payload(goal_id="compact-monitor", status="active", + quota_state="operator_gate", recommended_action="Wait for an owner decision", + agent_todos=summary, user_todo_items=[quota_todo_item( + todo_id="todo_peer_gate", title="Choose peer destination", role="user", + task_class="user_gate", blocks_agent="agent-b")], + coordination={"agent_model": "peer_v1", "registered_agents": ["agent-a", "agent-b"]}) + replay = build_quota_should_run(status, goal_id="compact-monitor", agent_id="agent-a", + receipt_bound_todo_id="todo_monitor", + receipt_bound_monitor_phase=ReceiptBoundMonitorPhase.SETTLED) + assert replay["selected_todo"]["todo_id"] == "todo_monitor" + assert replay["should_run"] is False + assert replay["effective_action"] == "heartbeat_settled_skip" + assert replay["safe_bypass_allowed"] is False + assert replay["execution_obligation"]["must_attempt_work"] is False + assert replay["interaction_contract"]["agent_channel"]["delivery_allowed"] is False @pytest.mark.parametrize("provider", ["legacy", "file", "sqlite"]) diff --git a/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts index 07b039d9b6..2fec58f10a 100644 --- a/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts +++ b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts @@ -72,8 +72,10 @@ test("fresh auxiliary admission requires exact advancement identity and ordinary ["other Turn", p => ({ ...p, turn_instance_id: "other-turn" }), "heartbeat_receipt_identity_conflict"], ["other binding", p => ({ ...p, observation: { ...p.observation as JsonObject, settlement_todo_id: "todo_other" } }), "heartbeat_receipt_identity_conflict"], ["missing lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, auxiliary_settlement_todo: null } }), "heartbeat_receipt_identity_conflict"], - ["blocked lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, - auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "blocked" } } }), "heartbeat_receipt_identity_conflict"], + ["invalid lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "unknown" } } }), "heartbeat_receipt_identity_conflict"], + ["deferred lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "deferred" } } }), "heartbeat_receipt_identity_conflict"], ["foreign primary", p => ({ ...p, decision: { ...p.decision as JsonObject, auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "done", claimed_by: "peer" } } }), "heartbeat_receipt_identity_conflict"], ["excluded actor", p => ({ ...p, decision: { ...p.decision as JsonObject, @@ -86,9 +88,11 @@ test("fresh auxiliary admission requires exact advancement identity and ordinary work_lane_contract: { must_attempt_work: false }, should_run: false } }), "monitor_poll_admission_rejected"], ["user action", p => ({ ...p, decision: { ...p.decision as JsonObject, requires_user_action: true } }), "monitor_poll_admission_rejected"], ]; - for (const settled of [false, true]) for (const [name, change, code] of cases) { - await t.test(`${settled ? "settled" : "pending"}: ${name}`, async st => { + for (const settled of [false, true]) for (const status of settled ? ["done", "blocked"] : ["done"]) for (const [name, change, code] of cases) { + await t.test(`${settled ? "settled" : "pending"} ${status}: ${name}`, async st => { const { runtime, params } = await fixture(st); + const decision = params.decision as JsonObject; + decision.auxiliary_settlement_todo = { ...decision.auxiliary_settlement_todo as JsonObject, status }; if (settled) await settlePrimary(runtime); const initialIndex = await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"); const changed = change(params); @@ -117,8 +121,11 @@ function providerReceipt(params: JsonObject): JsonObject { todo_update: { ok: true }, next_todos: [], successor_receipts: [] }; } -test("settled primary admits the first due auxiliary observation without another debit or delivery", async t => { +for (const status of ["done", "blocked"]) { +test(`settled ${status} primary admits the first due auxiliary observation without another debit or delivery`, async t => { const { runtime, params } = await fixture(t); + const decision = params.decision as JsonObject; + decision.auxiliary_settlement_todo = { ...decision.auxiliary_settlement_todo as JsonObject, status }; await settlePrimary(runtime); const index = await readFile(join(runtime, "goals", goal, "runs", "index.jsonl")); const request = { ...params, expected_index_digest: `sha256:${createHash("sha256").update(index).digest("hex")}` }; @@ -138,6 +145,21 @@ test("settled primary admits the first due auxiliary observation without another assert.equal(rows.filter(row => row.classification === "quota_slot_spent").length, 1); assert.equal(rows.length, 3); }); +} + +for (const status of ["blocked", "deferred"]) { +test(`unsettled ${status} primary does not admit a new auxiliary effect`, async t => { + const { runtime, params } = await fixture(t); + const decision = params.decision as JsonObject; + decision.auxiliary_settlement_todo = { ...decision.auxiliary_settlement_todo as JsonObject, status }; + await assert.rejects(evaluateQuotaMonitorPollCommit(params), error => { + assert.ok(error instanceof EffectRuntimeRequestError); + assert.equal(error.code, "heartbeat_receipt_identity_conflict"); + return true; + }); + assert.equal(await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"), ""); +}); +} test("completed primary preserves an admitted pending effect across settlement and replay reads current closeout", async t => { const { runtime, params } = await fixture(t); diff --git a/tests/control_plane_ts/causal_blocked_wait.test.ts b/tests/control_plane_ts/causal_blocked_wait.test.ts index ce50059f00..26cc90bdbb 100644 --- a/tests/control_plane_ts/causal_blocked_wait.test.ts +++ b/tests/control_plane_ts/causal_blocked_wait.test.ts @@ -1,6 +1,6 @@ import assert from "node:assert/strict"; import test from "node:test"; -import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait } from "../../loopx/control_plane/quota/blocked_wait.ts"; +import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait, projectReceiptBoundWait, RECEIPT_BOUND_WAIT_REQUEST_SCHEMA } from "../../loopx/control_plane/quota/blocked_wait.ts"; import { isBoundedBlockedRetry } from "../../loopx/control_plane/quota/settlement_phase.ts"; import { evaluateTodoResumeConditions, TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION } from "../../loopx/control_plane/todos/resume_condition.ts"; @@ -81,3 +81,31 @@ test("a PR wait gives a causal recovery route without accepting caller-authored request.todos[0].resume_when = "pr_merged:example/project#1"; assert.throws(() => prepareBlockedWait(request), /cannot qualify a PR merge wait.*registered monitor_changed or todo_done/); }); + +test("an unpolled blocked Monitor offers lifecycle repair without execution authority", () => { + const todo = {todo_id: "todo_monitor", role: "agent", task_class: "continuous_monitor", + status: "blocked", archive_state: "active", claimed_by: "agent-a"}; + const request = {schema_version: RECEIPT_BOUND_WAIT_REQUEST_SCHEMA, + turn_instance_id: "turn-a", agent_id: "agent-a", todo_id: todo.todo_id, + monitor_phase: "poll_due", todos: [todo]}; + const before = structuredClone(request); + const result = projectReceiptBoundWait(request); + assert.equal(result.status, "recovery_required"); + assert.deepEqual(result.recovery, {schema_version: "unsettled_host_turn_recovery_v0", + scope: "current_turn", turn_instance_id: "turn-a", binding_kind: "todo", + binding_id: "todo_monitor", repair: "lifecycle"}); + const obligation = result.obligation as Record; + assert.equal(obligation.delivery_allowed, false); + assert.equal(obligation.must_attempt_work, true); + assert.deepEqual(request, before); + for (const patch of [{claimed_by: "agent-b"}, {archive_state: "archived"}, + {role: "user"}, {status: "open"}, {status: "done"}, {task_class: "advancement_task"}]) { + assert.equal(projectReceiptBoundWait({...request, todos: [{...todo, ...patch}]}).status, "none"); + } + for (const monitor_phase of [undefined, "settled", "settlement_pending", "invalid"]) { + assert.equal(projectReceiptBoundWait({...request, monitor_phase}).status, "none"); + } + assert.equal(projectReceiptBoundWait({...request, todos: []}).status, "none"); + assert.equal(projectReceiptBoundWait({...request, todos: [todo, todo]}).status, "none"); + assert.throws(() => projectReceiptBoundWait({...request, turn_instance_id: null}), /bound Turn/); +});