diff --git a/loopx/control_plane/agents/agent_lane_recommendation.py b/loopx/control_plane/agents/agent_lane_recommendation.py index 15968b7c6c..6d78560ba4 100644 --- a/loopx/control_plane/agents/agent_lane_recommendation.py +++ b/loopx/control_plane/agents/agent_lane_recommendation.py @@ -13,6 +13,7 @@ from ..todos.summary_item import compact_todo_summary_item from ..work_items.primary_action import protocol_action_text from ..work_items.work_lane import ( + ReceiptBoundMonitorPhase, work_lane_contract_is_due_monitor_attempt, work_lane_contract_is_lark_inbox_reply_due, ) @@ -42,7 +43,12 @@ def build_receipt_bound_monitor_next_action( available_capabilities: Any, receipt_bound_todo_id: str | None, ) -> dict[str, Any] | None: - """Recover an exact due monitor omitted by priority-limited hot lanes.""" + """Recover an exact receipt-bound monitor omitted by compact hot lanes. + + A monitor remains the same-turn settlement identity after a successful poll + reschedules it into the future. That replay must preserve the parent Todo + without asking the agent to observe the target again. + """ if not isinstance(agent_identity, dict): return None @@ -53,13 +59,16 @@ def build_receipt_bound_monitor_next_action( for item in agent_todo_items: if normalize_todo_id(item.get("todo_id")) != todo_id: continue + monitor_due = todo_item_is_due_monitor(item) if ( not _todo_item_is_actionable_open(item) or _todo_task_class(item) != TODO_TASK_CLASS_MONITOR - or not todo_item_is_due_monitor(item) - or missing_required_capabilities( - item, - available_capabilities=available_capabilities, + or ( + monitor_due + and missing_required_capabilities( + item, + available_capabilities=available_capabilities, + ) ) or not agent_scope_item_claimed_by_agent_or_unclaimed( item, @@ -80,6 +89,11 @@ def build_receipt_bound_monitor_next_action( "confidence": "selected", "preserves_goal_next_action": True, "selection_binding": "heartbeat_receipt", + "receipt_bound_monitor_phase": ( + ReceiptBoundMonitorPhase.POLL_DUE + if monitor_due + else ReceiptBoundMonitorPhase.SETTLEMENT_PENDING + ).value, } ) if not agent_scope_item_claimed_by(item): diff --git a/loopx/control_plane/work_items/work_lane.py b/loopx/control_plane/work_items/work_lane.py index f511d48b44..d146c2bd6b 100644 --- a/loopx/control_plane/work_items/work_lane.py +++ b/loopx/control_plane/work_items/work_lane.py @@ -1,5 +1,6 @@ from __future__ import annotations +from enum import StrEnum from typing import Any from ..todos.contract import TODO_TASK_CLASS_MONITOR, normalize_todo_id @@ -7,16 +8,29 @@ WORK_LANE_CONTRACT_SCHEMA_VERSION = "work_lane_contract_v1" +WORK_LANE_RECEIPT_BOUND_MONITOR_SETTLEMENT_OBLIGATION = ( + "settle_receipt_bound_monitor" +) WORK_LANE_CURRENT_AGENT_MONITOR_REPAIR_OBLIGATIONS = { "attempt_due_monitor", "repair_monitor_schedule_metadata", "repair_resume_gate_or_close_standing_monitor", + WORK_LANE_RECEIPT_BOUND_MONITOR_SETTLEMENT_OBLIGATION, } WORK_LANE_EXTERNAL_EVIDENCE_OBSERVATION_OBLIGATION = ( "observe_external_evidence_or_blocker" ) WORK_LANE_LARK_INBOX_REPLY_DUE_OBLIGATION = "drain_lark_inbox_reply_due" WORK_LANE_TODO_MONITOR_DUE_KIND = "todo_monitor_due" + + +class ReceiptBoundMonitorPhase(StrEnum): + """Same-turn phase of a monitor selected by a heartbeat receipt.""" + + POLL_DUE = "poll_due" + SETTLEMENT_PENDING = "settlement_pending" + + WORK_LANE_TODO_ITEM_FIELDS = ( "index", "text", @@ -119,9 +133,7 @@ def preserve_heartbeat_receipt_bound_work_lane( ) -> dict[str, Any] | None: """Keep a committed same-turn Todo binding ahead of a newly due monitor.""" - if not isinstance(contract, dict) or not work_lane_contract_is_due_monitor_attempt( - contract - ): + if not isinstance(contract, dict): return contract if not isinstance(selected_todo, dict): return contract @@ -129,6 +141,41 @@ def preserve_heartbeat_receipt_bound_work_lane( if not todo_id or selected_todo.get("selection_binding") != "heartbeat_receipt": return contract if selected_todo.get("task_class") == TODO_TASK_CLASS_MONITOR: + try: + monitor_phase = ReceiptBoundMonitorPhase( + selected_todo.get("receipt_bound_monitor_phase") + ) + except (TypeError, ValueError) as exc: + raise ValueError( + "receipt-bound monitor selection requires an explicit monitor phase" + ) from exc + if monitor_phase is ReceiptBoundMonitorPhase.SETTLEMENT_PENDING: + return { + **contract, + "schema_version": WORK_LANE_CONTRACT_SCHEMA_VERSION, + "lane": "continuous_monitor", + "obligation": ( + WORK_LANE_RECEIPT_BOUND_MONITOR_SETTLEMENT_OBLIGATION + ), + "must_attempt_work": True, + "selection_binding": "heartbeat_receipt", + "selected_todo_id": todo_id, + "monitor_due_count": 0, + "monitor_due_items": [], + "receipt_bound_monitor_item": selected_todo, + "reason_codes": [ + "heartbeat_receipt_bound_replay", + "monitor_poll_already_recorded", + "same_turn_settlement_identity", + ], + "monitor_policy": "settle_receipt_bound_monitor_without_repoll", + "deferred_work_lane": contract, + "action": ( + "finish refresh and quota settlement for the already-polled " + "monitor bound to this heartbeat turn; do not poll again; " + "reconsider its successor on the next turn" + ), + } return { **contract, "schema_version": WORK_LANE_CONTRACT_SCHEMA_VERSION, @@ -150,6 +197,8 @@ def preserve_heartbeat_receipt_bound_work_lane( "reconsider newly runnable monitor priority on the next turn" ), } + if not work_lane_contract_is_due_monitor_attempt(contract): + return contract return { "schema_version": WORK_LANE_CONTRACT_SCHEMA_VERSION, "lane": "advancement_task", diff --git a/loopx/state_projection.py b/loopx/state_projection.py index 03010834b9..438bb88dbb 100644 --- a/loopx/state_projection.py +++ b/loopx/state_projection.py @@ -390,16 +390,29 @@ def _section_lines(state_text: str, heading: str) -> list[str]: return lines -def _section_entries(lines: list[str]) -> list[str]: +def _section_entries( + lines: list[str], + *, + text_limit: int | None = 220, +) -> list[str]: entries: list[str] = [] current: list[str] = [] + + def append_entry(parts: list[str]) -> None: + text = compact_todo_text(" ".join(parts)) + if not text: + return + entries.append( + text if text_limit is None else _compact_text(text, limit=text_limit) + ) + for line in lines: if line.strip().startswith("