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
80 changes: 19 additions & 61 deletions loopx/cli_commands/turn.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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]:
Expand Down
120 changes: 88 additions & 32 deletions loopx/cli_commands/turn_decision.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
"""
Expand All @@ -13,6 +19,7 @@

import argparse
from collections.abc import Callable, Mapping
from dataclasses import dataclass
from pathlib import Path
from typing import Any

Expand Down Expand Up @@ -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]:
Expand All @@ -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
Expand Down Expand Up @@ -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],
*,
Expand Down Expand Up @@ -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",
]
Loading
Loading