Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
181 changes: 45 additions & 136 deletions loopx/cli_commands/quota.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,11 @@
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 ..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,
Expand Down Expand Up @@ -203,110 +207,58 @@ def _heartbeat_quota_action_selection_bindings(
)


def _apply_requested_quota_action_selection_preflight(
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,
) -> bool:
) -> dict[str, object] | None:
if not requested_todo_id:
return False
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 False
selected_todo = payload.get("selected_todo")
selected_todo_id = (
normalize_todo_id(selected_todo.get("todo_id"))
if isinstance(selected_todo, Mapping)
return None
interaction = payload.get("interaction_contract")
agent_channel = (
interaction.get("agent_channel")
if isinstance(interaction, Mapping)
else None
)
qualification_value = payload.get("action_selection_qualification")
qualification: Mapping[str, object] = (
qualification_value if isinstance(qualification_value, Mapping) else {}
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 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
return None
if requested_todo_id == receipt_pending_action_todo_id:
return False
return None
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 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,
Expand All @@ -321,52 +273,7 @@ def _apply_requested_quota_action_selection_preflight(
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
return build_action_selection_recovery_fields(payload)


def _reconcile_requested_quota_action_selection(
Expand All @@ -380,27 +287,29 @@ def _reconcile_requested_quota_action_selection(
receipt_pending_action_todo_id: str | None,
receipt_identity_upgraded: bool,
) -> bool:
rejected = _apply_requested_quota_action_selection_preflight(
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 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
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(
Expand Down
17 changes: 17 additions & 0 deletions loopx/control_plane/quota/heartbeat_recommendation.py
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,23 @@ def _recommendation(*parts: dict[str, Any]) -> dict[str, Any]:
return payload


def build_action_selection_recovery_recommendation(
*, reason: str,
) -> dict[str, Any]:
"""Project closed heartbeat guidance after a typed selection refusal."""

return _recommendation(
{
"source": "action_selection_recovery",
"recommended_mode": "quota_skip",
"notify": "DONT_NOTIFY",
"spend_policy": "no quota spend until an eligible Todo is selected",
"reason": reason,
"agent_must_attempt": False,
}
)


def _stall_self_repair_rule(
facts: _HeartbeatRecommendationFacts,
) -> dict[str, Any] | None:
Expand Down
55 changes: 48 additions & 7 deletions loopx/control_plane/quota/should_run_packet.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
quota_execution_profile_summary as _quota_execution_profile_summary,
)
from ..quota.heartbeat_recommendation import (
build_action_selection_recovery_recommendation,
build_heartbeat_recommendation,
refine_heartbeat_recommendation,
)
Expand Down Expand Up @@ -102,9 +103,13 @@
qualify_action_selection_from_inventory,
)
from ..work_items.execution_obligation import build_execution_obligation
from ..work_items.action_selection_contract import (
apply_action_selection_recovery_projection,
)
from ..work_items.goal_route_hint import build_goal_route_hint
from ..work_items.interaction_contract import (
build_interaction_contract,
unadmitted_action_selection,
build_protocol_action_packet,
finalize_user_gate_notification_cooldown,
)
Expand Down Expand Up @@ -366,6 +371,40 @@ def _apply_agent_monitor_only_precedence(
clear_quota_action_projections(payload)


def _apply_unadmitted_action_selection_precedence(
payload: dict[str, Any],
*,
replay_phase: ReceiptBoundReplayPhase | None,
) -> None:
"""Finalize one closed selection recovery before shared projections."""

if replay_phase is ReceiptBoundReplayPhase.SETTLED or not (
unadmitted_action_selection(payload)
):
return
for field in (
"agent_lane_next_action",
"agent_scope_frontier",
"autonomous_replan_obligation",
"execution_profile",
"goal_route_hint",
"handoff_readiness",
"replan_action_packet",
"scoped_user_gate_fallback",
"selected_todo",
"task_orchestration_contract",
"todo_id",
"todo_write_hint",
"work_lane_contract",
"workspace_guard",
):
payload.pop(field, None)
apply_action_selection_recovery_projection(payload)
payload["heartbeat_recommendation"] = (
build_action_selection_recovery_recommendation(
reason=str(payload.get("reason") or "")
)
)
def _delivery_preemptions_for_route(
prepared: _QuotaDecisionPreparation,
*,
Expand Down Expand Up @@ -1310,12 +1349,8 @@ def _build_active_quota_payload(
next_action_warning=route.next_action_warning,
replan_obligation=prepared.replan_obligation,
)
bounded_research_frontier = (
prepared.status_payload.get("bounded_research_frontier")
if isinstance(
prepared.status_payload.get("bounded_research_frontier"), dict
)
else None
bounded_research_frontier = _dict_field(
prepared.status_payload, "bounded_research_frontier"
)
_attach_truthy_fields(
payload,
Expand All @@ -1326,7 +1361,13 @@ def _build_active_quota_payload(
monitor_only=prepared.agent_monitor_only,
inbox_priority_due=prepared.inbox_priority_due,
)
if isinstance(payload.get("autonomous_replan_obligation"), dict):
_apply_unadmitted_action_selection_precedence(
payload,
replay_phase=prepared.receipt_bound_replay_phase,
)
if isinstance(
payload.get("autonomous_replan_obligation"), dict
) and not unadmitted_action_selection(payload):
payload["replan_action_packet"] = build_replan_action_packet(
payload["autonomous_replan_obligation"],
goal_id=prepared.safe_goal_id,
Expand Down
Loading
Loading