From 774e16bcca2cdfd73d85b3cfb5db248181986110 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 16 Sep 2026 22:00:14 +0800 Subject: [PATCH] fix(replan): decide a lane's periodic review on lane-complete run history An agent lane's autonomous-replan triggers count only that lane's material runs since its own last replan ACK, while the Goal-wide status window is shared with every peer lane. On a multi-lane Goal, peer volume can push the lane's own rows (and its last ACK) out of that window, so the guarded Turn selects a different replan obligation generation than the state-refresh writeback derives from the complete run index. The guarded writeback then defers silently, no ACK is ever recorded, and the lane can never discharge the obligation: every wake re-selects the same unsatisfiable replan. Keep the goal-wide display window and add the requesting lane's own AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK rows, so lane-scoped decisions see the same history the writeback derives from. The projection cache key now includes the requested lane. Validation: focused lookback/window regressions, refresh-state replan gate, quota settlement CLI, quota/status/history suites, and the live lane (codex-side-bypass) guard now selects the same obligation id the writeback derives. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/quota_context.py | 3 + loopx/cli_commands/status.py | 4 + loopx/cli_commands/turn_decision.py | 1 + .../runtime/run_context_retention.py | 39 +++- .../runtime/status_projection_cache.py | 9 + loopx/control_plane/status/collection.py | 2 + .../work_items/autonomous_replan_ack.py | 6 + loopx/history.py | 7 +- loopx/status.py | 8 +- ...est_autonomous_replan_periodic_lookback.py | 171 ++++++++++++++++++ 10 files changed, 243 insertions(+), 7 deletions(-) diff --git a/loopx/cli_commands/quota_context.py b/loopx/cli_commands/quota_context.py index 064853d714..db797371b1 100644 --- a/loopx/cli_commands/quota_context.py +++ b/loopx/cli_commands/quota_context.py @@ -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 @@ -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( @@ -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"]) diff --git a/loopx/cli_commands/status.py b/loopx/cli_commands/status.py index 8f710d173b..ce3018d0f5 100644 --- a/loopx/cli_commands/status.py +++ b/loopx/cli_commands/status.py @@ -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( @@ -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( @@ -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: @@ -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) diff --git a/loopx/cli_commands/turn_decision.py b/loopx/cli_commands/turn_decision.py index 7ff89f6be8..1932a9bf53 100644 --- a/loopx/cli_commands/turn_decision.py +++ b/loopx/cli_commands/turn_decision.py @@ -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, ) diff --git a/loopx/control_plane/runtime/run_context_retention.py b/loopx/control_plane/runtime/run_context_retention.py index 48164ca7c4..fae8667df6 100644 --- a/loopx/control_plane/runtime/run_context_retention.py +++ b/loopx/control_plane/runtime/run_context_retention.py @@ -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 = ( @@ -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], + } + ) + ] diff --git a/loopx/control_plane/runtime/status_projection_cache.py b/loopx/control_plane/runtime/status_projection_cache.py index b40114346b..9954a64366 100644 --- a/loopx/control_plane/runtime/status_projection_cache.py +++ b/loopx/control_plane/runtime/status_projection_cache.py @@ -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, @@ -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 ), @@ -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, @@ -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 @@ -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], @@ -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, @@ -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 @@ -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, @@ -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) diff --git a/loopx/control_plane/status/collection.py b/loopx/control_plane/status/collection.py index 76ba2e4849..680303c834 100644 --- a/loopx/control_plane/status/collection.py +++ b/loopx/control_plane/status/collection.py @@ -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( @@ -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, diff --git a/loopx/control_plane/work_items/autonomous_replan_ack.py b/loopx/control_plane/work_items/autonomous_replan_ack.py index a89be06bc2..777a34f3ba 100644 --- a/loopx/control_plane/work_items/autonomous_replan_ack.py +++ b/loopx/control_plane/work_items/autonomous_replan_ack.py @@ -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: diff --git a/loopx/history.py b/loopx/history.py index 3515b0bed5..d9dd908604 100644 --- a/loopx/history.py +++ b/loopx/history.py @@ -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, @@ -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: diff --git a/loopx/status.py b/loopx/status.py index d0a671d727..d31aaaa2d8 100644 --- a/loopx/status.py +++ b/loopx/status.py @@ -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 ( @@ -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( @@ -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, @@ -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(), ) diff --git a/tests/control_plane/test_autonomous_replan_periodic_lookback.py b/tests/control_plane/test_autonomous_replan_periodic_lookback.py index 696e2e9dcc..75e46298b6 100644 --- a/tests/control_plane/test_autonomous_replan_periodic_lookback.py +++ b/tests/control_plane/test_autonomous_replan_periodic_lookback.py @@ -1,5 +1,9 @@ from __future__ import annotations +import json +from pathlib import Path +from typing import Any + from loopx.cli_commands.status import ( _status_collection_limit_for_agent_lane, _trim_run_history_for_status_display, @@ -7,9 +11,14 @@ from loopx.control_plane.goals.goal_frontier.ack_policy import ( autonomous_replan_ack_satisfies_obligation, ) +from loopx.control_plane.runtime.run_context_retention import ( + latest_runs_with_agent_context, +) +from loopx.history import collect_history from loopx.status import ( AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, AUTONOMOUS_REPLAN_PERIODIC_RUN_THRESHOLD, + autonomous_replan_obligation_from_runs, autonomous_replan_periodic_review_from_runs, ) @@ -129,3 +138,165 @@ def test_agent_lane_keeps_periodic_control_history_off_the_display_path() -> Non "status display limit" ), } + + +PEER_AGENT_ID = "codex-peer-lane" +LANE_AGENT_ID = "codex-fixture" +MULTI_LANE_GOAL_ID = "multi-lane-goal" + + +def _peer_lane_runs(*, count: int) -> list[dict[str, Any]]: + """Return newest-first peer-lane rows that fill the goal-wide window.""" + + return [ + { + "classification": f"peer_delivery_{index:02d}", + "generated_at": f"2026-09-01T00:{index:02d}:00Z", + "agent_id": PEER_AGENT_ID, + } + for index in range(count) + ] + + +def _lane_material_runs(*, count: int) -> list[dict[str, Any]]: + """Return newest-first material rows for the requesting lane.""" + + return [ + { + "classification": f"lane_delivery_{index:02d}", + "generated_at": f"2026-08-31T23:{59 - index:02d}:00Z", + "agent_id": LANE_AGENT_ID, + } + for index in range(count) + ] + + +def _lane_replan_ack_run() -> dict[str, Any]: + return { + "classification": "state_refreshed", + "generated_at": "2026-08-31T23:00:00Z", + "agent_id": LANE_AGENT_ID, + "autonomous_replan_ack": { + "schema_version": "autonomous_replan_ack_v0", + "recorded": True, + "source": "refresh_state_semantic_delta", + "semantic_delta": { + "accepted": True, + "obligation_id": f"replan-{'0' * 16}", + "outcomes": ["new_surface"], + }, + }, + } + + +def _multi_lane_runs() -> list[dict[str, Any]]: + """Peer volume fills the goal-wide window ahead of one lane's own review.""" + + return [ + *_peer_lane_runs(count=AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK), + *_lane_material_runs(count=AUTONOMOUS_REPLAN_PERIODIC_RUN_THRESHOLD), + _lane_replan_ack_run(), + ] + + +def _write_multi_lane_history(tmp_path: Path, runs: list[dict[str, Any]]) -> tuple[Path, Path]: + registry_path = tmp_path / "registry.json" + registry_path.write_text("{}\n", encoding="utf-8") + runtime_root = tmp_path / "runtime" + runs_dir = runtime_root / "goals" / MULTI_LANE_GOAL_ID / "runs" + runs_dir.mkdir(parents=True) + rows = [] + for position, run in enumerate(runs): + row = dict(run) + row.setdefault("json_path", f"artifacts/{MULTI_LANE_GOAL_ID}-{position}.json") + row.setdefault("markdown_path", f"artifacts/{MULTI_LANE_GOAL_ID}-{position}.md") + rows.append(json.dumps(row)) + (runs_dir / "index.jsonl").write_text("\n".join(rows) + "\n", encoding="utf-8") + return registry_path, runtime_root + + +def test_peer_lane_volume_cannot_hide_one_lane_periodic_review(tmp_path: Path) -> None: + """A lane-scoped replan trigger must not depend on peer-lane volume. + + The guarded Turn selects the replan obligation for one agent lane, and the + writeback derives it again from the full run index. When the goal-wide + window cannot decide the lane's material-run count, the two derivations + disagree and the guarded writeback can never discharge the obligation. + """ + + runs = _multi_lane_runs() + registry_path, runtime_root = _write_multi_lane_history(tmp_path, runs) + + shared_window = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=MULTI_LANE_GOAL_ID, + limit=AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, + ) + goal_window_runs = shared_window["goals"][0]["latest_runs"] + assert not [ + run + for run in goal_window_runs + if str(run.get("agent_id") or "") == LANE_AGENT_ID + ] + assert ( + autonomous_replan_obligation_from_runs( + goal_window_runs, + agent_todos=None, + agent_id=LANE_AGENT_ID, + ) + is None + ) + + lane_window = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=MULTI_LANE_GOAL_ID, + limit=AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, + agent_lane_id=LANE_AGENT_ID, + ) + lane_obligation = autonomous_replan_obligation_from_runs( + lane_window["goals"][0]["latest_runs"], + agent_todos=None, + agent_id=LANE_AGENT_ID, + ) + complete_history_obligation = autonomous_replan_obligation_from_runs( + runs, + agent_todos=None, + agent_id=LANE_AGENT_ID, + ) + + assert lane_obligation is not None + assert complete_history_obligation is not None + assert lane_obligation["triggers"][0]["kind"] == "periodic_review_due" + assert ( + lane_obligation["obligation_id"] + == complete_history_obligation["obligation_id"] + ) + + +def test_lane_window_keeps_goal_rows_and_adds_the_lane_rows() -> None: + runs = _multi_lane_runs() + + lane_window = latest_runs_with_agent_context( + runs, + limit=AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, + agent_lane_id=LANE_AGENT_ID, + ) + + assert lane_window[: AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK] == runs[ + : AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK + ] + lane_rows = [ + run for run in lane_window if str(run.get("agent_id") or "") == LANE_AGENT_ID + ] + assert len(lane_rows) == AUTONOMOUS_REPLAN_PERIODIC_RUN_THRESHOLD + 1 + assert lane_window[-1]["generated_at"] == "2026-08-31T23:00:00Z" + + +def test_lane_window_leaves_other_lanes_untouched() -> None: + runs = _lane_material_runs(count=2) + + assert latest_runs_with_agent_context( + runs, limit=1, agent_lane_id="codex-absent-lane" + ) == runs[:1]