From 7e49a13c14f2be9bcc5fe6971cbda42e98556cb7 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 21 Sep 2026 00:30:14 +0800 Subject: [PATCH] refactor(cli): extract quota action selection owner Signed-off-by: duanjialing.777 --- loopx/cli_commands/quota.py | 445 +++--------------- loopx/cli_commands/quota_action_selection.py | 415 ++++++++++++++++ loopx/semantics/vocabulary_v0.json | 2 +- .../test_quota_action_selection_conflict.py | 4 +- 4 files changed, 474 insertions(+), 392 deletions(-) create mode 100644 loopx/cli_commands/quota_action_selection.py diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index caa8594e1..8e4d8f157 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -1,5 +1,4 @@ from __future__ import annotations -from ..control_plane.quota.effective_action import EffectiveAction import argparse from collections.abc import Callable, Mapping @@ -20,22 +19,13 @@ compact_quota_monitor_poll_cli_payload, compact_quota_should_run_cli_payload, ) +from ..control_plane.quota.effective_action import EffectiveAction from ..control_plane.quota.effect_program import SettlementIdentity -from ..control_plane.quota.error_codes import ( - HeartbeatReceiptIdentityConflictError, - QuotaCommandValidationError, - QuotaActionSelectionConflictError, - QuotaActionSelectionConflictKind, -) +from ..control_plane.quota.error_codes import QuotaCommandValidationError from ..control_plane.quota.heartbeat_receipt import ( - 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 @@ -48,12 +38,6 @@ render_existing_heartbeat_receipt_payload, ) from ..control_plane.quota.turn_envelope import build_turn_envelope -from ..control_plane.scheduler.execution_context import ( - GUIDED_START_TURN_RUNTIME_PROFILES, - render_scheduler_execution_args, -) -from ..control_plane.todos.contract import normalize_todo_id -from ..control_plane.work_items.action_selection_contract import apply_action_selection_recovery from ..presentation.renderers.quota_event_markdown import ( render_quota_monitor_poll_markdown, render_quota_slot_preview_markdown, @@ -79,6 +63,13 @@ build_lark_operator_inbox_urgency_projector, dispatch_goal_lark_turn_start_hooks, ) +from .quota_action_selection import ( + RequestedQuotaActionSelection, + attach_uncommitted_action_selection_receipt, + commit_requested_action_selection, + load_requested_quota_action_selection, + reconcile_requested_quota_action_selection, +) from .quota_context import ( QuotaCommandContext, prepare_quota_command_context, @@ -106,13 +97,6 @@ None, ] RolloutEventAppender = Callable[..., dict[str, object]] -def _heartbeat_receipt_settlement_bindings( - event: Mapping[str, object], -) -> tuple[str | None, str | None]: - return ( - heartbeat_receipt_settlement_todo_id(event), - heartbeat_receipt_settlement_replan_obligation_id(event), - ) def _effective_spend_turn_instance_id( @@ -159,315 +143,6 @@ def _quota_renderer( }.get(command, render_quota_markdown) -def _requested_quota_action_todo_id( - args: argparse.Namespace, -) -> str | None: - if not ( - bool(args.codex_app) - or bool(getattr(args, "trae_app", False)) - or args.runtime_profile - in {profile.value for profile in GUIDED_START_TURN_RUNTIME_PROFILES} - ): - return None - return normalize_todo_id(args.todo_id) - - -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, bool]: - if not heartbeat_turn_id: - return None, None, None, None, False - existing = find_heartbeat_receipt( - runtime_root, - goal_id=args.goal_id, - agent_id=args.agent_id, - turn_instance_id=heartbeat_turn_id, - ) - if not existing: - 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", - ) - - -def _apply_requested_quota_action_selection_preflight( - payload: dict[str, object], - *, - 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: - 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 = ( - normalize_todo_id(selected_todo.get("todo_id")) - if isinstance(selected_todo, Mapping) - else None - ) - qualification_value = payload.get("action_selection_qualification") - qualification: Mapping[str, object] = ( - qualification_value if isinstance(qualification_value, Mapping) else {} - ) - if selected_todo_id is None and str(qualification.get("state") or "") == ( - "qualified" - ): - # An unsettled-host-turn recovery decision carries no top-level - # `selected_todo`: its qualification names the Todo that prior Turn - # has to settle, and binding the guard to that Todo is the documented - # closeout path rather than a conflict with the projection. - qualification_selected = qualification.get("selected_todo") - selected_todo_id = ( - normalize_todo_id(qualification_selected.get("todo_id")) - 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) - else None - ) - execution_obligation_value = payload.get("execution_obligation") - execution_obligation: Mapping[str, object] = ( - execution_obligation_value - if isinstance(execution_obligation_value, Mapping) - else {} - ) - interaction_value = payload.get("interaction_contract") - interaction: Mapping[str, object] = ( - interaction_value if isinstance(interaction_value, Mapping) else {} - ) - agent_channel_value = interaction.get("agent_channel") - agent_channel: Mapping[str, object] = ( - agent_channel_value if isinstance(agent_channel_value, Mapping) else {} - ) - pending_selection_delivery_qualified = ( - selection_binding == "pending_action_selection" - and payload.get("normal_delivery_allowed") is True - ) - pending_selection_workspace_repair_qualified = ( - selection_binding == "pending_action_selection" - and payload.get("workspace_repair_allowed") is True - and payload.get("effective_action") == EffectiveAction.AGENT_WORKSPACE_REPAIR.value - and execution_obligation.get("kind") == "agent_workspace_repair" - and execution_obligation.get("must_attempt_work") is True - and agent_channel.get("must_attempt") is True - and agent_channel.get("delivery_allowed") is False - ) - exact_current_obligation_qualified = ( - selection_binding != "pending_action_selection" - and execution_obligation.get("must_attempt_work") is True - and agent_channel.get("must_attempt") is True - ) - if ( - selected_todo_id == requested_todo_id - and payload.get("ok") is True - and payload.get("should_run") is True - and ( - pending_selection_delivery_qualified - or pending_selection_workspace_repair_qualified - or exact_current_obligation_qualified - ) - ): - return False - - if not isinstance(qualification_value, Mapping): - raise QuotaActionSelectionConflictError( - QuotaActionSelectionConflictKind.UNQUALIFIED, - requested_todo_id=requested_todo_id, - selected_todo_id=selected_todo_id, - ) - qualification_state = str(qualification.get("state") or "") - if qualification_state not in {"deferred", "rejected"}: - raise QuotaActionSelectionConflictError( - QuotaActionSelectionConflictKind.CONFLICT, - requested_todo_id=requested_todo_id, - selected_todo_id=selected_todo_id, - qualification_state=qualification_state, - ) - qualification_reason = str( - qualification.get("reason") or "candidate_not_currently_eligible" - ) - deferred = qualification_state == "deferred" - auxiliary_monitor = ( - qualification_reason - == "auxiliary_monitor_not_selectable_in_advancement_lane" - ) - error_code = ( - "quota_action_selection_deferred" - if deferred - else "quota_action_selection_rejected" - ) - payload.update( - { - "ok": False, - "decision": "skip", - "should_run": False, - "effective_action": EffectiveAction.QUOTA_SKIP.value, - "state": error_code, - "waiting_on": "codex", - "status": error_code, - "error_code": error_code, - "reason": ( - "explicit action selection was deferred by the current " - f"delivery frontier: {qualification_reason}" - if deferred - else "explicit action selection is not currently eligible: " - f"{qualification_reason}" - ), - "recommended_action": ( - "handle the current delivery preemption, then rerun quota " - "should-run with the same --turn-instance-id; omit --todo-id " - "first when a refreshed action portfolio is needed" - if deferred - else "the due monitor is visible as auxiliary context, not an " - "independently selectable action in the current advancement lane; " - "choose a current advancement Todo, or rerun after the monitor " - "becomes the hard lane" - if auxiliary_monitor - else "rerun quota should-run with the same --turn-instance-id " - "without --todo-id, then choose a currently eligible Todo" - ), - } - ) - return True - - -def _reconcile_requested_quota_action_selection( - payload: dict[str, object], - args: argparse.Namespace, - *, - registry_path: Path, - 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( - payload, registry_path=str(registry_path), runtime_root=str(context.runtime_root), - goal_id=args.goal_id, agent_id=args.agent_id, - turn_instance_id=context.heartbeat_turn_id, - available_capabilities=args.available_capabilities, - scheduler_args=render_scheduler_execution_args( - scheduler_execution_context=context.scheduler_context), - ) - obligation = payload.get("execution_obligation") - if isinstance(obligation, dict): - obligation.update(must_attempt_work=False, delivery_allowed=False, - reason=payload["recommended_action"]) - return rejected - - -def _attach_uncommitted_action_selection_receipt( - payload: dict[str, object], - *, - turn_instance_id: str, -) -> None: - """Expose an accurate non-durable receipt for a rejected preflight.""" - - payload["heartbeat_receipt"] = { - "schema_version": HEARTBEAT_RECEIPT_SCHEMA_VERSION, - "turn_instance_id": turn_instance_id, - "status": "not_committed", - "stall_observation": "not_evaluated", - "reason_code": str( - payload.get("error_code") or "quota_action_selection_rejected" - ), - } - - -def _commit_requested_action_selection( - payload: Mapping[str, object], - *, - requested_todo_id: str | None, -) -> None: - """Project the exact requested selection only after receipt reconciliation.""" - - selected_todo = payload.get("selected_todo") - if ( - requested_todo_id - and isinstance(selected_todo, dict) - and normalize_todo_id(selected_todo.get("todo_id")) == requested_todo_id - ): - 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], @@ -477,9 +152,7 @@ def _record_automatic_heartbeat_stall( 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, + action_selection: RequestedQuotaActionSelection, cache_metadata: object, ) -> tuple[dict[str, object], dict[str, object], object, str]: """Commit and reproject the automatic no-spend heartbeat observation.""" @@ -542,14 +215,12 @@ def _record_automatic_heartbeat_stall( 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, + receipt_bound_todo_id=action_selection.receipt_bound_todo_id, + receipt_bound_replan_obligation_id=( + action_selection.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 + action_selection.retained_todo_id_for_decision ), turn_instance_id=turn_id, interaction_projection_hooks=interaction_projection_hooks, @@ -662,6 +333,7 @@ def handle_quota_command( heartbeat_receipt_existing_appended = False heartbeat_receipt_ready = False action_selection_preflight_failed = False + action_selection: RequestedQuotaActionSelection | None = None heartbeat_stall_observation = "not_evaluated" detail_sections: frozenset[str] = frozenset() context: QuotaCommandContext | None = None @@ -714,17 +386,12 @@ def handle_quota_command( agent_id=args.agent_id, ), ) - ( - heartbeat_receipt_existing, - receipt_bound_todo_id, - receipt_bound_replan_obligation_id, - receipt_pending_action_todo_id, - receipt_identity_upgraded, - ) = _heartbeat_quota_action_selection_bindings( + action_selection = load_requested_quota_action_selection( + args, runtime_root=runtime_root, - args=args, - heartbeat_turn_id=heartbeat_turn_id, + turn_instance_id=heartbeat_turn_id, ) + heartbeat_receipt_existing = action_selection.receipt payload = build_live_quota_should_run_decision( status_payload, goal_id=args.goal_id, @@ -744,45 +411,39 @@ def handle_quota_command( bounded_research_frontier_projector=( project_live_explore_composition_frontier ), - receipt_bound_todo_id=receipt_bound_todo_id, + receipt_bound_todo_id=action_selection.receipt_bound_todo_id, requested_action_todo_id=( - _requested_quota_action_todo_id(args) - if receipt_bound_todo_id is None - else None + action_selection.requested_todo_id_for_decision + ), + receipt_bound_replan_obligation_id=( + action_selection.receipt_bound_replan_obligation_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 + action_selection.retained_todo_id_for_decision ), turn_instance_id=heartbeat_turn_id, interaction_projection_hooks=interaction_projection_hooks, turn_start_hook_dispatch=turn_start_hook_dispatch, ) _attach_turn_start_hook_dispatch(payload, turn_start_hook_dispatch) - action_selection_preflight_failed = _reconcile_requested_quota_action_selection( - 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: - ( - heartbeat_receipt_existing, - heartbeat_receipt_existing_status, - retained_selection_appended, - ) = _retain_deferred_action_selection( + action_selection_preflight = ( + reconcile_requested_quota_action_selection( payload, args, - runtime_root=runtime_root, - heartbeat_turn_id=heartbeat_turn_id, - existing=heartbeat_receipt_existing, + registry_path=registry_path, + context=context, + selection=action_selection, + ) + ) + action_selection_preflight_failed = action_selection_preflight.rejected + if action_selection_preflight.rejected: + heartbeat_receipt_existing = action_selection_preflight.receipt + heartbeat_receipt_existing_status = ( + action_selection_preflight.receipt_status + ) + heartbeat_receipt_existing_appended = ( + action_selection_preflight.receipt_appended ) - heartbeat_receipt_existing_appended = retained_selection_appended if heartbeat_turn_id: if action_selection_preflight_failed: heartbeat_receipt_ready = True @@ -814,11 +475,7 @@ def handle_quota_command( 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, + action_selection=action_selection, cache_metadata=cache_metadata, ) heartbeat_receipt_ready = True @@ -924,7 +581,7 @@ def handle_quota_command( appended=heartbeat_receipt_existing_appended, ) else: - _attach_uncommitted_action_selection_receipt( + attach_uncommitted_action_selection_receipt( payload, turn_instance_id=heartbeat_turn_id, ) @@ -948,9 +605,13 @@ def handle_quota_command( status=heartbeat_receipt_existing_status, appended=heartbeat_receipt_existing_appended, ) - _commit_requested_action_selection( + commit_requested_action_selection( payload, - requested_todo_id=_requested_quota_action_todo_id(args), + requested_todo_id=( + action_selection.requested_todo_id + if action_selection is not None + else None + ), ) else: settlement_identity = ( @@ -1015,9 +676,13 @@ def handle_quota_command( if rollout_event.get("appended") else "replayed", ) - _commit_requested_action_selection( + commit_requested_action_selection( payload, - requested_todo_id=_requested_quota_action_todo_id(args), + requested_todo_id=( + action_selection.requested_todo_id + if action_selection is not None + else None + ), ) else: fail_heartbeat_receipt( diff --git a/loopx/cli_commands/quota_action_selection.py b/loopx/cli_commands/quota_action_selection.py new file mode 100644 index 000000000..c0db6999d --- /dev/null +++ b/loopx/cli_commands/quota_action_selection.py @@ -0,0 +1,415 @@ +from __future__ import annotations + +import argparse +from collections.abc import Mapping +from dataclasses import dataclass +from pathlib import Path + +from ..control_plane.quota.effective_action import EffectiveAction +from ..control_plane.quota.error_codes import ( + HeartbeatReceiptIdentityConflictError, + QuotaActionSelectionConflictError, + QuotaActionSelectionConflictKind, +) +from ..control_plane.quota.heartbeat_receipt import ( + HEARTBEAT_RECEIPT_SCHEMA_VERSION, + find_heartbeat_receipt, + heartbeat_receipt_pending_action_todo_id, + heartbeat_receipt_settlement_replan_obligation_id, + heartbeat_receipt_settlement_todo_id, + retain_pending_heartbeat_action_selection, +) +from ..control_plane.scheduler.execution_context import ( + GUIDED_START_TURN_RUNTIME_PROFILES, + render_scheduler_execution_args, +) +from ..control_plane.todos.contract import normalize_todo_id +from ..control_plane.work_items.action_selection_contract import ( + apply_action_selection_recovery, +) +from .quota_context import QuotaCommandContext + + +@dataclass(frozen=True, slots=True) +class RequestedQuotaActionSelection: + requested_todo_id: str | None + receipt: dict[str, object] | 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 + + @property + def requested_todo_id_for_decision(self) -> str | None: + if self.receipt_bound_todo_id is not None: + return None + return self.requested_todo_id + + @property + def retained_todo_id_for_decision(self) -> str | None: + if ( + self.requested_todo_id is not None + or self.receipt_bound_todo_id is not None + or self.receipt_bound_replan_obligation_id is not None + ): + return None + return self.receipt_pending_action_todo_id + + +@dataclass(frozen=True, slots=True) +class ActionSelectionPreflightResult: + rejected: bool + receipt: dict[str, object] | None + receipt_status: str + receipt_appended: bool + + +def _requested_quota_action_todo_id( + args: argparse.Namespace, +) -> str | None: + if not ( + bool(args.codex_app) + or bool(getattr(args, "trae_app", False)) + or args.runtime_profile + in {profile.value for profile in GUIDED_START_TURN_RUNTIME_PROFILES} + ): + return None + return normalize_todo_id(args.todo_id) + + +def load_requested_quota_action_selection( + args: argparse.Namespace, + *, + runtime_root: Path, + turn_instance_id: str | None, +) -> RequestedQuotaActionSelection: + requested_todo_id = _requested_quota_action_todo_id(args) + if not turn_instance_id: + return RequestedQuotaActionSelection( + requested_todo_id=requested_todo_id, + receipt=None, + receipt_bound_todo_id=None, + receipt_bound_replan_obligation_id=None, + receipt_pending_action_todo_id=None, + receipt_identity_upgraded=False, + ) + existing = find_heartbeat_receipt( + runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=turn_instance_id, + ) + if not existing: + return RequestedQuotaActionSelection( + requested_todo_id=requested_todo_id, + receipt=None, + receipt_bound_todo_id=None, + receipt_bound_replan_obligation_id=None, + receipt_pending_action_todo_id=None, + receipt_identity_upgraded=False, + ) + details_value = existing.get("details") + details: Mapping[str, object] = ( + details_value if isinstance(details_value, Mapping) else {} + ) + return RequestedQuotaActionSelection( + requested_todo_id=requested_todo_id, + receipt=existing, + receipt_bound_todo_id=heartbeat_receipt_settlement_todo_id(existing), + receipt_bound_replan_obligation_id=( + heartbeat_receipt_settlement_replan_obligation_id(existing) + ), + receipt_pending_action_todo_id=( + heartbeat_receipt_pending_action_todo_id(existing) + ), + receipt_identity_upgraded=( + str(details.get("settlement_receipt_revision") or "") == "identity_upgrade" + ), + ) + + +def _apply_requested_quota_action_selection_preflight( + payload: dict[str, object], + *, + 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: + 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 = ( + normalize_todo_id(selected_todo.get("todo_id")) + if isinstance(selected_todo, Mapping) + else None + ) + qualification_value = payload.get("action_selection_qualification") + qualification: Mapping[str, object] = ( + qualification_value if isinstance(qualification_value, Mapping) else {} + ) + if selected_todo_id is None and str(qualification.get("state") or "") == ( + "qualified" + ): + # Recovery names the prior Turn's Todo in the qualification instead of + # projecting it as the current selected Todo. + qualification_selected = qualification.get("selected_todo") + selected_todo_id = ( + normalize_todo_id(qualification_selected.get("todo_id")) + if isinstance(qualification_selected, Mapping) + else None + ) + if receipt_bound_replan_obligation_id: + if not receipt_identity_upgraded: + # A Turn that began in autonomous replan has no Todo selection + # authority to replace. A later same-Turn selection is a replay. + 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) + else None + ) + execution_obligation_value = payload.get("execution_obligation") + execution_obligation: Mapping[str, object] = ( + execution_obligation_value + if isinstance(execution_obligation_value, Mapping) + else {} + ) + interaction_value = payload.get("interaction_contract") + interaction: Mapping[str, object] = ( + interaction_value if isinstance(interaction_value, Mapping) else {} + ) + agent_channel_value = interaction.get("agent_channel") + agent_channel: Mapping[str, object] = ( + agent_channel_value if isinstance(agent_channel_value, Mapping) else {} + ) + pending_selection_delivery_qualified = ( + selection_binding == "pending_action_selection" + and payload.get("normal_delivery_allowed") is True + ) + pending_selection_workspace_repair_qualified = ( + selection_binding == "pending_action_selection" + and payload.get("workspace_repair_allowed") is True + and payload.get("effective_action") + == EffectiveAction.AGENT_WORKSPACE_REPAIR.value + and execution_obligation.get("kind") == "agent_workspace_repair" + and execution_obligation.get("must_attempt_work") is True + and agent_channel.get("must_attempt") is True + and agent_channel.get("delivery_allowed") is False + ) + exact_current_obligation_qualified = ( + selection_binding != "pending_action_selection" + and execution_obligation.get("must_attempt_work") is True + and agent_channel.get("must_attempt") is True + ) + if ( + selected_todo_id == requested_todo_id + and payload.get("ok") is True + and payload.get("should_run") is True + and ( + pending_selection_delivery_qualified + or pending_selection_workspace_repair_qualified + or exact_current_obligation_qualified + ) + ): + return False + + if not isinstance(qualification_value, Mapping): + raise QuotaActionSelectionConflictError( + QuotaActionSelectionConflictKind.UNQUALIFIED, + requested_todo_id=requested_todo_id, + selected_todo_id=selected_todo_id, + ) + qualification_state = str(qualification.get("state") or "") + if qualification_state not in {"deferred", "rejected"}: + raise QuotaActionSelectionConflictError( + QuotaActionSelectionConflictKind.CONFLICT, + requested_todo_id=requested_todo_id, + selected_todo_id=selected_todo_id, + qualification_state=qualification_state, + ) + qualification_reason = str( + qualification.get("reason") or "candidate_not_currently_eligible" + ) + deferred = qualification_state == "deferred" + auxiliary_monitor = ( + qualification_reason == "auxiliary_monitor_not_selectable_in_advancement_lane" + ) + error_code = ( + "quota_action_selection_deferred" + if deferred + else "quota_action_selection_rejected" + ) + payload.update( + { + "ok": False, + "decision": "skip", + "should_run": False, + "effective_action": EffectiveAction.QUOTA_SKIP.value, + "state": error_code, + "waiting_on": "codex", + "status": error_code, + "error_code": error_code, + "reason": ( + "explicit action selection was deferred by the current " + f"delivery frontier: {qualification_reason}" + if deferred + else "explicit action selection is not currently eligible: " + f"{qualification_reason}" + ), + "recommended_action": ( + "handle the current delivery preemption, then rerun quota " + "should-run with the same --turn-instance-id; omit --todo-id " + "first when a refreshed action portfolio is needed" + if deferred + else "the due monitor is visible as auxiliary context, not an " + "independently selectable action in the current advancement lane; " + "choose a current advancement Todo, or rerun after the monitor " + "becomes the hard lane" + if auxiliary_monitor + else "rerun quota should-run with the same --turn-instance-id " + "without --todo-id, then choose a currently eligible Todo" + ), + } + ) + return True + + +def reconcile_requested_quota_action_selection( + payload: dict[str, object], + args: argparse.Namespace, + *, + registry_path: Path, + context: QuotaCommandContext, + selection: RequestedQuotaActionSelection, +) -> ActionSelectionPreflightResult: + rejected = _apply_requested_quota_action_selection_preflight( + payload, + requested_todo_id=selection.requested_todo_id, + receipt_bound_todo_id=selection.receipt_bound_todo_id, + receipt_bound_replan_obligation_id=( + selection.receipt_bound_replan_obligation_id + ), + receipt_pending_action_todo_id=selection.receipt_pending_action_todo_id, + receipt_identity_upgraded=selection.receipt_identity_upgraded, + ) + if not rejected: + return ActionSelectionPreflightResult( + rejected=False, + receipt=selection.receipt, + receipt_status="replayed", + receipt_appended=False, + ) + + apply_action_selection_recovery( + payload, + registry_path=str(registry_path), + runtime_root=str(context.runtime_root), + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=context.heartbeat_turn_id, + available_capabilities=args.available_capabilities, + scheduler_args=render_scheduler_execution_args( + scheduler_execution_context=context.scheduler_context + ), + ) + obligation = payload.get("execution_obligation") + if isinstance(obligation, dict): + obligation.update( + must_attempt_work=False, + delivery_allowed=False, + reason=payload["recommended_action"], + ) + receipt, receipt_status, receipt_appended = _retain_deferred_action_selection( + payload, + args, + runtime_root=context.runtime_root, + turn_instance_id=context.heartbeat_turn_id, + selection=selection, + ) + return ActionSelectionPreflightResult( + rejected=True, + receipt=receipt, + receipt_status=receipt_status, + receipt_appended=receipt_appended, + ) + + +def attach_uncommitted_action_selection_receipt( + payload: dict[str, object], + *, + turn_instance_id: str, +) -> None: + """Expose an accurate non-durable receipt for a rejected preflight.""" + + payload["heartbeat_receipt"] = { + "schema_version": HEARTBEAT_RECEIPT_SCHEMA_VERSION, + "turn_instance_id": turn_instance_id, + "status": "not_committed", + "stall_observation": "not_evaluated", + "reason_code": str( + payload.get("error_code") or "quota_action_selection_rejected" + ), + } + + +def commit_requested_action_selection( + payload: Mapping[str, object], + *, + requested_todo_id: str | None, +) -> None: + """Project the exact requested selection only after receipt reconciliation.""" + + selected_todo = payload.get("selected_todo") + if ( + requested_todo_id + and isinstance(selected_todo, dict) + and normalize_todo_id(selected_todo.get("todo_id")) == requested_todo_id + ): + selected_todo["selection_binding"] = "heartbeat_receipt" + + +def _retain_deferred_action_selection( + payload: Mapping[str, object], + args: argparse.Namespace, + *, + runtime_root: Path, + turn_instance_id: str | None, + selection: RequestedQuotaActionSelection, +) -> tuple[dict[str, object] | None, str, bool]: + """Append a deferred explicit choice without granting settlement authority.""" + + qualification = payload.get("action_selection_qualification") + if ( + not turn_instance_id + or selection.receipt is None + or selection.requested_todo_id is None + or not isinstance(qualification, Mapping) + or qualification.get("state") != "deferred" + ): + return selection.receipt, "replayed", False + retained, appended = retain_pending_heartbeat_action_selection( + runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_instance_id=turn_instance_id, + todo_id=selection.requested_todo_id, + reason=str(qualification.get("reason") or "current_delivery_gate"), + ) + return retained, "selection_retained" if appended else "replayed", appended diff --git a/loopx/semantics/vocabulary_v0.json b/loopx/semantics/vocabulary_v0.json index f0c4b66e2..a8caaf5e3 100644 --- a/loopx/semantics/vocabulary_v0.json +++ b/loopx/semantics/vocabulary_v0.json @@ -496,7 +496,7 @@ "unsettled_host_turn_recovery": "A host Turn is unsettled and must be recovered before anything else: the selected Todo and action portfolio are dropped, should_run is set, and normal, recovery and self-repair delivery are all refused." }, "producers": [ - "loopx/cli_commands/quota.py::_apply_requested_quota_action_selection_preflight", + "loopx/cli_commands/quota_action_selection.py::_apply_requested_quota_action_selection_preflight", "loopx/control_plane/quota/decision_summary.py::_task_orchestration_effective_action", "loopx/control_plane/quota/decision_summary.py::quota_effective_action", "loopx/control_plane/quota/decision_summary.py::resolve_quota_run_decision", diff --git a/tests/control_plane/test_quota_action_selection_conflict.py b/tests/control_plane/test_quota_action_selection_conflict.py index 0018475fc..00ba45a71 100644 --- a/tests/control_plane/test_quota_action_selection_conflict.py +++ b/tests/control_plane/test_quota_action_selection_conflict.py @@ -7,7 +7,9 @@ import pytest -from loopx.cli_commands.quota import _apply_requested_quota_action_selection_preflight +from loopx.cli_commands.quota_action_selection import ( + _apply_requested_quota_action_selection_preflight, +) from loopx.cli_commands.quota_failure_report import quota_failure_payload from loopx.control_plane.quota.error_codes import ( QuotaActionSelectionConflictError,