Skip to content

Commit 6e5d540

Browse files
authored
Merge pull request #4775 from songoow/codex/coalesce-live-packet-rendering
refactor(quota): coalesce live projection packet rendering
2 parents dc03c8c + 6598103 commit 6e5d540

2 files changed

Lines changed: 123 additions & 12 deletions

File tree

‎loopx/control_plane/quota/live_decision.py‎

Lines changed: 16 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -184,12 +184,12 @@ def _project_turn_start_required_reads(
184184
),
185185
turn_instance_id: str | None,
186186
runtime_root: Path,
187-
) -> None:
188-
"""Order fresh operator evidence before work without changing work selection."""
187+
) -> bool:
188+
"""Order evidence before work; return whether the packet needs rendering."""
189189

190190
projected = _turn_start_required_reads(dispatch)
191191
if not projected:
192-
return
192+
return False
193193
existing = payload.get("required_reads")
194194
required_reads = (
195195
[dict(item) for item in existing if isinstance(item, Mapping)]
@@ -227,7 +227,7 @@ def _project_turn_start_required_reads(
227227
turn_instance_id=turn_instance_id,
228228
runtime_root=str(runtime_root),
229229
)
230-
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
230+
return True
231231

232232

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

276276
if not isinstance(projection, Mapping) or projection.get("state") != "pending":
277-
return
277+
return False
278278
summary = str(projection.get("action_summary") or "").strip()
279279
command = str(projection.get("command") or "").strip()
280280
if not summary or not command:
281-
return
281+
return False
282282
payload.update(
283283
{
284284
"decision": "run",
@@ -334,7 +334,7 @@ def _apply_pending_capability_intent_precedence(
334334
scheduler_execution_context=scheduler_execution_context,
335335
turn_instance_id=turn_instance_id,
336336
)
337-
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
337+
return True
338338

339339

340340
def bind_scheduler_followup_cli_routes(
@@ -593,7 +593,7 @@ def build_live_quota_should_run_decision(
593593
turn_start_hook_dispatch, registry=registry_path, runtime_root=runtime_root,
594594
goal_id=goal_id, agent_id=agent_id,
595595
)
596-
_project_turn_start_required_reads(
596+
packet_changed = _project_turn_start_required_reads(
597597
payload,
598598
turn_start_hook_dispatch,
599599
available_capabilities=available_capabilities,
@@ -604,16 +604,21 @@ def build_live_quota_should_run_decision(
604604
hook_dispatch = dispatch_interaction_projection_hooks(interaction_projection_hooks)
605605
projections = hook_dispatch["projections"]
606606
if isinstance(projections, Mapping):
607-
_apply_pending_capability_intent_precedence(
607+
intent_changed = _apply_pending_capability_intent_precedence(
608608
payload,
609609
projections.get("pending_capability_intent"),
610610
available_capabilities=available_capabilities,
611611
scheduler_execution_context=resolved_context,
612612
turn_instance_id=turn_instance_id,
613613
)
614+
packet_changed = packet_changed or intent_changed
614615
interaction = payload.get("interaction_contract")
615616
if isinstance(interaction, dict):
616617
interaction.update(projections)
618+
# Neither projection consumes the intermediate packet. Recovery below does
619+
# require a complete decision, so finalize this projection stage here.
620+
if packet_changed:
621+
payload["protocol_action_packet"] = build_protocol_action_packet(payload)
617622
# A settled receipt owns this host Turn until it ends. Looking for an older
618623
# unsettled Turn here can overwrite the settled-skip route with a recovery
619624
# obligation and then select a successor against the immutable receipt

‎tests/control_plane/test_effect_turn_live_quota_decision.py‎

Lines changed: 107 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -190,8 +190,9 @@ def test_live_quota_decision_maps_to_effect_turn(tmp_path: Path) -> None:
190190
assert turn.next_effect.cli_actions[0].startswith("loopx --runtime-root ")
191191

192192

193+
@pytest.mark.parametrize("reads", [False, True])
193194
def test_managed_turn_projects_prior_unsettled_heartbeat_recovery(
194-
tmp_path: Path,
195+
tmp_path: Path, reads: bool,
195196
) -> None:
196197
runtime_root = tmp_path / "runtime"
197198
agent_id = "codex-fixture"
@@ -261,6 +262,7 @@ def test_managed_turn_projects_prior_unsettled_heartbeat_recovery(
261262
codex_app_current_rrule=None,
262263
registry_path=tmp_path / "registry.json",
263264
runtime_root=runtime_root,
265+
turn_start_hook_dispatch=_turn_start_dispatch(required=reads),
264266
route_source="loopx_turn_plan",
265267
turn_instance_id="managed-current-turn",
266268
scheduler_execution_context={
@@ -1002,3 +1004,107 @@ def test_prior_closeout_identity_conflict_fails_closed(
10021004
"execution_mode": "interactive",
10031005
},
10041006
)
1007+
1008+
1009+
@pytest.mark.parametrize("status_name", ["active", "paused"])
1010+
@pytest.mark.parametrize("reads", [False, True])
1011+
@pytest.mark.parametrize("intent", ["absent", "pending", "invalid"])
1012+
def test_live_projection_stage_renders_one_complete_packet(
1013+
tmp_path: Path, monkeypatch: pytest.MonkeyPatch,
1014+
status_name: str, reads: bool, intent: str,
1015+
) -> None:
1016+
"""Coalescing must preserve the full decision and signed envelope input."""
1017+
from copy import deepcopy
1018+
from loopx.control_plane.quota import live_decision
1019+
from loopx.control_plane.capability_hooks import (
1020+
InteractionProjectionHookRegistration,
1021+
INTERACTION_PROJECTION_HOOK_RESULT_SCHEMA_VERSION,
1022+
)
1023+
from loopx.control_plane.quota.turn_envelope import quota_action_signature_document
1024+
1025+
projection = {
1026+
"schema_version": "pending_capability_intent_projection_v0",
1027+
"capability_id": "periodic-report",
1028+
"intent_kind": "periodic_report.trigger_evaluation",
1029+
"idempotency_key": "periodic-report:fixture",
1030+
"intent_digest": "sha256:" + "a" * 64,
1031+
"goal_id": GOAL_ID,
1032+
"agent_id": "fixture-agent",
1033+
"state": "pending",
1034+
"action_kind": "consume_periodic_report_intent",
1035+
"action_summary": "Generate the exact report and queue configured delivery.",
1036+
"command": f"loopx periodic-report consume-pending --goal-id {GOAL_ID} --agent-id fixture-agent --execute",
1037+
"generation_authorized": True,
1038+
"external_delivery_authorized": True,
1039+
"agent_read_required": True,
1040+
}
1041+
if intent == "invalid":
1042+
projection["command"] = ""
1043+
hook = InteractionProjectionHookRegistration(
1044+
hook_id="periodic_report.pending_intent",
1045+
capability_id="periodic-report",
1046+
projection_slots=("pending_capability_intent",),
1047+
requested_read_scope=("post_writeback_intent_journal",),
1048+
producer=lambda: {
1049+
"schema_version": INTERACTION_PROJECTION_HOOK_RESULT_SCHEMA_VERSION,
1050+
"hook_id": "periodic_report.pending_intent",
1051+
"capability_id": "periodic-report",
1052+
"phase": "interaction_projection",
1053+
"status": "candidate",
1054+
"projection_slot": "pending_capability_intent",
1055+
"payload": projection,
1056+
},
1057+
)
1058+
status = _ordinary_status_payload()
1059+
status["attention_queue"]["items"][0]["status"] = status_name
1060+
status["run_history"]["goals"][0]["status"] = status_name
1061+
kwargs = dict(
1062+
goal_id=GOAL_ID, agent_id=None, available_capabilities=["shell"],
1063+
include_scheduler_detail=False, codex_app_current_rrule=None,
1064+
registry_path=tmp_path / "registry.json", runtime_root=tmp_path / "runtime",
1065+
interaction_projection_hooks=[] if intent == "absent" else [hook],
1066+
turn_start_hook_dispatch=_turn_start_dispatch(required=reads),
1067+
)
1068+
rendered = []
1069+
render = live_decision.build_protocol_action_packet
1070+
1071+
def record_render(payload):
1072+
rendered.append(deepcopy(payload))
1073+
return render(payload)
1074+
1075+
recover = live_decision.apply_unsettled_host_turn_recovery_if_required
1076+
1077+
def check_recovery_input(payload, **recovery_kwargs):
1078+
assert payload["protocol_action_packet"] == render(payload)
1079+
return recover(payload, **recovery_kwargs)
1080+
1081+
monkeypatch.setattr(
1082+
live_decision, "apply_unsettled_host_turn_recovery_if_required", check_recovery_input
1083+
)
1084+
monkeypatch.setattr(live_decision, "build_protocol_action_packet", record_render)
1085+
actual = build_live_quota_should_run_decision(deepcopy(status), **kwargs)
1086+
assert len(rendered) == int(reads or intent == "pending")
1087+
if rendered:
1088+
assert actual["protocol_action_packet"] == render(actual)
1089+
if intent == "pending":
1090+
assert actual["effective_action"] == "governed_capability_intent"
1091+
assert actual["interaction_contract"]["cli_channel"]["next_cli_actions"] == [projection["command"]]
1092+
assert actual["heartbeat_recommendation"]["notify"] == "DONT_NOTIFY"
1093+
if reads:
1094+
assert actual["interaction_contract"]["agent_channel"]["required_reads"]
1095+
1096+
# Reproduce the previous composition order: eagerly render after each
1097+
# changed helper, then compare every output field (including summary order).
1098+
for name in ("_project_turn_start_required_reads", "_apply_pending_capability_intent_precedence"):
1099+
original = getattr(live_decision, name)
1100+
1101+
def eager(payload, *args, _original=original, **helper_kwargs):
1102+
changed = _original(payload, *args, **helper_kwargs)
1103+
if changed:
1104+
payload["protocol_action_packet"] = render(payload)
1105+
return False
1106+
1107+
monkeypatch.setattr(live_decision, name, eager)
1108+
expected = build_live_quota_should_run_decision(deepcopy(status), **kwargs)
1109+
assert actual == expected
1110+
assert quota_action_signature_document(actual) == quota_action_signature_document(expected)

0 commit comments

Comments
 (0)