From e08fc5656b80a5dde46cbde83066eb5049c1cb1b Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 18:00:26 +0800 Subject: [PATCH] fix(quota): resume interrupted advancement under its original turn Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/heartbeat-automation-prompt.md | 14 ++-- .../quota/unsettled_host_turn_recovery.ts | 39 ++++++++--- .../unsettled_host_turn_contract.py | 25 ++++++- .../test_interrupted_turn_continuation.py | 66 +++++++++++++++++++ .../test_quota_settlement_cli.py | 3 +- .../unsettled_host_turn_recovery.test.ts | 37 ++++++++++- 6 files changed, 166 insertions(+), 18 deletions(-) create mode 100644 tests/control_plane/test_interrupted_turn_continuation.py diff --git a/docs/heartbeat-automation-prompt.md b/docs/heartbeat-automation-prompt.md index b41b1bdaa6..15c4473db9 100644 --- a/docs/heartbeat-automation-prompt.md +++ b/docs/heartbeat-automation-prompt.md @@ -295,10 +295,16 @@ whose capabilities are known when the automation is installed. `closeout_required=true`. A fresh heartbeat checks the immediately preceding flagged guard against its exact writeback/spend receipts and typed Todo lifecycle. If neither is present, `unsettled_host_turn_recovery_v0` preempts - ordinary work selection. The host must repair the prior closeout, rerun the - same current Turn, and then continue an eligible successor. Recovery is - idempotent and no-spend; receipts created before this explicit flag are not - retroactively treated as unsettled; + ordinary work selection. For an open advancement Todo without a resume gate, + the typed `resume_prior_turn` route re-enters the original Turn guard: inspect + existing effects first, follow its current eligibility and settlement contract, + and account only verified work under that original identity. Missing receipts + are not evidence of an external wait; never invent `monitor_changed` or a + successor to clear recovery. After legal closeout, rerun the same current Turn + and continue eligible work. Recovery selection itself is idempotent and + no-spend; actual validated delivery retains normal accounting. Genuine external + waits and monitor observations retain their existing typed closeouts. Receipts + created before the explicit flag are not retroactively treated as unsettled; - use `user_gate` only for an exact authority boundary such as approval to merge an aggregate branch into `main`, release, launch a benchmark, or perform a protected action; diff --git a/loopx/control_plane/quota/unsettled_host_turn_recovery.ts b/loopx/control_plane/quota/unsettled_host_turn_recovery.ts index 9a968c0783..2cc4a543a6 100644 --- a/loopx/control_plane/quota/unsettled_host_turn_recovery.ts +++ b/loopx/control_plane/quota/unsettled_host_turn_recovery.ts @@ -386,16 +386,36 @@ function acceptedLifecycleCloseout(todo: CandidateTodoFacts | null): AcceptedClo return settled ? "typed_blocker_or_lifecycle_transition" : null; } +type RecoveryRepair = "monitor_poll" | "resume_prior_turn" | "lifecycle"; + +function recoveryRepair( + bindingKind: HeartbeatReceiptBinding["binding_kind"], + todo: CandidateTodoFacts | null, +): RecoveryRepair { + if (bindingKind !== "todo" || todo === null) return "lifecycle"; + if (todo.task_class === MONITOR_TASK_CLASS) return "monitor_poll"; + // Missing receipts do not prove an external dependency. Re-enter + // the original guard to recheck eligibility and retain its settlement id. + // This is a recovery route, never an accepted closeout or delivery grant. + return todo.task_class === "advancement_task" && todo.status === "open" && + !todo.has_resume_when + ? "resume_prior_turn" + : "lifecycle"; +} + function recoveryObligation( - repair: "monitor_poll" | "lifecycle", + repair: RecoveryRepair, bindingId: string, ): JsonObject { + const obligation = repair === "resume_prior_turn" + ? "reenter_original_guard_then_settle" + : "author_typed_closeout_then_continue_successor"; return { lane: "control_plane_recovery", next_lane: "advancement_task", - obligation: "author_typed_closeout_then_continue_successor", + obligation, contract: "repair_prior_turn_closeout", - contract_obligation: "author_typed_closeout_then_continue_successor", + contract_obligation: obligation, must_attempt_work: true, delivery_allowed: false, notify: "DONT_NOTIFY", @@ -404,9 +424,12 @@ function recoveryObligation( reason: "a prior must-attempt heartbeat has no legal closeout receipt", recommendation_reason: "prior must-attempt host Turn is missing a legal closeout", unsettled_reason: "prior must-attempt host Turn remains unsettled", - recommended_action: - `Recover prior unsettled host Turn for ${bindingId}; use a typed ` + - "lifecycle observation, then rerun quota and continue independent work", + recommended_action: repair === "resume_prior_turn" + ? `Recover prior unsettled host Turn for ${bindingId}; re-enter its original ` + + "guard, inspect existing effects, then resume and settle verified work under " + + "that identity before rerunning the current Turn; do not invent an external wait" + : `Recover prior unsettled host Turn for ${bindingId}; use a typed ` + + "lifecycle observation, then rerun quota and continue eligible work", repair, }; } @@ -490,9 +513,7 @@ export function reduceUnsettledHostTurnRecovery(value: unknown): JsonObject { recovery.binding_target_key = todo.target_key; recovery.binding_cadence = todo.cadence; } - const repair = todo?.task_class === MONITOR_TASK_CLASS - ? "monitor_poll" as const - : "lifecycle" as const; + const repair = recoveryRepair(selected.binding_kind, todo); // The typed repair lane lets the CLI renderer project its command list // without re-deriving which closeout family the Turn belongs to. recovery.repair = repair; 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 4c2cd3e530..60d9a5dd26 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,27 @@ 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") == "resume_prior_turn": + prior_turn_id = str(recovery["prior_turn_instance_id"]) + return [ + ( + "inspect existing effects and persisted outcomes before retrying; " + "missing receipts do not establish an external wait or authorize " + "duplicate execution. Re-enter the original guard and follow its " + "fresh eligibility and settlement contract; record only verified " + "outcomes, never fabricate progress to clear recovery" + ), + ( + f"{typed_quota_guard} --turn-instance-id " + f"{shlex.quote(prior_turn_id)} --todo-id " + f"{shlex.quote(prior_todo_id)}" + ), + ( + "after the original Turn is legally settled, rerun the current " + "Turn below; recovery itself does not spend quota" + ), + f"{typed_quota_guard}{current_turn_arg}", + ] if recovery.get("repair") == "monitor_poll": prior_turn_id = str( recovery.get("prior_turn_instance_id") or "" @@ -54,7 +75,9 @@ def recovery_cli_actions( return [ ( "inspect unsettled_host_turn_recovery and supply a typed host " - "observation; never infer external state from Todo prose" + "observation; never infer external state from Todo prose. Only a " + "verified external-only wait may use the conditional transition below; " + "otherwise repair the actual lifecycle or missing binding" ), ( f"{command_prefix} todo update --goal-id {goal_id} --todo-id " diff --git a/tests/control_plane/test_interrupted_turn_continuation.py b/tests/control_plane/test_interrupted_turn_continuation.py new file mode 100644 index 0000000000..05534f9497 --- /dev/null +++ b/tests/control_plane/test_interrupted_turn_continuation.py @@ -0,0 +1,66 @@ +"""Interrupted work must remain runnable rather than acquire a fictitious wait.""" +from __future__ import annotations + +from pathlib import Path + +from test_quota_settlement_cli import ( + AGENT_ID, GOAL_ID, TODO_ID, _run_cli, _run_generated_cli, + _spend_run_count, _write_fixture, +) + + +def test_original_turn_can_resume_and_settle_without_mutating_todo(tmp_path: Path) -> None: + project, runtime, registry = _write_fixture(tmp_path) + prior_turn = "interrupted-original" + recovery_turn = "interrupted-recovery" + + def guard(turn: str): + rc, result = _run_cli( + registry, runtime, "quota", "should-run", "--codex-app", + "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--turn-instance-id", turn, "--todo-id", TODO_ID, + "--scan-path", str(project), + ) + assert rc == 0, result + return result + + admitted = guard(prior_turn) + assert admitted["heartbeat_receipt"]["closeout_required"] is True + interrupted = guard(recovery_turn) + assert interrupted["effective_action"] == "unsettled_host_turn_recovery" + assert interrupted["unsettled_host_turn_recovery"]["repair"] == "resume_prior_turn" + actions = interrupted["interaction_contract"]["cli_channel"]["next_cli_actions"] + assert not any("--resume-when" in action for action in actions) + state_path = project / f".codex/goals/{GOAL_ID}/ACTIVE_GOAL_STATE.md" + state_before = state_path.read_bytes() + rc, resumed = _run_generated_cli(actions[1], registry_path=registry) + assert rc == 0, resumed + assert state_path.read_bytes() == state_before + assert resumed["effective_action"] != "unsettled_host_turn_recovery", resumed + assert resumed["selected_todo"]["todo_id"] == TODO_ID + assert _spend_run_count(runtime) == 0 + + rc, writeback = _run_cli( + registry, runtime, "refresh-state", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--todo-id", TODO_ID, + "--turn-instance-id", prior_turn, + "--classification", "validated_progress", + "--delivery-batch-scale", "implementation", + "--delivery-outcome", "outcome_progress", + "--delivery-boundary", "in_flight_continuation", + "--no-global-sync", "--suppress-external-sinks", + ) + assert rc == 0, writeback + for _ in range(2): + rc, spent = _run_cli( + registry, runtime, "quota", "spend-slot", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--todo-id", TODO_ID, + "--turn-instance-id", prior_turn, "--source", "heartbeat", + "--slots", "1", "--execute", "--scan-path", str(project), + ) + assert rc == 0, spent + assert _spend_run_count(runtime) == 1 + rc, next_guard = _run_generated_cli(actions[-1], registry_path=registry) + assert rc == 0, next_guard + assert next_guard["effective_action"] != "unsettled_host_turn_recovery", next_guard + assert next_guard["selected_todo"]["todo_id"] == TODO_ID diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index c7cb59132e..84db541c13 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -1283,9 +1283,10 @@ def test_recovery_does_not_bind_current_replan_and_reenters_same_turn( "$.unsettled_host_turn_recovery" ) assert contract["cli_channel"]["spend_after_validation"] is False - assert "monitor_changed:" in contract["cli_channel"][ + assert "--turn-instance-id turn-unsettled-prior" in contract["cli_channel"][ "next_cli_actions" ][1] + assert packet["repair"] == "resume_prior_turn" assert "settlement_identity" not in recovery["heartbeat_receipt"] assert recovery["heartbeat_receipt"]["semantic_replan_obligation_id"] == ( replan_obligation_id diff --git a/tests/control_plane_ts/unsettled_host_turn_recovery.test.ts b/tests/control_plane_ts/unsettled_host_turn_recovery.test.ts index c75644ac9f..9cc926d31d 100644 --- a/tests/control_plane_ts/unsettled_host_turn_recovery.test.ts +++ b/tests/control_plane_ts/unsettled_host_turn_recovery.test.ts @@ -413,7 +413,7 @@ test("a committed monitor-poll closeout is not accepted for a non-monitor Turn", ), ); assert.equal(verdict.status, "recovery_required"); - assert.equal((verdict.obligation as JsonObject).repair, "lifecycle"); + assert.equal((verdict.obligation as JsonObject).repair, "resume_prior_turn"); } finally { await runtime.close(); } @@ -484,14 +484,45 @@ test("an unsettled Turn keeps the exact public recovery payload", async () => { binding_task_class: "advancement_task", binding_target_key: "key-1", binding_cadence: "1h", - repair: "lifecycle", + repair: "resume_prior_turn", }); const obligation = verdict.obligation as JsonObject; assert.equal(obligation.lane, "control_plane_recovery"); assert.equal(obligation.must_attempt_work, true); assert.equal(obligation.delivery_allowed, false); assert.equal(obligation.notify, "DONT_NOTIFY"); - assert.equal(obligation.repair, "lifecycle"); + assert.equal(obligation.repair, "resume_prior_turn"); + } finally { + await runtime.close(); + } +}); + +test("continuation is a guarded route, not a closeout or an inferred wait", async () => { + const runtime = await runtimeWith([ + receipt("turn-a", closeoutRequired("turn-a", "todo_alpha")), + ]); + try { + const candidate = candidateFrom(await preflight(runtime.root)); + for (const missing of [ + ["durable_writeback_receipt", "quota_spend_receipt"], + ["quota_spend_receipt"], + ]) { + for (const [todo, repair] of [ + [ADVANCEMENT_OPEN, "resume_prior_turn"], + [{...ADVANCEMENT_OPEN, has_successor_todo_ids: true}, "resume_prior_turn"], + [{...ADVANCEMENT_OPEN, has_resume_when: true}, "lifecycle"], + [{...ADVANCEMENT_OPEN, task_class: "user_gate"}, "lifecycle"], + [{...ADVANCEMENT_OPEN, status: "open "}, "lifecycle"], + [null, "lifecycle"], + ] as const) { + const verdict = reduce(candidate, missing, READ_TODO(todo)); + assert.equal(verdict.status, "recovery_required"); + assert.equal((verdict.recovery as JsonObject).repair, repair); + assert.equal((verdict.obligation as JsonObject).delivery_allowed, false); + assert.equal(verdict.accepted_closeout, undefined); + assert.deepEqual((verdict.recovery as JsonObject).missing_receipts, missing); + } + } } finally { await runtime.close(); }