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
3 changes: 3 additions & 0 deletions loopx/cli_commands/quota_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -319,6 +319,7 @@ def prepare_quota_command_context(
goal_id=status_goal_id,
max_age_seconds=projection_cache_ttl_seconds,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
if status_payload is None:
collector = status_collector or collect_status
Expand All @@ -329,6 +330,7 @@ def prepare_quota_command_context(
limit=status_limit,
goal_id=status_goal_id,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
if bool(getattr(args, "write_projection_cache", False)):
cache_metadata = write_status_projection_cache(
Expand All @@ -341,6 +343,7 @@ def prepare_quota_command_context(
payload=status_payload,
max_age_seconds=projection_cache_ttl_seconds,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
elif isinstance(status_payload.get("projection_cache"), dict):
cache_metadata = dict(status_payload["projection_cache"])
Expand Down
4 changes: 4 additions & 0 deletions loopx/cli_commands/status.py
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,7 @@ def handle_status_command(
goal_id=args.goal_id,
max_age_seconds=args.projection_cache_ttl_seconds,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
if payload is None:
payload = collect_status(
Expand All @@ -211,6 +212,7 @@ def handle_status_command(
include_task_graph=args.include_task_graph,
goal_id=args.goal_id,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
if args.write_projection_cache:
cache_metadata = write_status_projection_cache(
Expand All @@ -223,6 +225,7 @@ def handle_status_command(
payload=payload,
max_age_seconds=args.projection_cache_ttl_seconds,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
payload["projection_cache"] = cache_metadata
elif cache_metadata:
Expand Down Expand Up @@ -867,6 +870,7 @@ def handle_review_packet_command(
include_task_graph=not args.handoff_only,
goal_id=args.goal_id,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)
if args.agent_id:
attach_agent_lane_next_actions(status_payload, agent_id=args.agent_id)
Expand Down
1 change: 1 addition & 0 deletions loopx/cli_commands/turn_decision.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ def collect_turn_status_payload(
limit=max(max(0, args.limit), AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK),
goal_id=args.goal_id,
available_capabilities=args.available_capabilities,
agent_lane_id=args.agent_id,
)


Expand Down
39 changes: 37 additions & 2 deletions loopx/control_plane/runtime/run_context_retention.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@
PROGRESS_DELIVERY_OUTCOMES,
)
from ..work_items.progress_result import PROGRESS_OBSERVATION_SCHEMA_VERSION, ProgressResultClass
from ..work_items.autonomous_replan_ack import autonomous_replan_ack_recorded
from ..work_items.autonomous_replan_ack import (
AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK,
autonomous_replan_ack_recorded,
)
from ..work_items.autonomous_replan_obligation import run_history_agent_id

GOAL_SEMANTIC_HISTORY_SCHEMA_VERSION = "goal_semantic_history_v0"
SEMANTIC_CONTEXT_RUN_FIELDS = (
Expand Down Expand Up @@ -261,12 +265,43 @@ def latest_runs_with_agent_context(
runs: list[dict[str, Any]],
*,
limit: int,
agent_lane_id: str | None = None,
) -> list[dict[str, Any]]:
"""Return a strict recent-run drill-down window.

Durable control semantics live in ``goal_semantic_history_v0`` instead of
expanding this display list beyond its advertised limit.

``agent_lane_id`` adds one lane-scoped decision window next to the
goal-wide display window. Agent-lane replan triggers count only that lane's
material runs since its own last replan ACK, so a Goal with several active
lanes can interleave more than ``limit`` peer records between two rows of a
single lane and hide that lane's own threshold. The requesting lane keeps
the same ``AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK`` budget a single-lane Goal
would have, and the goal-wide rows stay in place.
"""

bounded = max(0, limit)
return list(runs[:bounded])
goal_window = list(runs[:bounded])
if not agent_lane_id:
return goal_window
lane_present = False
lane_positions: list[int] = []
for position, run in enumerate(runs):
attributed = run_history_agent_id(run)
if attributed == agent_lane_id:
lane_present = True
elif attributed is not None:
continue
lane_positions.append(position)
if not lane_present:
return goal_window
return [
runs[position]
for position in sorted(
{
*range(min(bounded, len(runs))),
*lane_positions[:AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK],
}
)
]
9 changes: 9 additions & 0 deletions loopx/control_plane/runtime/status_projection_cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ def status_projection_cache_key(
include_task_graph: bool,
goal_id: str | None,
available_capabilities: Any = None,
agent_lane_id: str | None = None,
) -> str:
request = {
"schema_version": STATUS_PROJECTION_CACHE_SCHEMA_VERSION,
Expand All @@ -67,6 +68,7 @@ def status_projection_cache_key(
"limit": max(0, int(limit)),
"include_task_graph": bool(include_task_graph),
"goal_id": str(goal_id or "").strip() or None,
"agent_lane_id": str(agent_lane_id or "").strip() or None,
"available_capabilities": _normalized_available_capabilities(
available_capabilities
),
Expand Down Expand Up @@ -94,6 +96,7 @@ def status_projection_cache_metadata(
goal_id: str | None,
max_age_seconds: int,
available_capabilities: Any = None,
agent_lane_id: str | None = None,
) -> dict[str, Any]:
key = status_projection_cache_key(
registry_path=registry_path,
Expand All @@ -103,6 +106,7 @@ def status_projection_cache_metadata(
include_task_graph=include_task_graph,
goal_id=goal_id,
available_capabilities=available_capabilities,
agent_lane_id=agent_lane_id,
)
normalized_capabilities = _normalized_available_capabilities(
available_capabilities
Expand All @@ -113,6 +117,7 @@ def status_projection_cache_metadata(
"path": str(status_projection_cache_path(runtime_root, key)),
"max_age_seconds": max(0, int(max_age_seconds)),
"goal_id": str(goal_id or "").strip() or None,
"agent_lane_id": str(agent_lane_id or "").strip() or None,
"limit": max(0, int(limit)),
"include_task_graph": bool(include_task_graph),
"scan_roots": [str(path.expanduser()) for path in scan_roots],
Expand All @@ -130,6 +135,7 @@ def load_status_projection_cache(
goal_id: str | None,
max_age_seconds: int,
available_capabilities: Any = None,
agent_lane_id: str | None = None,
) -> tuple[dict[str, Any] | None, dict[str, Any]]:
metadata = status_projection_cache_metadata(
registry_path=registry_path,
Expand All @@ -140,6 +146,7 @@ def load_status_projection_cache(
goal_id=goal_id,
max_age_seconds=max_age_seconds,
available_capabilities=available_capabilities,
agent_lane_id=agent_lane_id,
)
path = Path(str(metadata["path"]))
metadata["hit"] = False
Expand Down Expand Up @@ -193,6 +200,7 @@ def write_status_projection_cache(
payload: dict[str, Any],
max_age_seconds: int,
available_capabilities: Any = None,
agent_lane_id: str | None = None,
) -> dict[str, Any]:
metadata = status_projection_cache_metadata(
registry_path=registry_path,
Expand All @@ -203,6 +211,7 @@ def write_status_projection_cache(
goal_id=goal_id,
max_age_seconds=max_age_seconds,
available_capabilities=available_capabilities,
agent_lane_id=agent_lane_id,
)
path = Path(str(metadata["path"]))
path.parent.mkdir(parents=True, exist_ok=True)
Expand Down
2 changes: 2 additions & 0 deletions loopx/control_plane/status/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ def collect_status(
recent_run_limit: int | None = None,
include_goal_subagent_configuration: bool = False,
activation_state_filter: GoalActivationState | str | None = None,
agent_lane_id: str | None = None,
) -> dict[str, Any]:
display_limit = max(0, limit)
control_plane_limit = max(
Expand Down Expand Up @@ -107,6 +108,7 @@ def collect_status(
limit=control_plane_limit,
include_runtime_goals=include_runtime_goals,
activation_state_filter=activation_filter,
agent_lane_id=agent_lane_id,
)
contract = context.check_contract(
registry_path=registry_path,
Expand Down
6 changes: 6 additions & 0 deletions loopx/control_plane/work_items/autonomous_replan_ack.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,12 @@
from .progress_observation import FRESH_VISION_PATH_DISPOSITIONS

AUTONOMOUS_REPLAN_ACK_MATERIAL_RUN_WINDOW = 20
# A normal delivery appends both a durable run record and a neutral quota-spend
# record. Keep enough internal history to observe the full material-run
# threshold even when those records are interleaved, with headroom for other
# neutral events. The budget is per agent lane: a Goal-wide window that peer
# lanes also fill cannot decide one lane's material-run count.
AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK = AUTONOMOUS_REPLAN_ACK_MATERIAL_RUN_WINDOW * 3


def ack_binds_trigger_checkpoints(ack: dict[str, Any] | None) -> bool:
Expand Down
7 changes: 6 additions & 1 deletion loopx/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,7 @@ def collect_history(
limit: int,
include_runtime_goals: bool = True,
activation_state_filter: GoalActivationState | str | None = None,
agent_lane_id: str | None = None,
) -> dict[str, Any]:
from .capabilities.machine_configuration.builtins import (
build_builtin_machine_configuration_registry,
Expand Down Expand Up @@ -360,7 +361,11 @@ def collect_history(
"raw_index_records": raw_count,
"unique_runs": len(runs),
"latest_status_run": latest_status_run(runs),
"latest_runs": latest_runs_with_agent_context(runs, limit=limit),
"latest_runs": latest_runs_with_agent_context(
runs,
limit=limit,
agent_lane_id=agent_lane_id,
),
"semantic_history": goal_semantic_history_from_runs(runs),
}
if registry_member:
Expand Down
8 changes: 4 additions & 4 deletions loopx/status.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@
)
from .control_plane.work_items.autonomous_replan_ack import (
AUTONOMOUS_REPLAN_ACK_MATERIAL_RUN_WINDOW,
AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK as _AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK,
compact_autonomous_replan_ack,
)
from .control_plane.work_items.autonomous_replan_obligation import (
Expand Down Expand Up @@ -325,10 +326,7 @@
AUTONOMOUS_REPLAN_STALL_THRESHOLD = _AUTONOMOUS_REPLAN_STALL_THRESHOLD_READ_MODEL
DEAD_MONITOR_REPEAT_THRESHOLD = 6
AUTONOMOUS_REPLAN_PERIODIC_RUN_THRESHOLD = AUTONOMOUS_REPLAN_ACK_MATERIAL_RUN_WINDOW
# A normal delivery appends both a durable run and a neutral quota-spend run.
# Keep enough internal history to observe the full material-run threshold even
# when those records are interleaved, with headroom for other neutral events.
AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK = AUTONOMOUS_REPLAN_PERIODIC_RUN_THRESHOLD * 3
AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK = _AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK
BACKLOG_HYGIENE_SECTION_HEADINGS = ("Next Action", "Operating Lessons")
BACKLOG_HYGIENE_BULLET_PATTERN = re.compile(r"^\s*(?:[-*]|\d+[.)])\s+(.+?)\s*$")
BACKLOG_HYGIENE_HINT_PATTERN = re.compile(
Expand Down Expand Up @@ -1240,6 +1238,7 @@ def collect_status(
recent_run_limit: int | None = None,
include_goal_subagent_configuration: bool = False,
activation_state_filter: str | None = None,
agent_lane_id: str | None = None,
) -> dict[str, Any]:
return _collect_status_read_model(
registry_path=registry_path,
Expand All @@ -1255,5 +1254,6 @@ def collect_status(
include_goal_subagent_configuration
),
activation_state_filter=activation_state_filter,
agent_lane_id=agent_lane_id,
context=build_status_collection_context(),
)
Loading
Loading