diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 414749992..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,16 +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 ( - bind_action_selection_recovery_command, - build_action_selection_recovery_fields, - current_action_selection_admission, -) from ..presentation.renderers.quota_event_markdown import ( render_quota_monitor_poll_markdown, render_quota_slot_preview_markdown, @@ -83,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, @@ -110,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( @@ -163,220 +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 _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 = None, - receipt_identity_upgraded: bool = False, -) -> dict[str, object] | None: - if not requested_todo_id: - return None - 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 None - interaction = payload.get("interaction_contract") - agent_channel = ( - interaction.get("agent_channel") - if isinstance(interaction, Mapping) - else None - ) - agent_channel = agent_channel if isinstance(agent_channel, Mapping) else {} - selected_todo_id, admitted = current_action_selection_admission( - payload, - requested_todo_id=requested_todo_id, - agent_must_attempt=agent_channel.get("must_attempt") is True, - agent_delivery_refused=agent_channel.get("delivery_allowed") is False, - ) - 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 None - if requested_todo_id == receipt_pending_action_todo_id: - return None - raise QuotaActionSelectionConflictError( - QuotaActionSelectionConflictKind.CONFLICT, - requested_todo_id=requested_todo_id, - selected_todo_id=receipt_pending_action_todo_id, - qualification_state="retained_selection", - ) - if admitted: - return None - - qualification_value = payload.get("action_selection_qualification") - qualification: Mapping[str, object] = ( - qualification_value if isinstance(qualification_value, Mapping) else {} - ) - 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, - ) - return build_action_selection_recovery_fields(payload) - - -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: - recovery = _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 recovery is None: - return False - payload.update(recovery) - bind_action_selection_recovery_command( - 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 - ), - ) - return True - - -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], @@ -386,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.""" @@ -451,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, @@ -571,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 @@ -623,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, @@ -653,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 @@ -723,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 @@ -833,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, ) @@ -857,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 = ( @@ -924,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..a1cee4826 --- /dev/null +++ b/loopx/cli_commands/quota_action_selection.py @@ -0,0 +1,314 @@ +from __future__ import annotations + +import argparse +from collections.abc import Mapping +from dataclasses import dataclass +from pathlib import Path + +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 ( + bind_action_selection_recovery_command, + build_action_selection_recovery_fields, + current_action_selection_admission, +) +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 _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 = None, + receipt_identity_upgraded: bool = False, +) -> dict[str, object] | None: + if not requested_todo_id: + return None + 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 None + interaction = payload.get("interaction_contract") + agent_channel = ( + interaction.get("agent_channel") if isinstance(interaction, Mapping) else None + ) + agent_channel = agent_channel if isinstance(agent_channel, Mapping) else {} + selected_todo_id, admitted = current_action_selection_admission( + payload, + requested_todo_id=requested_todo_id, + agent_must_attempt=agent_channel.get("must_attempt") is True, + agent_delivery_refused=agent_channel.get("delivery_allowed") is False, + ) + 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 None + if requested_todo_id == receipt_pending_action_todo_id: + return None + raise QuotaActionSelectionConflictError( + QuotaActionSelectionConflictKind.CONFLICT, + requested_todo_id=requested_todo_id, + selected_todo_id=receipt_pending_action_todo_id, + qualification_state="retained_selection", + ) + if admitted: + return None + + qualification_value = payload.get("action_selection_qualification") + qualification: Mapping[str, object] = ( + qualification_value if isinstance(qualification_value, Mapping) else {} + ) + 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, + ) + return build_action_selection_recovery_fields(payload) + + +def reconcile_requested_quota_action_selection( + payload: dict[str, object], + args: argparse.Namespace, + *, + registry_path: Path, + context: QuotaCommandContext, + selection: RequestedQuotaActionSelection, +) -> ActionSelectionPreflightResult: + recovery = _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 recovery is None: + return ActionSelectionPreflightResult( + rejected=False, + receipt=selection.receipt, + receipt_status="replayed", + receipt_appended=False, + ) + + payload.update(recovery) + bind_action_selection_recovery_command( + 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 + ), + ) + 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/tests/control_plane/test_quota_action_selection_conflict.py b/tests/control_plane/test_quota_action_selection_conflict.py index edaa578cd..bf29f31b3 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 _requested_quota_action_selection_preflight +from loopx.cli_commands.quota_action_selection import ( + _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, diff --git a/tests/control_plane/test_unadmitted_selection_construction.py b/tests/control_plane/test_unadmitted_selection_construction.py index f38d74239..3c10774e4 100644 --- a/tests/control_plane/test_unadmitted_selection_construction.py +++ b/tests/control_plane/test_unadmitted_selection_construction.py @@ -6,7 +6,9 @@ import pytest -from loopx.cli_commands.quota import _requested_quota_action_selection_preflight +from loopx.cli_commands.quota_action_selection import ( + _requested_quota_action_selection_preflight, +) from loopx.control_plane.work_items import interaction_contract from loopx.control_plane.work_items.action_selection_contract import ( action_selection_needs_recovery,