From 39ac7210a9e77462e11ccdccdf5fd03036501e67 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:27:35 +0800 Subject: [PATCH] refactor(turn): derive every Turn decision from one shared owner `run-once` rebuilt the fresh Turn decision inline after #4443, so the status read, scheduler context, route source, capability hooks and decision arguments were maintained twice and the later fix only recovered part of them. Both commands now build through `build_fresh_turn_decision_owner`, which also carries the inputs a later spend or scheduler re-check must settle against, and the decision builder is private. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/turn.py | 80 ++++------------- loopx/cli_commands/turn_decision.py | 120 +++++++++++++++++++------- tests/test_loopx_turn_managed_step.py | 83 ++++++++++++++---- 3 files changed, 174 insertions(+), 109 deletions(-) diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index bf36f1e3c2..29adcdb506 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -18,7 +18,6 @@ run_configured_turn_outcome_ingest_fail_open, ) from ..capabilities.periodic_report.cadence_runtime import extend_cadence_turn_start_dispatch -from ..capabilities.periodic_report.pending_intent import periodic_report_pending_intent_interaction_hook from ..control_plane.quota.live_decision import build_live_quota_should_run_decision from ..control_plane.agents.workspace_guard import capture_delivery_workspace from ..control_plane.quota.heartbeat_receipt import ( @@ -40,9 +39,6 @@ read_persisted_todo_record, read_persisted_todo_record_with_source, ) -from ..control_plane.scheduler.execution_context import ( - scheduler_execution_context_for_turn, -) from ..control_plane.turn_driver import ( LOOPX_TURN_EXECUTION_SCHEMA_VERSION, TurnRecoveryBlockedError, @@ -57,13 +53,12 @@ from ..control_plane.turn_driver.host_binding import managed_executor_binding from ..quota import spend_quota_slot from ..state_refresh import refresh_state_run -from ..status import AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, collect_status from ..todos import resolve_todo_state_path -from .lark_inbox import ( - build_lark_operator_inbox_urgency_projector, - dispatch_goal_lark_turn_start_hooks, +from .lark_inbox import dispatch_goal_lark_turn_start_hooks +from .turn_decision import ( + build_fresh_turn_decision_owner, + collect_turn_status_payload, ) -from .turn_decision import apply_controller_advisory_primary from .turn_dsh_host import build_dsh_host_runner from .turn_registration import register_turn_commands as register_turn_commands from .turn_inspection import handle_turn_journal_inspection @@ -116,9 +111,6 @@ def handle_turn_command( output_format=output_format, print_payload=print_payload, ) try: - scan_roots = [Path(item).expanduser() for item in args.scan_path] - if not scan_roots: - scan_roots = [Path(args.scan_root).expanduser()] runtime_root = resolve_status_projection_cache_runtime_root( registry_path=registry_path, runtime_root_override=runtime_root_arg, @@ -148,51 +140,20 @@ def handle_turn_command( agent_id=args.agent_id, available=args.available_capabilities, ) - operator_inbox_urgency_projector = build_lark_operator_inbox_urgency_projector( - runtime_root_arg=runtime_root, - ) - status_payload = collect_status( + # `run-once` and `managed-step` must resolve the same governing decision + # from the same live status, scheduler context and capability hooks, so + # this Turn takes all of them -- and its later settle-against inputs -- + # from the shared decision owner instead of rebuilding them per command. + decision_owner = build_fresh_turn_decision_owner( + args, registry_path=registry_path, - runtime_root_override=runtime_root_arg, - scan_roots=scan_roots, - limit=max(max(0, args.limit), AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK), - goal_id=args.goal_id, - available_capabilities=args.available_capabilities, - ) - scheduler_context = scheduler_execution_context_for_turn( - host=args.host, - execution_mode=args.execution_mode, - scheduler_owner=args.scheduler_owner, + runtime_root=runtime_root, + runtime_root_arg=runtime_root_arg, + turn_start_hook_dispatch=turn_start_hook_dispatch, ) - def build_turn_decision( - *, requested_action_todo_id: str | None = None - ) -> dict[str, Any]: - return build_live_quota_should_run_decision( - status_payload, - goal_id=args.goal_id, - agent_id=args.agent_id, - available_capabilities=args.available_capabilities, - include_scheduler_detail=False, - codex_app_current_rrule=None, - registry_path=registry_path, - runtime_root=runtime_root, - route_source="loopx_turn_plan", - scheduler_execution_context=scheduler_context, - operator_inbox_urgency_projector=operator_inbox_urgency_projector, - bounded_research_frontier_projector=( - project_live_explore_composition_frontier - ), - requested_action_todo_id=requested_action_todo_id, - turn_start_hook_dispatch=turn_start_hook_dispatch, - interaction_projection_hooks=(periodic_report_pending_intent_interaction_hook( - registry_path=registry_path, runtime_root=runtime_root, - goal_id=args.goal_id, agent_id=args.agent_id),), - ) - - # `run-once` and `managed-step` must resolve the same governing decision, - # so the advisory-primary rebinding lives in the shared decision owner - # instead of being repeated per subcommand. - decision = apply_controller_advisory_primary(build_turn_decision) + operator_inbox_urgency_projector = decision_owner.operator_inbox_urgency_projector + scheduler_context = decision_owner.scheduler_execution_context + decision = decision_owner.resolve() resume_requested, session_binding = resolve_turn_resume_session_binding(args) turn_envelope = build_turn_envelope( decision, @@ -684,13 +645,10 @@ def terminal_closeout( return todo_completion(result, effect_ref=effect_ref) def current_status() -> dict[str, object]: - return collect_status( + return collect_turn_status_payload( + args, registry_path=registry_path, - runtime_root_override=runtime_root_arg, - scan_roots=scan_roots, - limit=max(max(0, args.limit), AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK), - goal_id=args.goal_id, - available_capabilities=args.available_capabilities, + runtime_root_arg=runtime_root_arg, ) def spend(*, effect_ref: str) -> dict[str, object]: diff --git a/loopx/cli_commands/turn_decision.py b/loopx/cli_commands/turn_decision.py index 8ad90e24d1..7ff89f6be8 100644 --- a/loopx/cli_commands/turn_decision.py +++ b/loopx/cli_commands/turn_decision.py @@ -5,6 +5,12 @@ primary, and sign the same ``loopx_turn_envelope_v0``. Duplicating that chain would let the two drift, so both owners build through this module. +The owner carries the inputs the decision was derived from as well, because a +Turn that later spends quota or re-checks the scheduler must use the same live +status, the same scheduler context and the same activation-bound capability +hooks as the decision it already took. A call site that re-derives one of those +inputs is free to drift from the decision it is settling. + Nothing here executes, writes, or spends: it only projects the control plane's current decision into the envelope the loop controller consumes. """ @@ -13,6 +19,7 @@ import argparse from collections.abc import Callable, Mapping +from dataclasses import dataclass from pathlib import Path from typing import Any @@ -56,35 +63,25 @@ def collect_turn_status_payload( ) -def build_turn_decision_builder( +def _build_turn_decision( args: argparse.Namespace, *, registry_path: Path, runtime_root: Path, - runtime_root_arg: str | None, status_payload: Mapping[str, Any], + scheduler_execution_context: Mapping[str, Any], + operator_inbox_urgency_projector: Callable[..., dict[str, Any]], turn_start_hook_dispatch: Mapping[str, Any] | None = None, ) -> Callable[..., dict[str, Any]]: - """Return the shared ``build_turn_decision`` used by the Turn owners. + """Return the ``build_turn_decision`` every Turn owner resolves through. ``turn_start_hook_dispatch`` is the caller's business: an executing Turn may publish Go/No-Go hooks before deciding, while a read-only managed step must not. Passing the projection in keeps that choice with the caller. + Every other input is the owner's, so the two commands cannot disagree about + the status, the scheduler context or the capability hooks behind a decision. """ - scheduler_context = scheduler_execution_context_for_turn( - host=args.host, - execution_mode=args.execution_mode, - scheduler_owner=args.scheduler_owner, - ) - # Use the resolved runtime root, not the raw CLI argument. When a registry - # declares `common_runtime_root` and the command omits `--runtime-root`, - # the raw value is None and the activation check would silently read the - # global default instead of this registry's own extension state. - operator_inbox_urgency_projector = build_lark_operator_inbox_urgency_projector( - runtime_root_arg=runtime_root, - ) - def build_turn_decision( *, requested_action_todo_id: str | None = None ) -> dict[str, Any]: @@ -98,7 +95,7 @@ def build_turn_decision( registry_path=registry_path, runtime_root=runtime_root, route_source=TURN_DECISION_ROUTE_SOURCE, - scheduler_execution_context=scheduler_context, + scheduler_execution_context=scheduler_execution_context, operator_inbox_urgency_projector=operator_inbox_urgency_projector, bounded_research_frontier_projector=( project_live_explore_composition_frontier @@ -141,6 +138,75 @@ def apply_controller_advisory_primary( return decision +@dataclass(frozen=True) +class FreshTurnDecisionOwner: + """The one owner of a fresh Turn decision and the inputs behind it. + + A Turn that already decided must settle against what it decided with: the + same live status, the same scheduler context and the same + activation-bound operator-inbox projector. Carrying them here is what keeps + a later spend or scheduler re-check on the decision's own inputs instead of + a second reading of the same arguments. + """ + + status_payload: Mapping[str, Any] + scheduler_execution_context: Mapping[str, Any] + operator_inbox_urgency_projector: Callable[..., dict[str, Any]] + build_turn_decision: Callable[..., dict[str, Any]] + + def resolve(self) -> dict[str, Any]: + """The current governing decision, controller advisory primary applied.""" + + return apply_controller_advisory_primary(self.build_turn_decision) + + +def build_fresh_turn_decision_owner( + args: argparse.Namespace, + *, + registry_path: Path, + runtime_root: Path, + runtime_root_arg: str | None, + turn_start_hook_dispatch: Mapping[str, Any] | None = None, +) -> FreshTurnDecisionOwner: + """Read the live status and derive the shared decision inputs from it. + + ``runtime_root_arg`` is only the override the status read forwards, as + ``collect_status`` documents. The root every activation check uses is the + resolved ``runtime_root``: when a registry declares + ``common_runtime_root`` and the command omits ``--runtime-root``, the raw + argument is None and the check would otherwise read the operator's global + extension state instead of this registry's own. + """ + + status_payload = collect_turn_status_payload( + args, + registry_path=registry_path, + runtime_root_arg=runtime_root_arg, + ) + scheduler_execution_context = scheduler_execution_context_for_turn( + host=args.host, + execution_mode=args.execution_mode, + scheduler_owner=args.scheduler_owner, + ) + operator_inbox_urgency_projector = build_lark_operator_inbox_urgency_projector( + runtime_root_arg=runtime_root, + ) + return FreshTurnDecisionOwner( + status_payload=status_payload, + scheduler_execution_context=scheduler_execution_context, + operator_inbox_urgency_projector=operator_inbox_urgency_projector, + build_turn_decision=_build_turn_decision( + args, + registry_path=registry_path, + runtime_root=runtime_root, + status_payload=status_payload, + scheduler_execution_context=scheduler_execution_context, + operator_inbox_urgency_projector=operator_inbox_urgency_projector, + turn_start_hook_dispatch=turn_start_hook_dispatch, + ), + ) + + def fresh_turn_envelope( decision: Mapping[str, Any], *, @@ -173,34 +239,24 @@ def build_fresh_envelope_for_managed_step( may wake the same failed Turn again. """ - status_payload = collect_turn_status_payload( - args, - registry_path=registry_path, - runtime_root_arg=runtime_root_arg, - ) - build_turn_decision = build_turn_decision_builder( + owner = build_fresh_turn_decision_owner( args, registry_path=registry_path, runtime_root=runtime_root, runtime_root_arg=runtime_root_arg, - status_payload=status_payload, ) - decision = apply_controller_advisory_primary(build_turn_decision) return fresh_turn_envelope( - decision, - scheduler_execution_context=scheduler_execution_context_for_turn( - host=args.host, - execution_mode=args.execution_mode, - scheduler_owner=args.scheduler_owner, - ), + owner.resolve(), + scheduler_execution_context=owner.scheduler_execution_context, ) __all__ = [ + "FreshTurnDecisionOwner", "TURN_DECISION_ROUTE_SOURCE", "apply_controller_advisory_primary", "build_fresh_envelope_for_managed_step", - "build_turn_decision_builder", + "build_fresh_turn_decision_owner", "collect_turn_status_payload", "fresh_turn_envelope", ] diff --git a/tests/test_loopx_turn_managed_step.py b/tests/test_loopx_turn_managed_step.py index 895fd8b92e..be9925b7c5 100644 --- a/tests/test_loopx_turn_managed_step.py +++ b/tests/test_loopx_turn_managed_step.py @@ -396,12 +396,14 @@ def test_journal_turn_key_mismatch_is_refused() -> None: _decide(journal) -def test_shared_decision_builder_reads_the_resolved_runtime_root(monkeypatch) -> None: - """The activation check must read this registry's root, not the global default. +def test_shared_decision_owner_binds_both_turn_owners_to_one_set_of_inputs(monkeypatch) -> None: + """``run-once`` and ``managed-step`` must resolve through the shared owner. A registry that declares ``common_runtime_root`` while the command omits - ``--runtime-root`` passes ``None`` as the raw argument. Wiring the check to - that raw value would silently read the operator's global extension state. + ``--runtime-root`` passes ``None`` as the raw argument, so the activation + check must receive the resolved root instead. Deriving the status, the + scheduler context or the capability hooks per command would additionally + let the two owners disagree about the current governing decision. """ import argparse @@ -409,14 +411,35 @@ def test_shared_decision_builder_reads_the_resolved_runtime_root(monkeypatch) -> from loopx.cli_commands import turn_decision - seen: list[object] = [] + roots: list[object] = [] + captured: list[dict[str, Any]] = [] + status_payload = { + "ok": True, + "attention_queue": {"items": []}, + "run_history": {"goals": []}, + } + + def _projector(**_kwargs: object) -> dict[str, object]: + return {"schema_version": "lark_event_inbox_urgency_v0"} + + def _record_projector(*, runtime_root_arg): + roots.append(runtime_root_arg) + return _projector - def _record(*, runtime_root_arg): - seen.append(runtime_root_arg) - return lambda **_: {"schema_version": "lark_event_inbox_urgency_v0"} + def _record_decision(payload, **kwargs): + captured.append({"status_payload": payload, **kwargs}) + return {"selected_todo": None} monkeypatch.setattr( - turn_decision, "build_lark_operator_inbox_urgency_projector", _record + turn_decision, "build_lark_operator_inbox_urgency_projector", _record_projector + ) + monkeypatch.setattr( + turn_decision, "build_live_quota_should_run_decision", _record_decision + ) + monkeypatch.setattr( + turn_decision, + "collect_turn_status_payload", + lambda *_args, **_kwargs: status_payload, ) args = argparse.Namespace( goal_id=GOAL_ID, @@ -427,16 +450,44 @@ def _record(*, runtime_root_arg): available_capabilities=[], ) resolved_root = Path("/tmp/registry-scoped-runtime-root") - - turn_decision.build_turn_decision_builder( + registry_path = Path("/tmp/registry.json") + owner = turn_decision.build_fresh_turn_decision_owner( args, - registry_path=Path("/tmp/registry.json"), + registry_path=registry_path, runtime_root=resolved_root, runtime_root_arg=None, - status_payload={"ok": True, "attention_queue": {"items": []}, "run_history": {"goals": []}}, ) - assert seen == [resolved_root], ( - "the activation check must receive the resolved runtime root, not the raw " - f"argument that is None when the registry declares common_runtime_root: {seen}" + assert owner.resolve() == {"selected_todo": None} + assert captured[-1]["status_payload"] is status_payload + assert ( + captured[-1]["scheduler_execution_context"] is owner.scheduler_execution_context + ) + assert ( + captured[-1]["operator_inbox_urgency_projector"] + is owner.operator_inbox_urgency_projector + ) + assert captured[-1]["route_source"] == turn_decision.TURN_DECISION_ROUTE_SOURCE + + envelopes: list[object] = [] + + def _record_envelope(decision, *, scheduler_execution_context): + envelopes.append((decision, scheduler_execution_context)) + return {"schema_version": "loopx_turn_envelope_v0"} + + monkeypatch.setattr(turn_decision, "fresh_turn_envelope", _record_envelope) + turn_decision.build_fresh_envelope_for_managed_step( + args, + registry_path=registry_path, + runtime_root=resolved_root, + runtime_root_arg=None, + ) + assert captured[-1]["status_payload"] is status_payload + assert len(envelopes) == 1 + signed_decision, signed_context = envelopes[0] + assert signed_decision == {"selected_todo": None} + assert signed_context == captured[-1]["scheduler_execution_context"] + assert roots == [resolved_root, resolved_root], ( + "both Turn owners must receive the resolved runtime root, not the raw " + f"argument that is None when the registry declares common_runtime_root: {roots}" )