From 606f48434f2eee14816d74c7ace585a8003220d2 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 04:07:29 +0800 Subject: [PATCH 1/6] fix(quota): preserve deferred turn selection Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota-monitor-observation-receipt-v0.md | 6 + docs/reference/protocols/turn-envelope-v0.md | 23 +++- loopx/cli_commands/quota.py | 65 +++++++++- .../control_plane/effect_runtime_handlers.ts | 5 + .../control_plane/quota/heartbeat_receipt.py | 115 +++++++++++++++++- loopx/control_plane/quota/live_decision.py | 106 ++++++++++++++++ .../quota/monitor_poll_commit.ts | 11 ++ loopx/control_plane/quota/settlement_cli.py | 23 ++++ .../work_items/action_portfolio.py | 40 ++++++ .../work_items/action_portfolio.ts | 73 +++++++++++ .../test_monitor_followthrough_contract.py | 11 ++ .../test_quota_settlement_cli.py | 44 ++++--- .../test_selection_replan_reentry.py | 79 +++++++++++- .../control_plane_ts/action_portfolio.test.ts | 45 +++++++ 14 files changed, 616 insertions(+), 30 deletions(-) diff --git a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md index 96e7fa640c..9481dca9eb 100644 --- a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md +++ b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md @@ -34,6 +34,9 @@ advancement work remains active. remains selected. A material observation may create its independently routed successor through the existing monitor contract, but it still does not replace the Turn's settlement identity. +- An executed turn-scoped poll returns `turn_continuation` declaring the current + Turn settled and requiring a fresh `--turn-instance-id` before unrelated work. + A no-spend closeout does not make the settled Turn reusable. ### Acceptance @@ -184,6 +187,9 @@ using a complete read-only snapshot with disposable File/SQLite/PostgreSQL arms. 替换为本 Turn 的结算 Todo;多个辅助回执也绝不改变既有结算身份。 - 辅助观察无变化后,原 advancement Todo 继续保持选中;若观察发生重大变化, 可按既有 monitor 契约创建独立路由的 successor,但仍不替换本 Turn 的结算身份。 +- 执行成功的 turn-scoped poll 会返回 `turn_continuation`,明确当前 Turn 已结算; + 开始无关工作前必须使用新的 `--turn-instance-id`。不计费 closeout 不代表原 Turn + 可以再次绑定另一份独立工作。 ### 验收 diff --git a/docs/reference/protocols/turn-envelope-v0.md b/docs/reference/protocols/turn-envelope-v0.md index a93b487510..2e65610b14 100644 --- a/docs/reference/protocols/turn-envelope-v0.md +++ b/docs/reference/protocols/turn-envelope-v0.md @@ -76,15 +76,25 @@ settlement-identity conflict. The full quota response preserves the TypeScript `action_selection_qualification_v0` result and returns `quota_action_selection_deferred` or `quota_action_selection_rejected`, including the exact current preemption or eligibility reason. An existing identity-less -receipt is replayed without mutation; a first-call rejection reports +receipt appends an identity-less `pending_action_selection` revision when an +otherwise eligible explicit choice is deferred. That revision is not delivery +or settlement authority: it preserves the explicit choice only so a no-argument +same-Turn reentry cannot replace it with the current recommendation. An +ineligible/rejected choice still replays the receipt without mutation; a +first-call rejection reports `heartbeat_receipt.status=not_committed` and writes no receipt event. The agent receives `recovery_action=reenter_guard_without_selection` and one executable same-Turn guard in the full decision's `cli_channel.next_cli_actions`; the compact envelope preserves the recovery in its action and writeback preview. The failed selection exposes no settlement plan, spend command, or unadmitted replan action packet. Execute that guard without -a Todo/replan argument before following the resulting binding or portfolio. A receipt already bound to a different Todo or autonomous -replan obligation remains a hard `heartbeat_receipt_identity_conflict`. +a Todo/replan argument before following the resulting binding or portfolio. A +receipt already bound to a different Todo or autonomous replan obligation +remains a hard `heartbeat_receipt_identity_conflict`. On reentry, an identical +projected Todo may bind normally. If a hard autonomous replan owns the current +lane, the replan receives the Turn's settlement identity and the retained Todo +is reported as `deferred_to_fresh_turn`; a different recommended Todo never +inherits the retained choice or its authority. When a due monitor is visible only as auxiliary context for an advancement lane, the typed reason is `auxiliary_monitor_not_selectable_in_advancement_lane`. The agent selects a @@ -98,6 +108,13 @@ worktree recovery instruction. Moving to an independent worktree and rerunning the guard with the same Turn id resumes the selected Todo; the wrapper must not rewrite this recoverable state as a settlement-identity conflict. +An executed, turn-scoped `quota monitor-poll` is a no-spend closeout, but it is +still that Turn's single settlement identity. Its response therefore includes +`turn_continuation.next_turn_required=true` and an instruction to rerun +`quota should-run` with a fresh `--turn-instance-id` before starting unrelated +work. A host must not interpret successful no-spend closeout as permission to +bind an independent Todo in the already-settled Turn. + Portfolio v2 preserves v1's selection policy, candidate ordering, and settlement rules, and adds an optional `continuation_hint` to each suggested action. The default quota producer and Turn controller now require v2. The diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index cd0b74ecbc..95c335fd72 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -29,9 +29,11 @@ HEARTBEAT_RECEIPT_SCHEMA_VERSION, fail_heartbeat_receipt, find_heartbeat_receipt, + heartbeat_receipt_pending_action_todo_id, heartbeat_receipt_settlement_replan_obligation_id, heartbeat_receipt_settlement_todo_id, heartbeat_receipt_view, + retain_pending_heartbeat_action_selection, ) from ..control_plane.quota.live_decision import build_live_quota_should_run_decision from ..control_plane.quota.monitor_poll import find_quota_monitor_poll_turn @@ -173,9 +175,9 @@ def _heartbeat_quota_action_selection_bindings( runtime_root: Path, args: argparse.Namespace, heartbeat_turn_id: str | None, -) -> tuple[dict[str, object] | None, str | None, str | None]: +) -> tuple[dict[str, object] | None, str | None, str | None, str | None]: if not heartbeat_turn_id: - return None, None, None + return None, None, None, None existing = find_heartbeat_receipt( runtime_root, goal_id=args.goal_id, @@ -183,9 +185,14 @@ def _heartbeat_quota_action_selection_bindings( turn_instance_id=heartbeat_turn_id, ) if not existing: - return None, None, None + return None, None, None, None todo_id, replan_obligation_id = _heartbeat_receipt_settlement_bindings(existing) - return existing, todo_id, replan_obligation_id + return ( + existing, + todo_id, + replan_obligation_id, + heartbeat_receipt_pending_action_todo_id(existing), + ) def _apply_requested_quota_action_selection_preflight( @@ -548,6 +555,7 @@ def handle_quota_command( heartbeat_receipt_existing, receipt_bound_todo_id, receipt_bound_replan_obligation_id, + receipt_pending_action_todo_id, ) = _heartbeat_quota_action_selection_bindings( runtime_root=runtime_root, args=args, @@ -579,6 +587,13 @@ def handle_quota_command( else None ), receipt_bound_replan_obligation_id=(receipt_bound_replan_obligation_id), + retained_action_selection_todo_id=( + receipt_pending_action_todo_id + if _requested_quota_action_todo_id(args) is None + and receipt_bound_todo_id is None + and receipt_bound_replan_obligation_id is None + else None + ), turn_instance_id=heartbeat_turn_id, interaction_projection_hooks=interaction_projection_hooks, turn_start_hook_dispatch=turn_start_hook_dispatch, @@ -589,6 +604,37 @@ def handle_quota_command( receipt_bound_todo_id=receipt_bound_todo_id, receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, ) + if ( + heartbeat_turn_id + and action_selection_preflight_failed + and heartbeat_receipt_existing is not None + and (requested_todo_id := _requested_quota_action_todo_id(args)) + and isinstance( + qualification := payload.get("action_selection_qualification"), + Mapping, + ) + and qualification.get("state") == "deferred" + ): + qualification_reason = str( + qualification.get("reason") or "current_delivery_gate" + ) + ( + heartbeat_receipt_existing, + retained_selection_appended, + ) = retain_pending_heartbeat_action_selection( + runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=heartbeat_turn_id, + todo_id=requested_todo_id, + reason=qualification_reason, + ) + heartbeat_receipt_existing_status = ( + "selection_retained" + if retained_selection_appended + else "replayed" + ) + heartbeat_receipt_existing_appended = retained_selection_appended if heartbeat_turn_id: if action_selection_preflight_failed: heartbeat_receipt_ready = True @@ -668,6 +714,13 @@ def handle_quota_command( receipt_bound_replan_obligation_id=( receipt_bound_replan_obligation_id ), + retained_action_selection_todo_id=( + receipt_pending_action_todo_id + if _requested_quota_action_todo_id(args) is None + and receipt_bound_todo_id is None + and receipt_bound_replan_obligation_id is None + else None + ), turn_instance_id=heartbeat_turn_id, interaction_projection_hooks=(interaction_projection_hooks), ) @@ -781,8 +834,8 @@ def handle_quota_command( payload, receipt=heartbeat_receipt_existing, turn_instance_id=heartbeat_turn_id, - status="replayed", - appended=False, + status=heartbeat_receipt_existing_status, + appended=heartbeat_receipt_existing_appended, ) else: _attach_uncommitted_action_selection_receipt( diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index eae441182f..95fcfe1647 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -118,6 +118,7 @@ import { import { projectQuotaActionPortfolio, qualifyActionSelection, + reconcileRetainedActionSelection, } from "./work_items/action_portfolio.ts"; import { projectQuotaPlanningHorizon } from "./work_items/planning_horizon.ts"; import { projectTaskGraphTopology } from "./work_items/task_graph.ts"; @@ -457,6 +458,10 @@ export function createEffectRuntimeHandlers( ["turn.delivery_route.evaluate", evaluateDeliveryRoute], ["work_item.action_portfolio.project", projectQuotaActionPortfolio], ["work_item.action_selection.qualify", qualifyActionSelection], + [ + "work_item.action_selection.reconcile_retained", + reconcileRetainedActionSelection, + ], ["work_item.planning_horizon.project", projectQuotaPlanningHorizon], ["work_item.task_graph.topology", projectTaskGraphTopology], ["work_item.planning_inventory.project", projectTodoPlanningInventory], diff --git a/loopx/control_plane/quota/heartbeat_receipt.py b/loopx/control_plane/quota/heartbeat_receipt.py index cbb14056e2..8fe2b37b63 100644 --- a/loopx/control_plane/quota/heartbeat_receipt.py +++ b/loopx/control_plane/quota/heartbeat_receipt.py @@ -12,7 +12,7 @@ load_rollout_events, rollout_event_log_path, ) -from ..todos.contract import normalize_todo_replan_obligation_id +from ..todos.contract import normalize_todo_id, normalize_todo_replan_obligation_id from .effect_program import SettlementBindingKind, SettlementIdentity from .error_codes import HeartbeatReceiptIdentityConflictError from .settlement_workspace_causality import ( @@ -116,6 +116,19 @@ def heartbeat_receipt_semantic_replan_obligation_id( ) +def heartbeat_receipt_pending_action_todo_id( + event: Mapping[str, object] | None, +) -> str | None: + """Return an explicit Todo choice retained before settlement binding.""" + + if not isinstance(event, Mapping): + return None + details = event.get("details") + if not isinstance(details, Mapping): + return None + return normalize_todo_id(details.get("pending_action_selection_todo_id")) + + def _effective_heartbeat_receipt( events: list[dict[str, object]], ) -> dict[str, object] | None: @@ -215,6 +228,91 @@ def ensure_turn_heartbeat_settlement_receipt( return receipt +def retain_pending_heartbeat_action_selection( + runtime_root: Path, + *, + goal_id: str, + agent_id: str, + turn_instance_id: str, + todo_id: str, + reason: str, +) -> tuple[dict[str, object], bool]: + """Append an identity-less revision retaining one explicit Todo choice. + + The choice is not settlement authority. It only fences a later no-argument + reentry from silently replacing the agent's explicit selection with a new + recommendation. A later explicit selection may replace it until the Turn + acquires its single settlement identity. + """ + + normalized_todo_id = normalize_todo_id(todo_id) + if normalized_todo_id is None: + raise HeartbeatReceiptIdentityConflictError( + "pending heartbeat action selection requires a legal Todo id" + ) + normalized_reason = str(reason or "current_delivery_gate").strip() + log_path = rollout_event_log_path(runtime_root, goal_id) + log_path.parent.mkdir(parents=True, exist_ok=True) + with exclusive_file_lock(log_path): + events = load_rollout_events(log_path) + matching = _heartbeat_receipt_events( + events, + goal_id=goal_id, + agent_id=agent_id, + turn_instance_id=turn_instance_id, + ) + effective = _effective_heartbeat_receipt(matching) + if effective is None: + raise ValueError( + "identity-less heartbeat receipt is missing; rerun the original guard" + ) + if _receipt_settlement_identity(effective) is not None: + raise HeartbeatReceiptIdentityConflictError( + "a settled heartbeat Turn cannot retain a different pending selection" + ) + details_value = effective.get("details") + details = ( + dict(details_value) if isinstance(details_value, Mapping) else {} + ) + if ( + normalize_todo_id(details.get("pending_action_selection_todo_id")) + == normalized_todo_id + and str(details.get("pending_action_selection_reason") or "").strip() + == normalized_reason + ): + return effective, False + details.update( + { + "pending_action_selection_todo_id": normalized_todo_id, + "pending_action_selection_state": "deferred", + "pending_action_selection_reason": normalized_reason, + "settlement_effect_id": "", + "todo_id": "", + "replan_obligation_id": "", + } + ) + source_event_id = str(effective.get("event_id") or "").strip() or None + retained = build_rollout_event( + goal_id=goal_id, + event_kind="quota_should_run", + agent_id=agent_id, + run_id=turn_instance_id, + status="action_selection_deferred", + summary=( + "heartbeat explicit action selection retained without settlement " + f"binding for turn={turn_instance_id}" + ), + source_event_id=source_event_id, + caused_by=source_event_id, + details=details, + ) + with log_path.open("a", encoding="utf-8") as handle: + handle.write( + json.dumps(retained, sort_keys=True, ensure_ascii=False) + "\n" + ) + return retained, True + + def upgrade_identityless_heartbeat_receipt( runtime_root: Path, *, @@ -410,6 +508,21 @@ def heartbeat_receipt_view( receipt["semantic_replan_obligation_id"] = ( semantic_replan_obligation_id ) + pending_action_todo_id = heartbeat_receipt_pending_action_todo_id(event) + if pending_action_todo_id: + receipt["pending_action_selection"] = { + "todo_id": pending_action_todo_id, + "state": str( + details.get("pending_action_selection_state") or "deferred" + ), + "reason": str( + details.get("pending_action_selection_reason") + or "current_delivery_gate" + ), + "settlement_bound": ( + details.get("pending_action_selection_settlement_bound") is True + ), + } return receipt diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index d364441f08..486da1f3bb 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -20,6 +20,11 @@ build_interaction_contract, build_protocol_action_packet, ) +from ..work_items.action_portfolio import reconcile_retained_action_selection +from ..work_items.autonomous_replan_obligation import ( + replan_obligation_id_from_packet, +) +from ..todos.contract import normalize_todo_id from ..scheduler.execution_context import ( SchedulerExecutionContextResolution, resolve_scheduler_execution_context, @@ -33,6 +38,98 @@ BoundedResearchFrontierProjector = Callable[..., Mapping[str, Any] | None] +def _apply_retained_action_selection_reentry( + payload: dict[str, Any], + *, + retained_todo_id: str | None, + available_capabilities: list[str] | None, + scheduler_execution_context: ( + Mapping[str, Any] | SchedulerExecutionContextResolution | None + ), + turn_instance_id: str | None, + runtime_root: Path, +) -> None: + """Fence a no-argument reentry with its last explicit Todo choice.""" + + normalized_retained = normalize_todo_id(retained_todo_id) + if normalized_retained is None: + return + selected = ( + payload.get("selected_todo") + if isinstance(payload.get("selected_todo"), Mapping) + else {} + ) + projected_todo_id = normalize_todo_id(selected.get("todo_id")) + verdict = reconcile_retained_action_selection( + retained_todo_id=normalized_retained, + projected_todo_id=projected_todo_id, + effective_action=str(payload.get("effective_action") or ""), + replan_obligation_id=replan_obligation_id_from_packet( + payload.get("replan_action_packet") + ), + ) + payload["retained_action_selection"] = verdict + disposition = verdict.get("disposition") + if disposition == "preserve_retained_todo": + return + if disposition == "bind_autonomous_replan": + payload.pop("selected_todo", None) + payload.pop("todo_id", None) + payload.pop("agent_lane_next_action", None) + payload["deferred_action_selection"] = { + "todo_id": normalized_retained, + "reason": "autonomous_replan_preemption", + "resume": "fresh_turn_after_replan_closeout", + } + elif disposition == "require_explicit_selection": + payload.update( + { + "ok": False, + "decision": "skip", + "should_run": False, + "effective_action": EffectiveAction.QUOTA_SKIP.value, + "normal_delivery_allowed": False, + "recovery_delivery_allowed": False, + "self_repair_allowed": False, + "state": "action_selection_required", + "reason": ( + "the current projected default differs from the explicit " + "Todo retained by this Turn" + ), + "recommended_action": ( + "rerun quota should-run with the same --turn-instance-id " + "and an explicit eligible --todo-id" + ), + } + ) + obligation = ( + dict(payload.get("execution_obligation") or {}) + if isinstance(payload.get("execution_obligation"), Mapping) + else {} + ) + obligation.update( + must_attempt_work=False, + delivery_allowed=False, + reason=payload["recommended_action"], + ) + payload["execution_obligation"] = obligation + payload.pop("selected_todo", None) + payload.pop("todo_id", None) + payload.pop("agent_lane_next_action", None) + else: + raise RuntimeError( + "TypeScript retained action-selection disposition is unsupported" + ) + payload["interaction_contract"] = build_interaction_contract( + payload, + available_capabilities=available_capabilities, + scheduler_execution_context=scheduler_execution_context, + turn_instance_id=turn_instance_id, + runtime_root=str(runtime_root), + ) + payload["protocol_action_packet"] = build_protocol_action_packet(payload) + + def _fresh_read_covers_all_pending_material( dispatch: Mapping[str, Any] | None, projector: Callable[..., dict[str, Any]] | None, @@ -392,6 +489,7 @@ def build_live_quota_should_run_decision( receipt_bound_todo_id: str | None = None, requested_action_todo_id: str | None = None, receipt_bound_replan_obligation_id: str | None = None, + retained_action_selection_todo_id: str | None = None, turn_instance_id: str | None = None, interaction_projection_hooks: Sequence[InteractionProjectionHookRegistration] | None = None, @@ -479,6 +577,14 @@ def build_live_quota_should_run_decision( turn_instance_id=turn_instance_id, runtime_root=runtime_root, ) + _apply_retained_action_selection_reentry( + payload, + retained_todo_id=retained_action_selection_todo_id, + available_capabilities=available_capabilities, + scheduler_execution_context=resolved_context, + turn_instance_id=turn_instance_id, + runtime_root=runtime_root, + ) remembered_runtime = (payload.get("agent_identity") or {}).get( "runtime_available_capabilities" ) diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index 986388e99c..962e70543f 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -1463,6 +1463,17 @@ function payloadFor( if (request.turn_instance_id) { payload.turn_instance_id = request.turn_instance_id; payload.replayed = options.replayed; + if (request.execute) { + payload.turn_continuation = { + schema_version: "quota_turn_continuation_v0", + current_turn_settled: true, + same_turn_independent_settlement_allowed: false, + next_turn_required: true, + next_action: + "rerun quota should-run with a fresh --turn-instance-id before independent work", + reason: "the committed monitor-poll is this Turn's single settlement identity", + }; + } } if (request.status_reload_warning) { payload.status_reload_warning = request.status_reload_warning; diff --git a/loopx/control_plane/quota/settlement_cli.py b/loopx/control_plane/quota/settlement_cli.py index 4722bcdd74..c7722218b8 100644 --- a/loopx/control_plane/quota/settlement_cli.py +++ b/loopx/control_plane/quota/settlement_cli.py @@ -309,6 +309,12 @@ def quota_rollout_details( semantic_replan_obligation_id = replan_obligation_id_from_packet( payload.get("replan_action_packet") ) + retained_selection = ( + payload.get("retained_action_selection") + if isinstance(payload.get("retained_action_selection"), Mapping) + else {} + ) + retained_disposition = str(retained_selection.get("disposition") or "") workspace_causality = build_delivery_workspace_causality(selected_todo) interaction = ( payload.get("interaction_contract") @@ -357,6 +363,23 @@ def quota_rollout_details( "quiet_noop_allowed": bool(agent_channel.get("quiet_noop_allowed")), "closeout_required": closeout_required, } + retained_todo_id = normalize_todo_id(retained_selection.get("retained_todo_id")) + if retained_todo_id: + retained_bound = retained_disposition == "preserve_retained_todo" + details.update( + { + "pending_action_selection_todo_id": retained_todo_id, + "pending_action_selection_state": ( + "bound" if retained_bound else "deferred_to_fresh_turn" + ), + "pending_action_selection_reason": ( + "explicit_choice_preserved" + if retained_bound + else "autonomous_replan_preemption" + ), + "pending_action_selection_settlement_bound": retained_bound, + } + ) if workspace_causality: details.update(delivery_workspace_causality_event_fields(workspace_causality)) return details diff --git a/loopx/control_plane/work_items/action_portfolio.py b/loopx/control_plane/work_items/action_portfolio.py index c301eda3fc..1d6fef2e26 100644 --- a/loopx/control_plane/work_items/action_portfolio.py +++ b/loopx/control_plane/work_items/action_portfolio.py @@ -20,6 +20,12 @@ ACTION_SELECTION_QUALIFICATION_SCHEMA_VERSION = ACTION_PORTFOLIO_SELECTION_RESULT_SCHEMA QUOTA_PLANNING_PACKET_REQUEST_SCHEMA_VERSION = ACTION_PORTFOLIO_PLANNING_PACKET_REQUEST_SCHEMA QUOTA_PLANNING_PACKET_SCHEMA_VERSION = ACTION_PORTFOLIO_PLANNING_PACKET_RESULT_SCHEMA +RETAINED_ACTION_SELECTION_REENTRY_REQUEST_SCHEMA_VERSION = ( + "retained_action_selection_reentry_request_v0" +) +RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION = ( + "retained_action_selection_reentry_v0" +) def _compact_candidate(value: Mapping[str, Any]) -> dict[str, Any] | None: @@ -167,3 +173,37 @@ def qualify_action_selection_from_inventory( normal_delivery_allowed=normal_delivery_allowed, delivery_preemptions=delivery_preemptions, ) + + +def reconcile_retained_action_selection( + *, + retained_todo_id: str, + projected_todo_id: str | None, + effective_action: str, + replan_obligation_id: str | None, +) -> dict[str, Any]: + """Ask the typed owner whether a reentry may bind its new projection.""" + + try: + result = effect_runtime_result( + "work_item.action_selection.reconcile_retained", + { + "schema_version": ( + RETAINED_ACTION_SELECTION_REENTRY_REQUEST_SCHEMA_VERSION + ), + "retained_todo_id": retained_todo_id, + "projected_todo_id": projected_todo_id, + "effective_action": effective_action, + "replan_obligation_id": replan_obligation_id, + }, + ) + except EffectRuntimeRejected as exc: + raise ValueError(str(exc)) from None + if not isinstance(result, Mapping) or ( + result.get("schema_version") + != RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION + ): + raise RuntimeError( + "TypeScript retained action-selection reentry shape mismatch" + ) + return dict(result) diff --git a/loopx/control_plane/work_items/action_portfolio.ts b/loopx/control_plane/work_items/action_portfolio.ts index 9db9d039df..9c67b34dc3 100644 --- a/loopx/control_plane/work_items/action_portfolio.ts +++ b/loopx/control_plane/work_items/action_portfolio.ts @@ -32,6 +32,10 @@ export const ACTION_SELECTION_QUALIFICATION_SCHEMA_VERSION = ACTION_PORTFOLIO_SELECTION_RESULT_SCHEMA; export const ACTION_SELECTION_QUALIFICATION_REQUEST_SCHEMA_VERSION = ACTION_PORTFOLIO_SELECTION_REQUEST_SCHEMA; +export const RETAINED_ACTION_SELECTION_REENTRY_REQUEST_SCHEMA_VERSION = + "retained_action_selection_reentry_request_v0"; +export const RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION = + "retained_action_selection_reentry_v0"; export const QUOTA_PLANNING_PACKET_SCHEMA_VERSION = ACTION_PORTFOLIO_PLANNING_PACKET_RESULT_SCHEMA; export const QUOTA_PLANNING_PACKET_REQUEST_SCHEMA_VERSION = @@ -411,3 +415,72 @@ export function qualifyActionSelection(value: unknown): ActionSelectionQualifica }, }; } + +/** + * Reconcile a previously explicit, still-unbound Todo choice on guard reentry. + * + * A projected default is recommendation authority only. It may inherit the + * retained choice when the identities agree, but it may never replace that + * choice. A hard autonomous replan is different: it is the current control + * lane and therefore receives its own replan settlement identity while the + * Todo choice remains deferred for a fresh Turn. + */ +export function reconcileRetainedActionSelection(value: unknown): JsonObject { + const request = requireJsonObject( + value, + "retained_action_selection_reentry_request", + ); + if ( + request.schema_version !== + RETAINED_ACTION_SELECTION_REENTRY_REQUEST_SCHEMA_VERSION + ) { + throw new EffectRuntimeRequestError( + `retained_action_selection_reentry_request.schema_version must be ${RETAINED_ACTION_SELECTION_REENTRY_REQUEST_SCHEMA_VERSION}`, + ); + } + const retainedTodoId = requireNonEmptyString( + request.retained_todo_id, + "retained_action_selection_reentry_request.retained_todo_id", + ); + const projectedTodoId = optionalNonEmptyString( + request.projected_todo_id, + "retained_action_selection_reentry_request.projected_todo_id", + ); + const effectiveAction = requireNonEmptyString( + request.effective_action, + "retained_action_selection_reentry_request.effective_action", + ); + const replanObligationId = optionalNonEmptyString( + request.replan_obligation_id, + "retained_action_selection_reentry_request.replan_obligation_id", + ); + + if (projectedTodoId === retainedTodoId) { + return { + schema_version: RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION, + disposition: "preserve_retained_todo", + retained_todo_id: retainedTodoId, + projected_todo_id: projectedTodoId, + }; + } + if ( + effectiveAction === "autonomous_replan_required" && + replanObligationId !== null + ) { + return { + schema_version: RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION, + disposition: "bind_autonomous_replan", + retained_todo_id: retainedTodoId, + projected_todo_id: projectedTodoId, + replan_obligation_id: replanObligationId, + continuation: "fresh_turn_after_replan_closeout", + }; + } + return { + schema_version: RETAINED_ACTION_SELECTION_REENTRY_SCHEMA_VERSION, + disposition: "require_explicit_selection", + retained_todo_id: retainedTodoId, + projected_todo_id: projectedTodoId, + reason: "projected_default_differs_from_retained_explicit_choice", + }; +} diff --git a/tests/control_plane/test_monitor_followthrough_contract.py b/tests/control_plane/test_monitor_followthrough_contract.py index 9118830c0d..cbf6da0824 100644 --- a/tests/control_plane/test_monitor_followthrough_contract.py +++ b/tests/control_plane/test_monitor_followthrough_contract.py @@ -595,6 +595,17 @@ def test_same_turn_material_monitor_poll_is_no_spend_closeout_before_successor( ) successor_id = poll["successor_todo_ids"][0] assert poll["after"]["selected_todo"]["todo_id"] == admitted["todo_id"] + assert poll["turn_continuation"] == { + "schema_version": "quota_turn_continuation_v0", + "current_turn_settled": True, + "same_turn_independent_settlement_allowed": False, + "next_turn_required": True, + "next_action": ( + "rerun quota should-run with a fresh --turn-instance-id before " + "independent work" + ), + "reason": "the committed monitor-poll is this Turn's single settlement identity", + } # The production CLI must not confuse an observation row with a committed # closeout. Keep the exact guard and Todo fixed while corrupting only the diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index 5f8f612c98..5d34344027 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -3396,8 +3396,13 @@ def test_selection_added_after_pending_guard_reports_final_boundary( added["todo_id"] ) assert rejected["normal_delivery_allowed"] is False - assert rejected["heartbeat_receipt"]["status"] == "replayed" - assert _heartbeat_receipt_events(runtime, turn_instance_id) == [first_receipt] + assert rejected["heartbeat_receipt"]["status"] == "selection_retained" + assert rejected["heartbeat_receipt"]["pending_action_selection"]["todo_id"] == ( + added["todo_id"] + ) + retained_events = _heartbeat_receipt_events(runtime, turn_instance_id) + assert retained_events[0] == first_receipt + assert len(retained_events) == 2 repeated_rc, repeated = _run_cli( registry_path, @@ -3460,12 +3465,15 @@ def test_pending_action_selection_does_not_preempt_newly_due_monitor( assert selected["action_selection_qualification"]["reason"] == ( "blocking_work_lane" ) - assert selected["heartbeat_receipt"]["status"] == "replayed" - assert selected["rollout_event"]["appended"] is False + assert selected["heartbeat_receipt"]["status"] == "selection_retained" + assert selected["heartbeat_receipt"]["pending_action_selection"]["todo_id"] == ( + ALTERNATIVE_TODO_ID + ) + assert selected["rollout_event"]["appended"] is True events = _heartbeat_receipt_events(runtime, turn_instance_id) - assert len(events) == 1 - assert not events[0]["details"].get("todo_id") - assert not events[0]["details"].get("settlement_effect_id") + assert len(events) == 2 + assert all(not event["details"].get("todo_id") for event in events) + assert all(not event["details"].get("settlement_effect_id") for event in events) def test_pending_action_selection_reports_autonomous_replan_preemption( @@ -3514,9 +3522,12 @@ def test_pending_action_selection_reports_autonomous_replan_preemption( "reason": "autonomous_replan", "delivery_preemptions": ["autonomous_replan", "delivery_not_allowed"], } - assert selected["heartbeat_receipt"]["status"] == "replayed" - assert selected["rollout_event"]["appended"] is False - assert _heartbeat_receipt_count(runtime, turn_instance_id) == 1 + assert selected["heartbeat_receipt"]["status"] == "selection_retained" + assert selected["heartbeat_receipt"]["pending_action_selection"]["todo_id"] == ( + ALTERNATIVE_TODO_ID + ) + assert selected["rollout_event"]["appended"] is True + assert _heartbeat_receipt_count(runtime, turn_instance_id) == 2 def test_due_monitor_auxiliary_context_has_typed_selection_rejection( @@ -4017,12 +4028,15 @@ def test_pending_action_selection_does_not_commit_after_new_user_gate( assert selected["action_selection_qualification"]["reason"] == ( "delivery_not_allowed" ) - assert selected["heartbeat_receipt"]["status"] == "replayed" - assert selected["rollout_event"]["appended"] is False + assert selected["heartbeat_receipt"]["status"] == "selection_retained" + assert selected["heartbeat_receipt"]["pending_action_selection"]["todo_id"] == ( + ALTERNATIVE_TODO_ID + ) + assert selected["rollout_event"]["appended"] is True events = _heartbeat_receipt_events(runtime, turn_instance_id) - assert len(events) == 1 - assert not events[0]["details"].get("todo_id") - assert not events[0]["details"].get("settlement_effect_id") + assert len(events) == 2 + assert all(not event["details"].get("todo_id") for event in events) + assert all(not event["details"].get("settlement_effect_id") for event in events) def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( diff --git a/tests/control_plane/test_selection_replan_reentry.py b/tests/control_plane/test_selection_replan_reentry.py index 58e73eafb5..d7d18a8a67 100644 --- a/tests/control_plane/test_selection_replan_reentry.py +++ b/tests/control_plane/test_selection_replan_reentry.py @@ -30,8 +30,22 @@ def test_deferred_selection_recovers_same_turn_and_settles_once(tmp_path, bindin selected_id = TODO_ID rc, deferred = _run_cli(registry, runtime, *guard, "--todo-id", selected_id) assert rc == 1 and not deferred["should_run"] - assert deferred["heartbeat_receipt"]["event_id"] == first["heartbeat_receipt"]["event_id"] - assert _heartbeat_receipt_count(runtime, turn) == 1 + selection_deferred = deferred["action_selection_qualification"]["state"] == "deferred" + if selection_deferred: + assert deferred["heartbeat_receipt"]["event_id"] != first["heartbeat_receipt"]["event_id"] + assert deferred["heartbeat_receipt"]["status"] == "selection_retained" + assert deferred["heartbeat_receipt"]["pending_action_selection"] == { + "todo_id": selected_id, + "state": "deferred", + "reason": deferred["action_selection_qualification"]["reason"], + "settlement_bound": False, + } + else: + assert deferred["heartbeat_receipt"]["event_id"] == first["heartbeat_receipt"]["event_id"] + assert deferred["heartbeat_receipt"]["status"] == "replayed" + assert "pending_action_selection" not in deferred["heartbeat_receipt"] + expected_before_resume = 2 if selection_deferred else 1 + assert _heartbeat_receipt_count(runtime, turn) == expected_before_resume channel = deferred["interaction_contract"]["cli_channel"] assert "settlement_plan" not in channel assert "replan_settlement_contract" not in channel @@ -50,7 +64,7 @@ def test_deferred_selection_recovers_same_turn_and_settles_once(tmp_path, bindin [recovery_preview] = envelope["writeback"]["next_cli_actions"] assert recovery_preview.startswith("loopx ") assert not envelope["writeback"]["spend_after_validation"] - assert _heartbeat_receipt_count(runtime, turn) == 1 + assert _heartbeat_receipt_count(runtime, turn) == expected_before_resume argv = shlex.split(command) assert argv[argv.index("--turn-instance-id") + 1] == turn assert "--todo-id" not in argv and "--replan-obligation-id" not in argv @@ -63,7 +77,15 @@ def test_deferred_selection_recovers_same_turn_and_settles_once(tmp_path, bindin identity = receipt["settlement_identity"] assert identity["turn_instance_id"] == turn assert ("todo_id" in identity) == (binding == "todo") - assert _heartbeat_receipt_count(runtime, turn) == 2 + if selection_deferred: + assert receipt["pending_action_selection"]["todo_id"] == selected_id + assert receipt["pending_action_selection"]["settlement_bound"] is ( + binding == "todo" + ) + else: + assert "pending_action_selection" not in receipt + expected_after_resume = expected_before_resume + 1 + assert _heartbeat_receipt_count(runtime, turn) == expected_after_resume cli = resumed["interaction_contract"]["cli_channel"] assert cli["settlement_plan"]["identity"] == identity refresh = next(c for c in cli["next_cli_actions"] if "refresh-state" in c) @@ -90,5 +112,52 @@ def test_deferred_selection_recovers_same_turn_and_settles_once(tmp_path, bindin assert settled["heartbeat_receipt"]["settlement_identity"] == identity rc, conflict = _run_cli(registry, runtime, *guard, "--todo-id", "todo_another_selection") assert rc == 1 and conflict["ok"] is False - assert _heartbeat_receipt_count(runtime, turn) == 2 + assert _heartbeat_receipt_count(runtime, turn) == expected_after_resume assert _spend_run_count(runtime) == 1 + + +def test_reentry_never_replaces_retained_selection_with_recommended_todo(tmp_path): + project, runtime, registry = _write_fixture(tmp_path) + _configure_selectable_alternative(project) + turn = "turn-selection-recommendation-drift" + guard = ("quota", "should-run", "--codex-app", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--turn-instance-id", turn, "--scan-path", str(project)) + rc, first = _run_cli(registry, runtime, *guard) + assert rc == 0, first + assert first["interaction_contract"]["cli_channel"]["selection_required"] + + # The original explicit choice leaves the refreshed frontier while a + # different recommended Todo becomes visible under a hard replan. + _configure_selected_todo_replan_fixture(project, registry) + retained_todo_id = "todo_chain_000000000001" + rc, deferred = _run_cli( + registry, runtime, *guard, "--todo-id", retained_todo_id + ) + assert rc == 1, deferred + assert deferred["action_selection_qualification"]["state"] == "deferred" + assert deferred["heartbeat_receipt"]["pending_action_selection"]["todo_id"] == retained_todo_id + [command] = deferred["interaction_contract"]["cli_channel"]["next_cli_actions"] + + rc, resumed = _run_generated_cli(command, registry_path=registry) + assert rc == 0, resumed + assert resumed["decision"] == "autonomous_replan_required" + assert resumed.get("selected_todo") is None + retained = resumed["retained_action_selection"] + assert retained["disposition"] == "bind_autonomous_replan" + assert retained["retained_todo_id"] == retained_todo_id + assert retained["projected_todo_id"] == SELECTED_REPLAN_TODO_ID + identity = resumed["heartbeat_receipt"]["settlement_identity"] + assert identity["binding_kind"] == "autonomous_replan" + assert "todo_id" not in identity + assert resumed["heartbeat_receipt"]["pending_action_selection"] == { + "todo_id": retained_todo_id, + "state": "deferred_to_fresh_turn", + "reason": "autonomous_replan_preemption", + "settlement_bound": False, + } + plan = resumed["interaction_contract"]["cli_channel"]["settlement_plan"] + assert plan["identity"] == identity + assert all( + "--replan-obligation-id" in action and "--todo-id" not in action + for action in resumed["interaction_contract"]["cli_channel"]["next_cli_actions"] + ) diff --git a/tests/control_plane_ts/action_portfolio.test.ts b/tests/control_plane_ts/action_portfolio.test.ts index 0b9930a451..fe186edd7e 100644 --- a/tests/control_plane_ts/action_portfolio.test.ts +++ b/tests/control_plane_ts/action_portfolio.test.ts @@ -7,6 +7,7 @@ import { QUOTA_PLANNING_PACKET_REQUEST_SCHEMA_VERSION, projectQuotaActionPortfolio, qualifyActionSelection, + reconcileRetainedActionSelection, } from "../../loopx/control_plane/work_items/action_portfolio.ts"; import { PLANNING_HORIZON_REQUEST_SCHEMA_VERSION, @@ -392,3 +393,47 @@ test("pending selection explains an auxiliary monitor outside the advancement la reason: "auxiliary_monitor_not_selectable_in_advancement_lane", }); }); + +test("retained explicit selection cannot be replaced by a projected default", () => { + const base = { + schema_version: "retained_action_selection_reentry_request_v0", + retained_todo_id: "todo_explicit001", + effective_action: "autonomous_replan_required", + replan_obligation_id: "replan-0123456789abcdef", + }; + + assert.deepEqual(reconcileRetainedActionSelection({ + ...base, + projected_todo_id: "todo_explicit001", + }), { + schema_version: "retained_action_selection_reentry_v0", + disposition: "preserve_retained_todo", + retained_todo_id: "todo_explicit001", + projected_todo_id: "todo_explicit001", + }); + + assert.deepEqual(reconcileRetainedActionSelection({ + ...base, + projected_todo_id: "todo_recommended001", + }), { + schema_version: "retained_action_selection_reentry_v0", + disposition: "bind_autonomous_replan", + retained_todo_id: "todo_explicit001", + projected_todo_id: "todo_recommended001", + replan_obligation_id: "replan-0123456789abcdef", + continuation: "fresh_turn_after_replan_closeout", + }); + + assert.deepEqual(reconcileRetainedActionSelection({ + ...base, + effective_action: "normal_run", + replan_obligation_id: null, + projected_todo_id: "todo_recommended001", + }), { + schema_version: "retained_action_selection_reentry_v0", + disposition: "require_explicit_selection", + retained_todo_id: "todo_explicit001", + projected_todo_id: "todo_recommended001", + reason: "projected_default_differs_from_retained_explicit_choice", + }); +}); From 1455cdb67f53715de35c94d184e138002a45ce57 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 05:58:04 +0800 Subject: [PATCH 2/6] fix(quota): preserve auxiliary turn settlement Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota-monitor-observation-receipt-v0.md | 19 +- docs/reference/protocols/turn-envelope-v0.md | 14 +- loopx/cli_commands/quota.py | 265 +++++++++++------- loopx/control_plane/quota/live_decision.py | 49 ++-- .../quota/monitor_poll_commit.ts | 47 +++- .../work_items/action_portfolio.ts | 24 ++ loopx/semantics/vocabulary_v0.json | 1 + .../test_monitor_followthrough_contract.py | 1 + .../test_monitor_observation_admission.py | 18 ++ .../control_plane_ts/action_portfolio.test.ts | 23 ++ .../quota_monitor_poll_commit.test.ts | 9 + 11 files changed, 314 insertions(+), 156 deletions(-) diff --git a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md index 9481dca9eb..6c52411415 100644 --- a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md +++ b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md @@ -34,9 +34,13 @@ advancement work remains active. remains selected. A material observation may create its independently routed successor through the existing monitor contract, but it still does not replace the Turn's settlement identity. -- An executed turn-scoped poll returns `turn_continuation` declaring the current - Turn settled and requiring a fresh `--turn-instance-id` before unrelated work. - A no-spend closeout does not make the settled Turn reusable. +- An executed turn-scoped poll returns `turn_continuation`. An exact match + between `settlement_todo_id` and the observed `todo_id` closes the no-spend + monitor Turn and requires a fresh `--turn-instance-id`. A different admitted + monitor Todo is auxiliary: it records the observation but leaves the original + advancement Turn open for its durable writeback and single spend. Without an + exact or typed auxiliary binding the response fails closed from claiming the + Turn settled. ### Acceptance @@ -187,9 +191,12 @@ using a complete read-only snapshot with disposable File/SQLite/PostgreSQL arms. 替换为本 Turn 的结算 Todo;多个辅助回执也绝不改变既有结算身份。 - 辅助观察无变化后,原 advancement Todo 继续保持选中;若观察发生重大变化, 可按既有 monitor 契约创建独立路由的 successor,但仍不替换本 Turn 的结算身份。 -- 执行成功的 turn-scoped poll 会返回 `turn_continuation`,明确当前 Turn 已结算; - 开始无关工作前必须使用新的 `--turn-instance-id`。不计费 closeout 不代表原 Turn - 可以再次绑定另一份独立工作。 +- 执行成功的 turn-scoped poll 会返回 `turn_continuation`。仅当 + `settlement_todo_id` 与被观察的 `todo_id` 精确一致时,才完成该 monitor Turn 的 + 不计费结算,并要求后续使用新的 `--turn-instance-id`。不同但已准入的 monitor + Todo 属于辅助观察:只写观察回执,原 advancement Turn 仍需完成 durable + writeback 与唯一一次 spend。既非精确匹配、也无 typed auxiliary binding 时,响应 + 必须失败关闭,不能宣称 Turn 已结算。 ### 验收 diff --git a/docs/reference/protocols/turn-envelope-v0.md b/docs/reference/protocols/turn-envelope-v0.md index 2e65610b14..206de32e5f 100644 --- a/docs/reference/protocols/turn-envelope-v0.md +++ b/docs/reference/protocols/turn-envelope-v0.md @@ -108,12 +108,14 @@ worktree recovery instruction. Moving to an independent worktree and rerunning the guard with the same Turn id resumes the selected Todo; the wrapper must not rewrite this recoverable state as a settlement-identity conflict. -An executed, turn-scoped `quota monitor-poll` is a no-spend closeout, but it is -still that Turn's single settlement identity. Its response therefore includes -`turn_continuation.next_turn_required=true` and an instruction to rerun -`quota should-run` with a fresh `--turn-instance-id` before starting unrelated -work. A host must not interpret successful no-spend closeout as permission to -bind an independent Todo in the already-settled Turn. +An executed, turn-scoped `quota monitor-poll` is a no-spend closeout only when +its observed Todo exactly matches the Turn's `settlement_todo_id`. That response +includes `turn_continuation.next_turn_required=true` and requires a fresh Turn +before unrelated work. An admitted auxiliary monitor uses its own observed Todo +while retaining the advancement Todo as `settlement_todo_id`; its continuation +keeps `current_turn_settled=false` and `next_turn_required=false` so the original +writeback and spend can finish. A missing exact or typed auxiliary binding never +claims settlement. Portfolio v2 preserves v1's selection policy, candidate ordering, and settlement rules, and adds an optional `continuation_hint` to each suggested diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 95c335fd72..0abef659e6 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -15,6 +15,7 @@ repository_delivery_interaction_hook, ) from ..control_plane.effect_runtime import EffectRuntimeRejected +from ..control_plane.capability_hooks import InteractionProjectionHookRegistration from ..control_plane.quota.cli_projection import ( compact_quota_monitor_poll_cli_payload, compact_quota_should_run_cli_payload, @@ -404,6 +405,136 @@ def _commit_requested_action_selection( selected_todo["selection_binding"] = "heartbeat_receipt" +def _retain_deferred_action_selection( + payload: Mapping[str, object], + args: argparse.Namespace, + *, + runtime_root: Path, + heartbeat_turn_id: str | None, + existing: dict[str, object] | None, +) -> tuple[dict[str, object] | None, str, bool]: + """Append a deferred explicit choice without granting settlement authority.""" + + requested_todo_id = _requested_quota_action_todo_id(args) + qualification = payload.get("action_selection_qualification") + if ( + not heartbeat_turn_id + or existing is None + or requested_todo_id is None + or not isinstance(qualification, Mapping) + or qualification.get("state") != "deferred" + ): + return existing, "replayed", False + retained, appended = retain_pending_heartbeat_action_selection( + runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=heartbeat_turn_id, + todo_id=requested_todo_id, + reason=str(qualification.get("reason") or "current_delivery_gate"), + ) + return retained, "selection_retained" if appended else "replayed", appended + + +def _record_automatic_heartbeat_stall( + payload: dict[str, object], + status_payload: dict[str, object], + args: argparse.Namespace, + *, + registry_path: Path, + runtime_root_arg: str | None, + context: QuotaCommandContext, + interaction_projection_hooks: tuple[InteractionProjectionHookRegistration, ...], + receipt_bound_todo_id: str | None, + receipt_bound_replan_obligation_id: str | None, + receipt_pending_action_todo_id: str | None, + cache_metadata: object, +) -> tuple[dict[str, object], dict[str, object], object, str]: + """Commit and reproject the automatic no-spend heartbeat observation.""" + + turn_id = context.heartbeat_turn_id + if turn_id is None: + return payload, status_payload, cache_metadata, "not_applicable" + existing_stall = find_quota_monitor_poll_turn( + context.runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=turn_id, + ) + if ( + payload.get("effective_action") + != EffectiveAction.MONITOR_QUIET_SKIP.value + and existing_stall is None + ): + return payload, status_payload, cache_metadata, "not_applicable" + poll = record_quota_monitor_poll( + status_payload, + goal_id=args.goal_id, + registry_path=registry_path, + execute=True, + source="heartbeat", + agent_id=args.agent_id, + available_capabilities=args.available_capabilities, + turn_instance_id=turn_id, + scheduler_execution_context=context.scheduler_context, + operator_inbox_urgency_projector=context.operator_inbox_urgency_projector, + bounded_research_frontier_projector=project_live_explore_composition_frontier, + ) + if not poll.get("ok"): + raise RuntimeError( + "heartbeat stall observation writeback failed: " + f"{poll.get('reason') or 'missing follow-up quota decision'}" + ) + reloaded = collect_status( + registry_path=registry_path, + runtime_root_override=runtime_root_arg, + scan_roots=context.scan_roots, + limit=context.status_limit, + goal_id=context.status_goal_id, + available_capabilities=args.available_capabilities, + ) + rebuilt = build_live_quota_should_run_decision( + reloaded, + goal_id=args.goal_id, + agent_id=args.agent_id, + available_capabilities=args.available_capabilities, + include_scheduler_detail="scheduler" in context.detail_sections, + include_agent_todo_detail=( + "agent-todos" in context.detail_sections + and not bool(getattr(args, "turn_envelope", False)) + ), + codex_app_current_rrule=args.app_automation_current_rrule, + registry_path=registry_path, + runtime_root=context.runtime_root, + host_observation_resolver=resolve_codex_app_automation_rrule, + scheduler_execution_context=context.scheduler_context, + operator_inbox_urgency_projector=context.operator_inbox_urgency_projector, + bounded_research_frontier_projector=project_live_explore_composition_frontier, + receipt_bound_todo_id=receipt_bound_todo_id, + receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, + retained_action_selection_todo_id=( + receipt_pending_action_todo_id + if _requested_quota_action_todo_id(args) is None + and receipt_bound_todo_id is None + and receipt_bound_replan_obligation_id is None + else None + ), + turn_instance_id=turn_id, + interaction_projection_hooks=interaction_projection_hooks, + ) + rebuilt["heartbeat_stall_writeback"] = { + "turn_instance_id": turn_id, + "status": "replayed" if poll.get("replayed") else "appended", + "generated_at": poll.get("generated_at"), + } + return ( + rebuilt, + reloaded, + None, + "replayed" if poll.get("replayed") else "appended", + ) + + def _dispatch_quota_turn_start_hooks( args: argparse.Namespace, *, @@ -604,35 +735,17 @@ def handle_quota_command( receipt_bound_todo_id=receipt_bound_todo_id, receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, ) - if ( - heartbeat_turn_id - and action_selection_preflight_failed - and heartbeat_receipt_existing is not None - and (requested_todo_id := _requested_quota_action_todo_id(args)) - and isinstance( - qualification := payload.get("action_selection_qualification"), - Mapping, - ) - and qualification.get("state") == "deferred" - ): - qualification_reason = str( - qualification.get("reason") or "current_delivery_gate" - ) + if action_selection_preflight_failed: ( heartbeat_receipt_existing, + heartbeat_receipt_existing_status, retained_selection_appended, - ) = retain_pending_heartbeat_action_selection( - runtime_root, - goal_id=args.goal_id, - agent_id=args.agent_id, - turn_instance_id=heartbeat_turn_id, - todo_id=requested_todo_id, - reason=qualification_reason, - ) - heartbeat_receipt_existing_status = ( - "selection_retained" - if retained_selection_appended - else "replayed" + ) = _retain_deferred_action_selection( + payload, + args, + runtime_root=runtime_root, + heartbeat_turn_id=heartbeat_turn_id, + existing=heartbeat_receipt_existing, ) heartbeat_receipt_existing_appended = retained_selection_appended if heartbeat_turn_id: @@ -653,88 +766,26 @@ def handle_quota_command( existing=heartbeat_receipt_existing, ) else: - existing_stall = find_quota_monitor_poll_turn( - runtime_root, - goal_id=args.goal_id, - agent_id=args.agent_id, - turn_instance_id=heartbeat_turn_id, + ( + payload, + status_payload, + cache_metadata, + heartbeat_stall_observation, + ) = _record_automatic_heartbeat_stall( + payload, + status_payload, + args, + registry_path=registry_path, + runtime_root_arg=runtime_root_arg, + context=context, + interaction_projection_hooks=interaction_projection_hooks, + receipt_bound_todo_id=receipt_bound_todo_id, + receipt_bound_replan_obligation_id=( + receipt_bound_replan_obligation_id + ), + receipt_pending_action_todo_id=receipt_pending_action_todo_id, + cache_metadata=cache_metadata, ) - if ( - payload.get("effective_action") == EffectiveAction.MONITOR_QUIET_SKIP.value - or existing_stall is not None - ): - poll = record_quota_monitor_poll( - status_payload, - goal_id=args.goal_id, - registry_path=registry_path, - execute=True, - source="heartbeat", - agent_id=args.agent_id, - available_capabilities=args.available_capabilities, - turn_instance_id=heartbeat_turn_id, - scheduler_execution_context=scheduler_context, - operator_inbox_urgency_projector=operator_inbox_urgency_projector, - bounded_research_frontier_projector=( - project_live_explore_composition_frontier - ), - ) - if not poll.get("ok"): - raise RuntimeError( - "heartbeat stall observation writeback failed: " - f"{poll.get('reason') or 'missing follow-up quota decision'}" - ) - status_payload = collect_status( - registry_path=registry_path, - runtime_root_override=runtime_root_arg, - scan_roots=scan_roots, - limit=status_limit, - goal_id=status_goal_id, - available_capabilities=args.available_capabilities, - ) - payload = build_live_quota_should_run_decision( - status_payload, - goal_id=args.goal_id, - agent_id=args.agent_id, - available_capabilities=args.available_capabilities, - include_scheduler_detail="scheduler" in detail_sections, - include_agent_todo_detail=( - "agent-todos" in detail_sections - and not bool(getattr(args, "turn_envelope", False)) - ), - codex_app_current_rrule=args.app_automation_current_rrule, - registry_path=registry_path, - runtime_root=runtime_root, - host_observation_resolver=resolve_codex_app_automation_rrule, - scheduler_execution_context=scheduler_context, - operator_inbox_urgency_projector=operator_inbox_urgency_projector, - bounded_research_frontier_projector=( - project_live_explore_composition_frontier - ), - receipt_bound_todo_id=receipt_bound_todo_id, - receipt_bound_replan_obligation_id=( - receipt_bound_replan_obligation_id - ), - retained_action_selection_todo_id=( - receipt_pending_action_todo_id - if _requested_quota_action_todo_id(args) is None - and receipt_bound_todo_id is None - and receipt_bound_replan_obligation_id is None - else None - ), - turn_instance_id=heartbeat_turn_id, - interaction_projection_hooks=(interaction_projection_hooks), - ) - cache_metadata = None - heartbeat_stall_observation = ( - "replayed" if poll.get("replayed") else "appended" - ) - payload["heartbeat_stall_writeback"] = { - "turn_instance_id": heartbeat_turn_id, - "status": heartbeat_stall_observation, - "generated_at": poll.get("generated_at"), - } - else: - heartbeat_stall_observation = "not_applicable" heartbeat_receipt_ready = True elif args.quota_command == "monitor-poll": payload = record_quota_monitor_poll_for_cli( diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index 486da1f3bb..05dd3c16fa 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -82,40 +82,33 @@ def _apply_retained_action_selection_reentry( "resume": "fresh_turn_after_replan_closeout", } elif disposition == "require_explicit_selection": - payload.update( - { - "ok": False, - "decision": "skip", - "should_run": False, - "effective_action": EffectiveAction.QUOTA_SKIP.value, - "normal_delivery_allowed": False, - "recovery_delivery_allowed": False, - "self_repair_allowed": False, - "state": "action_selection_required", - "reason": ( - "the current projected default differs from the explicit " - "Todo retained by this Turn" - ), - "recommended_action": ( - "rerun quota should-run with the same --turn-instance-id " - "and an explicit eligible --todo-id" - ), - } - ) + projection = verdict.get("projection") + if not isinstance(projection, Mapping): + raise RuntimeError( + "TypeScript retained action-selection projection is missing" + ) + decision_patch = projection.get("decision_patch") + obligation_patch = projection.get("execution_obligation_patch") + clear_fields = projection.get("clear_fields") + if ( + not isinstance(decision_patch, Mapping) + or not isinstance(obligation_patch, Mapping) + or not isinstance(clear_fields, list) + or not all(isinstance(field, str) for field in clear_fields) + ): + raise RuntimeError( + "TypeScript retained action-selection projection is malformed" + ) + payload.update(decision_patch) obligation = ( dict(payload.get("execution_obligation") or {}) if isinstance(payload.get("execution_obligation"), Mapping) else {} ) - obligation.update( - must_attempt_work=False, - delivery_allowed=False, - reason=payload["recommended_action"], - ) + obligation.update(obligation_patch) payload["execution_obligation"] = obligation - payload.pop("selected_todo", None) - payload.pop("todo_id", None) - payload.pop("agent_lane_next_action", None) + for field in clear_fields: + payload.pop(field, None) else: raise RuntimeError( "TypeScript retained action-selection disposition is unsupported" diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index 962e70543f..f908e0bbdb 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -1464,15 +1464,44 @@ function payloadFor( payload.turn_instance_id = request.turn_instance_id; payload.replayed = options.replayed; if (request.execute) { - payload.turn_continuation = { - schema_version: "quota_turn_continuation_v0", - current_turn_settled: true, - same_turn_independent_settlement_allowed: false, - next_turn_required: true, - next_action: - "rerun quota should-run with a fresh --turn-instance-id before independent work", - reason: "the committed monitor-poll is this Turn's single settlement identity", - }; + const settlementTodoId = optionalString( + event.settlement_todo_id, + "monitor_event.settlement_todo_id", + )?.trim() ?? null; + const observedTodoId = optionalString( + event.todo_id, + "monitor_event.todo_id", + )?.trim() ?? null; + const exactSettlement = settlementTodoId !== null && + observedTodoId === settlementTodoId; + const auxiliaryObservation = settlementTodoId !== null && + observedTodoId !== null && observedTodoId !== settlementTodoId; + payload.turn_continuation = exactSettlement + ? { + schema_version: "quota_turn_continuation_v0", + settlement_binding_matches_observation: true, + current_turn_settled: true, + same_turn_independent_settlement_allowed: false, + next_turn_required: true, + next_action: + "rerun quota should-run with a fresh --turn-instance-id before independent work", + reason: "the committed monitor-poll is this Turn's single settlement identity", + } + : { + schema_version: "quota_turn_continuation_v0", + settlement_binding_matches_observation: auxiliaryObservation + ? false + : null, + current_turn_settled: false, + same_turn_independent_settlement_allowed: false, + next_turn_required: false, + next_action: auxiliaryObservation + ? "continue the original advancement settlement in this Turn" + : "obtain a typed settlement binding before claiming this Turn settled", + reason: auxiliaryObservation + ? "the committed monitor-poll is an auxiliary observation and does not settle the advancement Turn" + : "the committed monitor-poll has no exact Todo settlement binding", + }; } } if (request.status_reload_warning) { diff --git a/loopx/control_plane/work_items/action_portfolio.ts b/loopx/control_plane/work_items/action_portfolio.ts index 9c67b34dc3..b4e2c89d6d 100644 --- a/loopx/control_plane/work_items/action_portfolio.ts +++ b/loopx/control_plane/work_items/action_portfolio.ts @@ -24,6 +24,7 @@ import { PLANNING_HORIZON_REQUEST_SCHEMA_VERSION, projectQuotaPlanningHorizon, } from "./planning_horizon.ts"; +import { EffectiveAction } from "../quota/effective_action.generated.ts"; export const ACTION_PORTFOLIO_SCHEMA_VERSION = "quota_action_portfolio_v2"; export const ACTION_PORTFOLIO_REQUEST_SCHEMA_VERSION = @@ -482,5 +483,28 @@ export function reconcileRetainedActionSelection(value: unknown): JsonObject { retained_todo_id: retainedTodoId, projected_todo_id: projectedTodoId, reason: "projected_default_differs_from_retained_explicit_choice", + projection: { + decision_patch: { + ok: false, + decision: "skip", + should_run: false, + effective_action: EffectiveAction.QUOTA_SKIP, + normal_delivery_allowed: false, + recovery_delivery_allowed: false, + self_repair_allowed: false, + state: "action_selection_required", + reason: + "the current projected default differs from the explicit Todo retained by this Turn", + recommended_action: + "rerun quota should-run with the same --turn-instance-id and an explicit eligible --todo-id", + }, + execution_obligation_patch: { + must_attempt_work: false, + delivery_allowed: false, + reason: + "rerun quota should-run with the same --turn-instance-id and an explicit eligible --todo-id", + }, + clear_fields: ["selected_todo", "todo_id", "agent_lane_next_action"], + }, }; } diff --git a/loopx/semantics/vocabulary_v0.json b/loopx/semantics/vocabulary_v0.json index c026aae8c8..5e3cbbe686 100644 --- a/loopx/semantics/vocabulary_v0.json +++ b/loopx/semantics/vocabulary_v0.json @@ -500,6 +500,7 @@ "loopx/control_plane/quota/decision_summary.py::resolve_quota_run_decision", "loopx/control_plane/quota/heartbeat_receipt.py::fail_heartbeat_receipt", "loopx/control_plane/quota/live_decision.py::_apply_pending_capability_intent_precedence", + "loopx/control_plane/work_items/action_portfolio.ts::reconcileRetainedActionSelection", "loopx/control_plane/quota/projection_repair.py::build_boundary_projection_repair_hint", "loopx/control_plane/quota/projection_repair.py::build_state_projection_gap_repair_hint", "loopx/control_plane/quota/settlement_precedence.py::apply_settled_replay_payload_precedence", diff --git a/tests/control_plane/test_monitor_followthrough_contract.py b/tests/control_plane/test_monitor_followthrough_contract.py index cbf6da0824..47d0009b5d 100644 --- a/tests/control_plane/test_monitor_followthrough_contract.py +++ b/tests/control_plane/test_monitor_followthrough_contract.py @@ -597,6 +597,7 @@ def test_same_turn_material_monitor_poll_is_no_spend_closeout_before_successor( assert poll["after"]["selected_todo"]["todo_id"] == admitted["todo_id"] assert poll["turn_continuation"] == { "schema_version": "quota_turn_continuation_v0", + "settlement_binding_matches_observation": True, "current_turn_settled": True, "same_turn_independent_settlement_allowed": False, "next_turn_required": True, diff --git a/tests/control_plane/test_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index 7e6f26c0fe..8549e29c06 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -88,6 +88,24 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( assert poll["before"]["selected_todo"]["todo_id"] == TODO_ID assert poll["after"]["selected_todo"]["todo_id"] == TODO_ID assert poll["material_change"] is False + assert poll["turn_continuation"] == { + "schema_version": "quota_turn_continuation_v0", + "settlement_binding_matches_observation": False, + "current_turn_settled": False, + "same_turn_independent_settlement_allowed": False, + "next_turn_required": False, + "next_action": "continue the original advancement settlement in this Turn", + "reason": ( + "the committed monitor-poll is an auxiliary observation and does not " + "settle the advancement Turn" + ), + } + assert ( + poll["after"]["interaction_contract"]["cli_channel"][ + "spend_after_validation" + ] + is True + ) assert poll_replay_rc == 0, poll_replay assert poll_replay["replayed"] is True assert _classification_count(runtime, "quota_monitor_poll") == 1 diff --git a/tests/control_plane_ts/action_portfolio.test.ts b/tests/control_plane_ts/action_portfolio.test.ts index fe186edd7e..1137ab41d0 100644 --- a/tests/control_plane_ts/action_portfolio.test.ts +++ b/tests/control_plane_ts/action_portfolio.test.ts @@ -435,5 +435,28 @@ test("retained explicit selection cannot be replaced by a projected default", () retained_todo_id: "todo_explicit001", projected_todo_id: "todo_recommended001", reason: "projected_default_differs_from_retained_explicit_choice", + projection: { + decision_patch: { + ok: false, + decision: "skip", + should_run: false, + effective_action: "quota_skip", + normal_delivery_allowed: false, + recovery_delivery_allowed: false, + self_repair_allowed: false, + state: "action_selection_required", + reason: + "the current projected default differs from the explicit Todo retained by this Turn", + recommended_action: + "rerun quota should-run with the same --turn-instance-id and an explicit eligible --todo-id", + }, + execution_obligation_patch: { + must_attempt_work: false, + delivery_allowed: false, + reason: + "rerun quota should-run with the same --turn-instance-id and an explicit eligible --todo-id", + }, + clear_fields: ["selected_todo", "todo_id", "agent_lane_next_action"], + }, }); }); diff --git a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts index 9949c0eda2..5aa87ff20a 100644 --- a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts +++ b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts @@ -269,6 +269,15 @@ test("commit owns the repairable run artifacts and exact-effect replay", async ( const written = await evaluateQuotaMonitorPollCommit(params); assert.equal(written.status, "written"); assert.equal(written.payload.appended, true); + assert.deepEqual(written.payload.turn_continuation, { + schema_version: "quota_turn_continuation_v0", + settlement_binding_matches_observation: null, + current_turn_settled: false, + same_turn_independent_settlement_allowed: false, + next_turn_required: false, + next_action: "obtain a typed settlement binding before claiming this Turn settled", + reason: "the committed monitor-poll has no exact Todo settlement binding", + }); const jsonPath = String(written.payload.json_path); const markdownPath = String(written.payload.markdown_path); const indexPath = String(written.payload.index_path); From aaf570c1f39bddff704587c531f84a45303b33f3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 07:17:57 +0800 Subject: [PATCH 3/6] fix(quota): settle validation-bearing replans Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota/settlement_validation.py | 49 ++++ loopx/control_plane/quota/slot_accounting.py | 6 + .../completion_validation_accountability.py | 8 + loopx/state_refresh.py | 21 +- .../test_replan_completion_validation.py | 217 ++++++++++++++++++ 5 files changed, 292 insertions(+), 9 deletions(-) create mode 100644 tests/control_plane/test_replan_completion_validation.py diff --git a/loopx/control_plane/quota/settlement_validation.py b/loopx/control_plane/quota/settlement_validation.py index 74569198b5..9b9c9bfa37 100644 --- a/loopx/control_plane/quota/settlement_validation.py +++ b/loopx/control_plane/quota/settlement_validation.py @@ -10,6 +10,49 @@ from ..todos.contract import normalize_todo_id +def _is_qualified_semantic_replan_writeback( + delivery_run: dict[str, Any] | None, + *, + todo_id: str, +) -> bool: + """Recognize the durable proof that an open Todo changed its path. + + The refresh path writes this acknowledgement only after typed replan + qualification. Spend may trust it only through the exact Turn settlement + readback; callers must not pass an unrelated latest run. + """ + + if not isinstance(delivery_run, dict): + return False + settlement_identity = delivery_run.get("settlement_identity") + if ( + normalize_todo_id(delivery_run.get("todo_id")) != todo_id + or not isinstance(settlement_identity, dict) + or normalize_todo_id(settlement_identity.get("todo_id")) != todo_id + ): + return False + ack = delivery_run.get("autonomous_replan_ack") + if ( + not isinstance(ack, dict) + or ack.get("schema_version") != "autonomous_replan_ack_v0" + or ack.get("recorded") is not True + ): + return False + semantic_delta = ack.get("semantic_delta") + if ( + isinstance(semantic_delta, dict) + and semantic_delta.get("schema_version") == "replan_semantic_delta_v0" + and semantic_delta.get("accepted") is True + ): + return True + delta_contract = ack.get("delta_contract") + return bool( + isinstance(delta_contract, dict) + and delta_contract.get("schema_version") == "repair_delta_contract_v0" + and delta_contract.get("delta_present") is True + ) + + def completion_validation_spend_error( status_payload: dict[str, Any], *, @@ -17,12 +60,18 @@ def completion_validation_spend_error( todo_id: str | None, agent_id: str | None, selected_todo: dict[str, Any] | None, + delivery_run: dict[str, Any] | None = None, ) -> str | None: """Explain why one requested settlement still lacks validation authority.""" normalized_todo_id = normalize_todo_id(todo_id) if not normalized_todo_id: return None + if _is_qualified_semantic_replan_writeback( + delivery_run, + todo_id=normalized_todo_id, + ): + return None summaries: list[dict[str, Any]] = [] if ( selected_todo diff --git a/loopx/control_plane/quota/slot_accounting.py b/loopx/control_plane/quota/slot_accounting.py index 4166d99b4a..8d0d6f265f 100644 --- a/loopx/control_plane/quota/slot_accounting.py +++ b/loopx/control_plane/quota/slot_accounting.py @@ -70,6 +70,7 @@ def _todo_binding_error( requested_replan_obligation_id: str | None, agent_id: str | None, settlement_identity: SettlementIdentity | None = None, + delivery_run: dict[str, Any] | None = None, ) -> str | None: selected = ( before.get("selected_todo") @@ -93,6 +94,7 @@ def _todo_binding_error( todo_id=requested_todo_id, agent_id=agent_id, selected_todo=selected, + delivery_run=delivery_run, ) selected_todo_id = normalize_todo_id(selected.get("todo_id")) if requested_todo_id and selected_todo_id and requested_todo_id != selected_todo_id: @@ -117,6 +119,7 @@ def _todo_binding_error( todo_id=requested_todo_id, agent_id=agent_id, selected_todo=selected, + delivery_run=None, ) @@ -644,6 +647,9 @@ def build_quota_slot_preview_for_decision( requested_replan_obligation_id=normalized_replan_obligation_id, agent_id=safe_requested_agent_id, settlement_identity=settlement_identity, + delivery_run=( + delivery_completion_run if settlement_identity is not None else None + ), ) if binding_error: return { diff --git a/loopx/control_plane/todos/completion_validation_accountability.py b/loopx/control_plane/todos/completion_validation_accountability.py index d826a3c1a8..d08fb5ade9 100644 --- a/loopx/control_plane/todos/completion_validation_accountability.py +++ b/loopx/control_plane/todos/completion_validation_accountability.py @@ -18,9 +18,17 @@ def require_accountable_completion_validation( todo_fields: dict[str, Any] | None = None, delivery_boundary: str | None = None, delivery_outcome: str | None = None, + semantic_replan_recorded: bool = False, ) -> None: """Reject accountable evidence while its exact validation Todo is open.""" + # A qualified semantic replan updates the path for an open Todo; it does + # not claim that Todo completed. Its write-time qualification already + # rejects zero-effect replans, so terminal completion validation remains + # attached to the later Todo completion instead of circularly fencing the + # path change needed to reach it. + if semantic_replan_recorded: + return if ( delivery_boundary == DELIVERY_BOUNDARY_IN_FLIGHT and delivery_outcome == DeliveryOutcome.OUTCOME_PROGRESS.value diff --git a/loopx/state_refresh.py b/loopx/state_refresh.py index 6bc332cb51..def232489d 100644 --- a/loopx/state_refresh.py +++ b/loopx/state_refresh.py @@ -951,15 +951,6 @@ def refresh_state_run( runtime_root, safe_goal_id, resolved_state_file, require_display=bool(next_action) ) expected_write_state_text = state_text - if normalized_delivery_outcome in ACCOUNTABLE_DELIVERY_OUTCOMES: - require_accountable_completion_validation( - state_text, - todo_fields=todo_fields, - todo_id=(settlement_identity.todo_id if settlement_identity else None), - agent_id=normalized_agent_id or None, - delivery_boundary=normalized_delivery_boundary, - delivery_outcome=normalized_delivery_outcome, - ) normalized_next_action = normalize_next_action_text(next_action) if next_action else None registered_agents = registered_agents_for_goal(registry_goal) known_agents = {agent for agent in registered_agents if agent} @@ -1116,6 +1107,18 @@ def refresh_state_run( effective_autonomous_replan_recorded = ( replan_qualification.autonomous_replan_recorded ) + if normalized_delivery_outcome in ACCOUNTABLE_DELIVERY_OUTCOMES: + require_accountable_completion_validation( + state_text, + todo_fields=todo_fields, + todo_id=(settlement_identity.todo_id if settlement_identity else None), + agent_id=normalized_agent_id or None, + delivery_boundary=normalized_delivery_boundary, + delivery_outcome=normalized_delivery_outcome, + semantic_replan_recorded=( + effective_autonomous_replan_recorded + ), + ) vision_checkpoint = build_vision_checkpoint( agent_id=normalized_agent_id or None, agent_vision=agent_vision, diff --git a/tests/control_plane/test_replan_completion_validation.py b/tests/control_plane/test_replan_completion_validation.py new file mode 100644 index 0000000000..d63cc9b184 --- /dev/null +++ b/tests/control_plane/test_replan_completion_validation.py @@ -0,0 +1,217 @@ +"""Semantic replans must not be mistaken for Todo completion. + +Controller-declared completion validation protects the terminal Todo +transition. A qualified, evidence-linked path replan keeps that Todo open and +must still be able to settle its exact Turn; otherwise the path change needed +to reach validation is circularly blocked. +""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from loopx.control_plane.quota.settlement_validation import ( + completion_validation_spend_error, +) +from loopx.control_plane.todos.completion_validation_accountability import ( + require_accountable_completion_validation, +) +from loopx.control_plane.turn_driver.delivery_continuity import ( + DELIVERY_BOUNDARY_SEMANTIC_CLOSEOUT, +) +from test_quota_settlement_cli import ( + AGENT_ID, + GOAL_ID, + TODO_ID, + _configure_completion_validation_todo, + _run_cli, + _spend_run_count, + _write_fixture, +) + + +def _validation_todo() -> dict[str, object]: + return { + "todo_id": TODO_ID, + "status": "open", + "claimed_by": AGENT_ID, + "completion_validation_required": True, + } + + +def _vision_path(path: Path) -> Path: + packet = { + "schema_version": "goal_vision_replan_contract_v0", + "agent_id": AGENT_ID, + "state": "active", + "vision_patch": { + "vision_summary": "Use the evidence-backed successor path.", + "acceptance_summary": "Validate the successor before Todo completion.", + "replan_trigger_summary": "The prior path cannot reach validation.", + }, + "todo_delta": [], + "path_delta": { + "schema_version": "goal_path_delta_v0", + "outcome": "replan", + "prior_assumption": "The selected path could reach terminal validation.", + "observed_reality": "A control-plane fence blocks that path before validation.", + "retained": ["Controller-declared terminal completion validation."], + "changed": ["Use the independently validated successor path."], + "stopped": ["Retrying the fenced writeback without a path change."], + "evidence_refs": ["evidence:replan-validation-fence"], + }, + } + path.write_text(json.dumps(packet), encoding="utf-8") + return path + + +def test_qualified_path_replan_settles_without_completing_validation_todo( + tmp_path: Path, +) -> None: + project, runtime, registry = _write_fixture(tmp_path) + state_path = _configure_completion_validation_todo(project) + initial_state = state_path.read_text(encoding="utf-8") + turn_id = "turn-validation-bearing-replan" + + guard_rc, guard = _run_cli( + registry, + runtime, + "quota", + "should-run", + "--codex-app", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--todo-id", + TODO_ID, + "--turn-instance-id", + turn_id, + "--scan-path", + str(project), + ) + assert guard_rc == 0, guard + assert guard["heartbeat_receipt"]["settlement_identity"]["todo_id"] == TODO_ID + + refresh_rc, refresh = _run_cli( + registry, + runtime, + "refresh-state", + "--goal-id", + GOAL_ID, + "--classification", + "evidence_linked_path_replan", + "--delivery-batch-scale", + "implementation", + "--delivery-outcome", + "outcome_progress", + "--delivery-boundary", + DELIVERY_BOUNDARY_SEMANTIC_CLOSEOUT, + "--agent-id", + AGENT_ID, + "--todo-id", + TODO_ID, + "--turn-instance-id", + turn_id, + "--autonomous-replan-recorded", + "--repair-delta-kind", + "goal_vision_patch", + "--agent-vision-json", + str(_vision_path(tmp_path / "vision.json")), + "--no-global-sync", + "--suppress-external-sinks", + ) + assert refresh_rc == 0, refresh + assert refresh["settlement_result"]["ok"] is True + persisted = json.loads( + Path(refresh["json_path"]).read_text(encoding="utf-8") + ) + assert persisted["autonomous_replan_ack"]["recorded"] is True + assert persisted["autonomous_replan_ack"]["delta_contract"][ + "delta_present" + ] is True + assert persisted["agent_vision"]["path_delta"]["outcome"] == "replan" + assert state_path.read_text(encoding="utf-8") == initial_state + + spend_args = ( + "quota", + "spend-slot", + "--goal-id", + GOAL_ID, + "--slots", + "1", + "--source", + "heartbeat", + "--execute", + "--agent-id", + AGENT_ID, + "--todo-id", + TODO_ID, + "--turn-instance-id", + turn_id, + "--scan-path", + str(project), + ) + spend_rc, spend = _run_cli(registry, runtime, *spend_args) + assert spend_rc == 0, spend + assert spend["appended"] is True + assert spend["settlement_result"]["ok"] is True + replay_rc, replay = _run_cli(registry, runtime, *spend_args) + assert replay_rc == 0, replay + assert replay["idempotent_replay"] is True + assert _spend_run_count(runtime) == 1 + + +def test_unqualified_primary_outcome_keeps_completion_and_spend_fences() -> None: + todo_fields = {"agent_todos": {"items": [_validation_todo()]}} + with pytest.raises(ValueError, match="completion validation"): + require_accountable_completion_validation( + "", + todo_id=TODO_ID, + agent_id=AGENT_ID, + todo_fields=todo_fields, + delivery_boundary=DELIVERY_BOUNDARY_SEMANTIC_CLOSEOUT, + delivery_outcome="primary_goal_outcome", + semantic_replan_recorded=False, + ) + + zero_effect_run = { + "todo_id": TODO_ID, + "settlement_identity": {"todo_id": TODO_ID}, + "autonomous_replan_ack": { + "schema_version": "autonomous_replan_ack_v0", + "recorded": False, + "delta_contract": { + "schema_version": "repair_delta_contract_v0", + "delta_present": False, + }, + }, + } + wrong_identity_run = { + "todo_id": "todo_other_settlement", + "settlement_identity": {"todo_id": "todo_other_settlement"}, + "autonomous_replan_ack": { + "schema_version": "autonomous_replan_ack_v0", + "recorded": True, + "delta_contract": { + "schema_version": "repair_delta_contract_v0", + "delta_present": True, + }, + }, + } + expected = ( + "quota spend is blocked until controller-declared completion " + f"validation durably completes todo {TODO_ID}" + ) + for delivery_run in (zero_effect_run, wrong_identity_run): + assert completion_validation_spend_error( + {"attention_queue": {"items": []}}, + goal_id=GOAL_ID, + todo_id=TODO_ID, + agent_id=AGENT_ID, + selected_todo=_validation_todo(), + delivery_run=delivery_run, + ) == expected From 903047f04c622cec6f3fa0abc1f6fd4f5c5ffd98 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 07:40:36 +0800 Subject: [PATCH 4/6] fix(quota): preserve settled turn precedence Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/effect_program.py | 3 + .../control_plane/effect_runtime_handlers.ts | 2 + loopx/control_plane/quota/live_decision.py | 26 ++++--- loopx/control_plane/quota/settlement_phase.ts | 5 +- .../quota/settlement_precedence.py | 2 +- .../quota/settlement_readback.ts | 7 +- .../test_effect_turn_live_quota_decision.py | 78 +++++++++++++++++++ .../test_quota_settlement_cli.py | 24 ++++++ tests/control_plane_ts/effect_program.test.ts | 10 +++ 9 files changed, 144 insertions(+), 13 deletions(-) diff --git a/loopx/control_plane/effect_program.py b/loopx/control_plane/effect_program.py index 8c76270ca9..d00e8c73b8 100644 --- a/loopx/control_plane/effect_program.py +++ b/loopx/control_plane/effect_program.py @@ -59,6 +59,7 @@ def receipt_bound_monitor_phase( def receipt_bound_replay_phase( *, binding_kind: SettlementBindingKind | str | None = None, + writeback_completes_binding: bool = False, completion_receipt_present: bool, durable_writeback_present: bool, quota_spend_present: bool, @@ -70,6 +71,8 @@ def receipt_bound_replay_phase( } if binding_kind is not None: params["binding_kind"] = str(binding_kind) + if writeback_completes_binding: + params["writeback_completes_binding"] = True result = effect_runtime_result( "settlement.receipt_bound_replay_phase", params, diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 95fcfe1647..fc7c9e649b 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -698,6 +698,8 @@ export function createEffectRuntimeHandlers( "binding_kind", "binding_kind has an unsupported settlement binding kind", ), + writeback_completes_binding: + params.writeback_completes_binding === true, completion_receipt_present: params.completion_receipt_present === true, durable_writeback_present: params.durable_writeback_present === true, quota_spend_present: params.quota_spend_present === true, diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index 05dd3c16fa..94800b8db3 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -16,6 +16,7 @@ from .settlement import ( read_heartbeat_settlement, ) +from .effect_program import ReceiptBoundReplayPhase from ..work_items.interaction_contract import ( build_interaction_contract, build_protocol_action_packet, @@ -613,16 +614,21 @@ def build_live_quota_should_run_decision( interaction = payload.get("interaction_contract") if isinstance(interaction, dict): interaction.update(projections) - apply_unsettled_host_turn_recovery_if_required( - payload, - registry_path=registry_path, - runtime_root=runtime_root, - goal_id=goal_id, - agent_id=agent_id, - current_turn_instance_id=turn_instance_id, - available_capabilities=available_capabilities, - scheduler_execution_context=resolved_context, - ) + # A settled receipt owns this host Turn until it ends. Looking for an older + # unsettled Turn here can overwrite the settled-skip route with a recovery + # obligation and then select a successor against the immutable receipt + # identity. Leave prior-Turn recovery to the next fresh Turn instead. + if receipt_bound_replay_phase is not ReceiptBoundReplayPhase.SETTLED: + apply_unsettled_host_turn_recovery_if_required( + payload, + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + agent_id=agent_id, + current_turn_instance_id=turn_instance_id, + available_capabilities=available_capabilities, + scheduler_execution_context=resolved_context, + ) if hook_dispatch["failures"]: payload["capability_hook_dispatch"] = { key: value for key, value in hook_dispatch.items() if key != "projections" diff --git a/loopx/control_plane/quota/settlement_phase.ts b/loopx/control_plane/quota/settlement_phase.ts index 36a45b102f..fe4c834280 100644 --- a/loopx/control_plane/quota/settlement_phase.ts +++ b/loopx/control_plane/quota/settlement_phase.ts @@ -48,6 +48,8 @@ export type ReceiptBoundReplayPhase = export interface ReceiptBoundReplaySettlementState { binding_kind?: "todo" | "autonomous_replan" | "unbound"; + /** A validated writeback discharges the binding without Todo completion. */ + writeback_completes_binding?: boolean; completion_receipt_present: boolean; durable_writeback_present: boolean; quota_spend_present: boolean; @@ -56,7 +58,8 @@ export interface ReceiptBoundReplaySettlementState { export function receiptBoundReplayPhase( state: ReceiptBoundReplaySettlementState, ): ReceiptBoundReplayPhase { - const bindingComplete = state.binding_kind === "autonomous_replan" + const bindingComplete = state.binding_kind === "autonomous_replan" || + state.writeback_completes_binding === true ? state.durable_writeback_present : state.completion_receipt_present; if (!bindingComplete) return "open"; diff --git a/loopx/control_plane/quota/settlement_precedence.py b/loopx/control_plane/quota/settlement_precedence.py index f1cf97b179..a9894339e4 100644 --- a/loopx/control_plane/quota/settlement_precedence.py +++ b/loopx/control_plane/quota/settlement_precedence.py @@ -7,7 +7,7 @@ HEARTBEAT_SETTLED_REPLAY_REASON = ( - "the receipt-bound Todo completion and required settlement receipts " + "the receipt-bound work binding and required settlement receipts " "are complete for this heartbeat turn; defer successor selection to a new turn" ) diff --git a/loopx/control_plane/quota/settlement_readback.ts b/loopx/control_plane/quota/settlement_readback.ts index 69eb4de85e..1f135abdff 100644 --- a/loopx/control_plane/quota/settlement_readback.ts +++ b/loopx/control_plane/quota/settlement_readback.ts @@ -754,6 +754,10 @@ export async function readQuotaSettlement(value: unknown): Promise { const workspaceCausality: DeliveryWorkspaceCausality | null = normalizeDeliveryWorkspaceCausality(nestedCausality, identity.todo_id) ?? normalizeDeliveryWorkspaceCausality(flatCausality, identity.todo_id); + const semanticReplanGuard = projectSemanticReplanGuard(receiptDetails); + const todoBoundReplan = identity.binding_kind === "todo" && + semanticReplanGuard.scope === "turn_guard" && + semanticReplanGuard.selected_obligation_id !== null; const recovery = request.refresh_retry === null ? null : refreshRecovery( request.refresh_retry, writebackRun, writeback.failure === null, @@ -775,7 +779,7 @@ export async function readQuotaSettlement(value: unknown): Promise { terminal_closeout: bundle(terminalCloseout), terminal_settlement: bundle(terminalSettlement), workspace_causality: workspaceCausality, - semantic_replan_guard: projectSemanticReplanGuard(receiptDetails), + semantic_replan_guard: semanticReplanGuard, writeback_run: writebackRun, refresh_recovery: recovery, external_delivery: request.refresh_retry === null ? null : refreshExternalDelivery( @@ -795,6 +799,7 @@ export async function readQuotaSettlement(value: unknown): Promise { }), replay_phase: receiptBoundReplayPhase({ binding_kind: identity.binding_kind, + writeback_completes_binding: todoBoundReplan, completion_receipt_present: completionEvent !== null, durable_writeback_present: writeback.failure === null, quota_spend_present: spend.failure === null, diff --git a/tests/control_plane/test_effect_turn_live_quota_decision.py b/tests/control_plane/test_effect_turn_live_quota_decision.py index d83c7a9853..b98e2f5470 100644 --- a/tests/control_plane/test_effect_turn_live_quota_decision.py +++ b/tests/control_plane/test_effect_turn_live_quota_decision.py @@ -2,9 +2,12 @@ import json from pathlib import Path +from types import SimpleNamespace import pytest +from loopx.control_plane.quota import live_decision +from loopx.control_plane.quota.effect_program import ReceiptBoundReplayPhase from loopx.control_plane.effect_program import ( interpret_quota_should_run_packet, ) @@ -278,6 +281,81 @@ def test_managed_turn_projects_prior_unsettled_heartbeat_recovery( assert contract["cli_channel"]["spend_after_validation"] is False +def test_settled_turn_defers_prior_turn_recovery_to_a_fresh_turn( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + runtime_root = tmp_path / "runtime" + registry_path = tmp_path / "registry.json" + registry_path.write_text(json.dumps({"goals": []}), encoding="utf-8") + + monkeypatch.setattr( + live_decision, + "read_heartbeat_settlement", + lambda *_args, **_kwargs: SimpleNamespace( + monitor_phase=None, + replay_phase=ReceiptBoundReplayPhase.SETTLED, + ), + ) + + def fail_if_recovery_runs(*_args: object, **_kwargs: object) -> bool: + pytest.fail("a settled host Turn must not inspect prior-Turn recovery") + + monkeypatch.setattr( + live_decision, + "apply_unsettled_host_turn_recovery_if_required", + fail_if_recovery_runs, + ) + todo_text = "[P1] Keep advancing the selected task." + status = quota_status_payload( + goal_id=GOAL_ID, + status="active", + agent_todo_items=[ + { + "todo_id": "todo_ordinary_work", + "index": 1, + "text": todo_text, + "role": "agent", + "status": "open", + "priority": "P1", + "task_class": "advancement_task", + } + ], + recommended_action=todo_text, + next_action=todo_text, + coordination={ + "registered_agents": ["codex-fixture"], + "agent_model": "peer_v1", + }, + claim_scope_agent_id="codex-fixture", + ) + + packet = build_live_quota_should_run_decision( + status, + goal_id=GOAL_ID, + agent_id="codex-fixture", + available_capabilities=["shell"], + include_scheduler_detail=False, + codex_app_current_rrule=None, + registry_path=registry_path, + runtime_root=runtime_root, + route_source="loopx_turn_plan", + receipt_bound_todo_id="todo_ordinary_work", + turn_instance_id="managed-settled-turn", + scheduler_execution_context={ + "host_surface": "generic_cli", + "scheduler_owner": "agent_cli_loop", + "execution_mode": "interactive", + }, + ) + + assert packet["decision"] == "skip" + assert packet["effective_action"] == "heartbeat_settled_skip" + assert packet["should_run"] is False + assert packet.get("selected_todo") is None + assert packet.get("unsettled_host_turn_recovery") is None + + def test_managed_turn_accepts_exact_material_monitor_poll_closeout( tmp_path: Path, ) -> None: diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index 5d34344027..c8f9ba7107 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -4345,6 +4345,30 @@ def test_todoless_blocked_replan_settles_read_only_external_evidence_without_wor ] == ["validation", "durable_writeback", "quota_spend"] assert _spend_run_count(runtime) == 1 + replay_rc, replay = _run_cli( + registry_path, + runtime, + "quota", + "should-run", + "--codex-app", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--turn-instance-id", + turn_instance_id, + "--todo-id", + SELECTED_REPLAN_TODO_ID, + "--scan-path", + str(project), + ) + + assert replay_rc == 0, replay + assert replay["effective_action"] == "heartbeat_settled_skip" + assert replay["should_run"] is False + assert replay.get("selected_todo") is None + assert replay.get("unsettled_host_turn_recovery") is None + def test_unbound_visible_goal_todoless_replan_reenters_through_guided_turn( tmp_path: Path, diff --git a/tests/control_plane_ts/effect_program.test.ts b/tests/control_plane_ts/effect_program.test.ts index 13cafa16cc..172195be15 100644 --- a/tests/control_plane_ts/effect_program.test.ts +++ b/tests/control_plane_ts/effect_program.test.ts @@ -233,6 +233,16 @@ test("receipt-bound replay settlement follows its binding and full chain", () => }), "settled", ); + assert.equal( + receiptBoundReplayPhase({ + binding_kind: "todo", + writeback_completes_binding: true, + completion_receipt_present: false, + durable_writeback_present: true, + quota_spend_present: true, + }), + "settled", + ); assert.equal( receiptBoundReplayPhase({ completion_receipt_present: true, From f415da1b56aa9feecd76af705e951c2421d5a470 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 08:20:44 +0800 Subject: [PATCH 5/6] fix(quota): preserve settled replan replay Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/quota/settlement_cli.py | 65 ++++++++++++--------- 1 file changed, 38 insertions(+), 27 deletions(-) diff --git a/loopx/control_plane/quota/settlement_cli.py b/loopx/control_plane/quota/settlement_cli.py index c7722218b8..437e1dc907 100644 --- a/loopx/control_plane/quota/settlement_cli.py +++ b/loopx/control_plane/quota/settlement_cli.py @@ -67,35 +67,46 @@ def reconcile_existing_heartbeat_receipt( replan_obligation_id=rollout_replan_obligation_id, ).effect_id if (existing_todo_id or existing_replan_obligation_id) and existing_effect_id: - requested_todo_id = normalize_todo_id(args.todo_id) - if requested_todo_id and requested_todo_id != existing_todo_id: - raise HeartbeatReceiptIdentityConflictError( - "heartbeat receipt settlement identity conflicts with the " - "current selected Todo: explicitly requested Todo differs" - ) - requested_replan_obligation_id = normalize_todo_replan_obligation_id( - getattr(args, "replan_obligation_id", None) + # A fully settled Turn owns its immutable receipt regardless of a + # later explicit CLI choice. The live decision has already read + # the exact settlement and projected ``heartbeat_settled_skip``; + # rechecking that new choice here would turn a harmless same-Turn + # replay into an identity conflict and reopen successor routing. + settled_replay = ( + payload.get("effective_action") + == EffectiveAction.HEARTBEAT_SETTLED_SKIP.value ) - if ( - requested_replan_obligation_id - and requested_replan_obligation_id != existing_replan_obligation_id - ): - raise HeartbeatReceiptIdentityConflictError( - "heartbeat receipt settlement identity conflicts with the " - "explicitly requested autonomous replan obligation" - ) - if ( - existing_todo_id != rollout_todo_id - or existing_replan_obligation_id != rollout_replan_obligation_id - or existing_effect_id != expected_effect_id - ): - raise HeartbeatReceiptIdentityConflictError( - "heartbeat receipt settlement identity conflicts with the " - "current settlement binding: receipt=" - f"{existing_todo_id or existing_replan_obligation_id}, " - "current=" - f"{rollout_todo_id or rollout_replan_obligation_id}" + if not settled_replay: + requested_todo_id = normalize_todo_id(args.todo_id) + if requested_todo_id and requested_todo_id != existing_todo_id: + raise HeartbeatReceiptIdentityConflictError( + "heartbeat receipt settlement identity conflicts with the " + "current selected Todo: explicitly requested Todo differs" + ) + requested_replan_obligation_id = normalize_todo_replan_obligation_id( + getattr(args, "replan_obligation_id", None) ) + if ( + requested_replan_obligation_id + and requested_replan_obligation_id + != existing_replan_obligation_id + ): + raise HeartbeatReceiptIdentityConflictError( + "heartbeat receipt settlement identity conflicts with the " + "explicitly requested autonomous replan obligation" + ) + if ( + existing_todo_id != rollout_todo_id + or existing_replan_obligation_id != rollout_replan_obligation_id + or existing_effect_id != expected_effect_id + ): + raise HeartbeatReceiptIdentityConflictError( + "heartbeat receipt settlement identity conflicts with the " + "current settlement binding: receipt=" + f"{existing_todo_id or existing_replan_obligation_id}, " + "current=" + f"{rollout_todo_id or rollout_replan_obligation_id}" + ) else: rollout_details = quota_rollout_details( payload, From 5c371b77c54098da15bea056c2b31582719aa562 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 09:58:20 +0800 Subject: [PATCH 6/6] fix(quota): reject conflicting settled turn selection Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/quota.py | 47 ++++++++++++++++++--- loopx/control_plane/quota/settlement_cli.py | 4 +- 2 files changed, 43 insertions(+), 8 deletions(-) diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 0abef659e6..caa8594e1b 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -22,6 +22,7 @@ ) from ..control_plane.quota.effect_program import SettlementIdentity from ..control_plane.quota.error_codes import ( + HeartbeatReceiptIdentityConflictError, QuotaCommandValidationError, QuotaActionSelectionConflictError, QuotaActionSelectionConflictKind, @@ -176,9 +177,9 @@ def _heartbeat_quota_action_selection_bindings( runtime_root: Path, args: argparse.Namespace, heartbeat_turn_id: str | None, -) -> tuple[dict[str, object] | None, str | None, str | None, str | None]: +) -> tuple[dict[str, object] | None, str | None, str | None, str | None, bool]: if not heartbeat_turn_id: - return None, None, None, None + return None, None, None, None, False existing = find_heartbeat_receipt( runtime_root, goal_id=args.goal_id, @@ -186,13 +187,19 @@ def _heartbeat_quota_action_selection_bindings( turn_instance_id=heartbeat_turn_id, ) if not existing: - return None, None, None, None + return None, None, None, None, False todo_id, replan_obligation_id = _heartbeat_receipt_settlement_bindings(existing) + details_value = existing.get("details") + details: Mapping[str, object] = ( + details_value if isinstance(details_value, Mapping) else {} + ) return ( existing, todo_id, replan_obligation_id, heartbeat_receipt_pending_action_todo_id(existing), + str(details.get("settlement_receipt_revision") or "") + == "identity_upgrade", ) @@ -202,10 +209,17 @@ def _apply_requested_quota_action_selection_preflight( requested_todo_id: str | None, receipt_bound_todo_id: str | None, receipt_bound_replan_obligation_id: str | None, + receipt_pending_action_todo_id: str | None, + receipt_identity_upgraded: bool, ) -> bool: - if not requested_todo_id or ( - receipt_bound_todo_id or receipt_bound_replan_obligation_id - ): + if not requested_todo_id: + return False + if receipt_bound_todo_id: + if requested_todo_id != receipt_bound_todo_id: + raise HeartbeatReceiptIdentityConflictError( + "heartbeat receipt settlement identity conflicts with the " + "current selected Todo: explicitly requested Todo differs" + ) return False selected_todo = payload.get("selected_todo") selected_todo_id = ( @@ -230,6 +244,20 @@ def _apply_requested_quota_action_selection_preflight( if isinstance(qualification_selected, Mapping) else None ) + if receipt_bound_replan_obligation_id: + if not receipt_identity_upgraded: + # A Turn that started directly in autonomous replan has no Todo + # selection authority to replace. A later same-Turn --todo-id is + # therefore a harmless settled replay, not successor delivery. + return False + if requested_todo_id == receipt_pending_action_todo_id: + return False + raise QuotaActionSelectionConflictError( + QuotaActionSelectionConflictKind.CONFLICT, + requested_todo_id=requested_todo_id, + selected_todo_id=receipt_pending_action_todo_id, + qualification_state="retained_selection", + ) selection_binding = ( selected_todo.get("selection_binding") if isinstance(selected_todo, Mapping) @@ -349,11 +377,15 @@ def _reconcile_requested_quota_action_selection( context: QuotaCommandContext, receipt_bound_todo_id: str | None, receipt_bound_replan_obligation_id: str | None, + receipt_pending_action_todo_id: str | None, + receipt_identity_upgraded: bool, ) -> bool: rejected = _apply_requested_quota_action_selection_preflight( payload, requested_todo_id=_requested_quota_action_todo_id(args), receipt_bound_todo_id=receipt_bound_todo_id, receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, + receipt_pending_action_todo_id=receipt_pending_action_todo_id, + receipt_identity_upgraded=receipt_identity_upgraded, ) if rejected: apply_action_selection_recovery( @@ -687,6 +719,7 @@ def handle_quota_command( receipt_bound_todo_id, receipt_bound_replan_obligation_id, receipt_pending_action_todo_id, + receipt_identity_upgraded, ) = _heartbeat_quota_action_selection_bindings( runtime_root=runtime_root, args=args, @@ -734,6 +767,8 @@ def handle_quota_command( payload, args, registry_path=registry_path, context=context, receipt_bound_todo_id=receipt_bound_todo_id, receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, + receipt_pending_action_todo_id=receipt_pending_action_todo_id, + receipt_identity_upgraded=receipt_identity_upgraded, ) if action_selection_preflight_failed: ( diff --git a/loopx/control_plane/quota/settlement_cli.py b/loopx/control_plane/quota/settlement_cli.py index 437e1dc907..a616d6eede 100644 --- a/loopx/control_plane/quota/settlement_cli.py +++ b/loopx/control_plane/quota/settlement_cli.py @@ -70,8 +70,8 @@ def reconcile_existing_heartbeat_receipt( # A fully settled Turn owns its immutable receipt regardless of a # later explicit CLI choice. The live decision has already read # the exact settlement and projected ``heartbeat_settled_skip``; - # rechecking that new choice here would turn a harmless same-Turn - # replay into an identity conflict and reopen successor routing. + # action-selection preflight distinguishes a qualified current + # Todo from a conflicting explicit choice before this replay path. settled_replay = ( payload.get("effective_action") == EffectiveAction.HEARTBEAT_SETTLED_SKIP.value