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
27 changes: 16 additions & 11 deletions loopx/control_plane/quota/live_decision.py
Original file line number Diff line number Diff line change
Expand Up @@ -184,12 +184,12 @@ def _project_turn_start_required_reads(
),
turn_instance_id: str | None,
runtime_root: Path,
) -> None:
"""Order fresh operator evidence before work without changing work selection."""
) -> bool:
"""Order evidence before work; return whether the packet needs rendering."""

projected = _turn_start_required_reads(dispatch)
if not projected:
return
return False
existing = payload.get("required_reads")
required_reads = (
[dict(item) for item in existing if isinstance(item, Mapping)]
Expand Down Expand Up @@ -227,7 +227,7 @@ def _project_turn_start_required_reads(
turn_instance_id=turn_instance_id,
runtime_root=str(runtime_root),
)
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
return True


def _fresh_operator_inbox_observation_count(
Expand Down Expand Up @@ -270,15 +270,15 @@ def _apply_pending_capability_intent_precedence(
Mapping[str, Any] | SchedulerExecutionContextResolution | None
) = None,
turn_instance_id: str | None = None,
) -> None:
"""Wake one governed local capability action ahead of quiet/terminal routes."""
) -> bool:
"""Apply intent precedence; return whether the packet needs rendering."""

if not isinstance(projection, Mapping) or projection.get("state") != "pending":
return
return False
summary = str(projection.get("action_summary") or "").strip()
command = str(projection.get("command") or "").strip()
if not summary or not command:
return
return False
payload.update(
{
"decision": "run",
Expand Down Expand Up @@ -334,7 +334,7 @@ def _apply_pending_capability_intent_precedence(
scheduler_execution_context=scheduler_execution_context,
turn_instance_id=turn_instance_id,
)
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
return True


def bind_scheduler_followup_cli_routes(
Expand Down Expand Up @@ -593,7 +593,7 @@ def build_live_quota_should_run_decision(
turn_start_hook_dispatch, registry=registry_path, runtime_root=runtime_root,
goal_id=goal_id, agent_id=agent_id,
)
_project_turn_start_required_reads(
packet_changed = _project_turn_start_required_reads(
payload,
turn_start_hook_dispatch,
available_capabilities=available_capabilities,
Expand All @@ -604,16 +604,21 @@ def build_live_quota_should_run_decision(
hook_dispatch = dispatch_interaction_projection_hooks(interaction_projection_hooks)
projections = hook_dispatch["projections"]
if isinstance(projections, Mapping):
_apply_pending_capability_intent_precedence(
intent_changed = _apply_pending_capability_intent_precedence(
payload,
projections.get("pending_capability_intent"),
available_capabilities=available_capabilities,
scheduler_execution_context=resolved_context,
turn_instance_id=turn_instance_id,
)
packet_changed = packet_changed or intent_changed
interaction = payload.get("interaction_contract")
if isinstance(interaction, dict):
interaction.update(projections)
# Neither projection consumes the intermediate packet. Recovery below does
# require a complete decision, so finalize this projection stage here.
if packet_changed:
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
# 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
Expand Down
108 changes: 107 additions & 1 deletion tests/control_plane/test_effect_turn_live_quota_decision.py
Original file line number Diff line number Diff line change
Expand Up @@ -190,8 +190,9 @@ def test_live_quota_decision_maps_to_effect_turn(tmp_path: Path) -> None:
assert turn.next_effect.cli_actions[0].startswith("loopx --runtime-root ")


@pytest.mark.parametrize("reads", [False, True])
def test_managed_turn_projects_prior_unsettled_heartbeat_recovery(
tmp_path: Path,
tmp_path: Path, reads: bool,
) -> None:
runtime_root = tmp_path / "runtime"
agent_id = "codex-fixture"
Expand Down Expand Up @@ -261,6 +262,7 @@ def test_managed_turn_projects_prior_unsettled_heartbeat_recovery(
codex_app_current_rrule=None,
registry_path=tmp_path / "registry.json",
runtime_root=runtime_root,
turn_start_hook_dispatch=_turn_start_dispatch(required=reads),
route_source="loopx_turn_plan",
turn_instance_id="managed-current-turn",
scheduler_execution_context={
Expand Down Expand Up @@ -1002,3 +1004,107 @@ def test_prior_closeout_identity_conflict_fails_closed(
"execution_mode": "interactive",
},
)


@pytest.mark.parametrize("status_name", ["active", "paused"])
@pytest.mark.parametrize("reads", [False, True])
@pytest.mark.parametrize("intent", ["absent", "pending", "invalid"])
def test_live_projection_stage_renders_one_complete_packet(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch,
status_name: str, reads: bool, intent: str,
) -> None:
"""Coalescing must preserve the full decision and signed envelope input."""
from copy import deepcopy
from loopx.control_plane.quota import live_decision
from loopx.control_plane.capability_hooks import (
InteractionProjectionHookRegistration,
INTERACTION_PROJECTION_HOOK_RESULT_SCHEMA_VERSION,
)
from loopx.control_plane.quota.turn_envelope import quota_action_signature_document

projection = {
"schema_version": "pending_capability_intent_projection_v0",
"capability_id": "periodic-report",
"intent_kind": "periodic_report.trigger_evaluation",
"idempotency_key": "periodic-report:fixture",
"intent_digest": "sha256:" + "a" * 64,
"goal_id": GOAL_ID,
"agent_id": "fixture-agent",
"state": "pending",
"action_kind": "consume_periodic_report_intent",
"action_summary": "Generate the exact report and queue configured delivery.",
"command": f"loopx periodic-report consume-pending --goal-id {GOAL_ID} --agent-id fixture-agent --execute",
"generation_authorized": True,
"external_delivery_authorized": True,
"agent_read_required": True,
}
if intent == "invalid":
projection["command"] = ""
hook = InteractionProjectionHookRegistration(
hook_id="periodic_report.pending_intent",
capability_id="periodic-report",
projection_slots=("pending_capability_intent",),
requested_read_scope=("post_writeback_intent_journal",),
producer=lambda: {
"schema_version": INTERACTION_PROJECTION_HOOK_RESULT_SCHEMA_VERSION,
"hook_id": "periodic_report.pending_intent",
"capability_id": "periodic-report",
"phase": "interaction_projection",
"status": "candidate",
"projection_slot": "pending_capability_intent",
"payload": projection,
},
)
status = _ordinary_status_payload()
status["attention_queue"]["items"][0]["status"] = status_name
status["run_history"]["goals"][0]["status"] = status_name
kwargs = dict(
goal_id=GOAL_ID, agent_id=None, available_capabilities=["shell"],
include_scheduler_detail=False, codex_app_current_rrule=None,
registry_path=tmp_path / "registry.json", runtime_root=tmp_path / "runtime",
interaction_projection_hooks=[] if intent == "absent" else [hook],
turn_start_hook_dispatch=_turn_start_dispatch(required=reads),
)
rendered = []
render = live_decision.build_protocol_action_packet

def record_render(payload):
rendered.append(deepcopy(payload))
return render(payload)

recover = live_decision.apply_unsettled_host_turn_recovery_if_required

def check_recovery_input(payload, **recovery_kwargs):
assert payload["protocol_action_packet"] == render(payload)
return recover(payload, **recovery_kwargs)

monkeypatch.setattr(
live_decision, "apply_unsettled_host_turn_recovery_if_required", check_recovery_input
)
monkeypatch.setattr(live_decision, "build_protocol_action_packet", record_render)
actual = build_live_quota_should_run_decision(deepcopy(status), **kwargs)
assert len(rendered) == int(reads or intent == "pending")
if rendered:
assert actual["protocol_action_packet"] == render(actual)
if intent == "pending":
assert actual["effective_action"] == "governed_capability_intent"
assert actual["interaction_contract"]["cli_channel"]["next_cli_actions"] == [projection["command"]]
assert actual["heartbeat_recommendation"]["notify"] == "DONT_NOTIFY"
if reads:
assert actual["interaction_contract"]["agent_channel"]["required_reads"]

# Reproduce the previous composition order: eagerly render after each
# changed helper, then compare every output field (including summary order).
for name in ("_project_turn_start_required_reads", "_apply_pending_capability_intent_precedence"):
original = getattr(live_decision, name)

def eager(payload, *args, _original=original, **helper_kwargs):
changed = _original(payload, *args, **helper_kwargs)
if changed:
payload["protocol_action_packet"] = render(payload)
return False

monkeypatch.setattr(live_decision, name, eager)
expected = build_live_quota_should_run_decision(deepcopy(status), **kwargs)
assert actual == expected
assert quota_action_signature_document(actual) == quota_action_signature_document(expected)
Loading