diff --git a/examples/control_plane/status-collection-readmodel-smoke.py b/examples/control_plane/status-collection-readmodel-smoke.py index b8051b5b69..713990be71 100644 --- a/examples/control_plane/status-collection-readmodel-smoke.py +++ b/examples/control_plane/status-collection-readmodel-smoke.py @@ -7,6 +7,7 @@ from pathlib import Path import sys import tempfile +from types import SimpleNamespace from typing import Any @@ -117,6 +118,7 @@ def assert_wrapper_parity(registry_path: Path, runtime_root: Path, scan_root: Pa def assert_context_orchestration() -> None: calls: list[tuple[str, dict[str, Any]]] = [] runtime_root = Path("/tmp/status-collection-runtime") + history_audit = object() def record(name: str, value: Any) -> Any: calls.append((name, value if isinstance(value, dict) else {"value": value})) @@ -136,12 +138,12 @@ def resolve_runtime_root( assert registry_path == Path("registry.json"), registry_path return runtime_root - def collect_history(**kwargs: Any) -> dict[str, Any]: + def collect_status_history(**kwargs: Any) -> SimpleNamespace: assert kwargs["limit"] == 20, kwargs assert kwargs["goal_id"] == GOAL_ID, kwargs - assert kwargs["include_runtime_goals"] is True, kwargs - return record( - "collect_history", + assert kwargs["status_include_runtime_goals"] is True, kwargs + history = record( + "collect_status_history", { "goal_count": 1, "run_count": 0, @@ -149,10 +151,15 @@ def collect_history(**kwargs: Any) -> dict[str, Any]: "activation_state_filter": kwargs.get("activation_state_filter"), }, ) + return SimpleNamespace( + status_history=history, + contract_audit=history_audit, + ) def check_contract(**kwargs: Any) -> dict[str, Any]: assert kwargs["limit"] == 2, kwargs assert kwargs["goal_id_filter"] == GOAL_ID, kwargs + assert kwargs["history_audit"] is history_audit, kwargs return record( "check_contract", { @@ -186,7 +193,7 @@ def build_attention_queue(**kwargs: Any) -> dict[str, Any]: "collect_global_registry_health", {"ok": True, "current_registry_is_global": True}, ), - collect_history=collect_history, + collect_status_history=collect_status_history, check_contract=check_contract, build_attention_queue=build_attention_queue, build_runtime_summaries=lambda **kwargs: record( @@ -229,7 +236,9 @@ def build_attention_queue(**kwargs: Any) -> dict[str, Any]: context=context, ) - history_call = next(call for call in calls if call[0] == "collect_history") + history_call = next( + call for call in calls if call[0] == "collect_status_history" + ) assert history_call[1]["activation_state_filter"] == "active", history_call contract_call = next(call for call in calls if call[0] == "check_contract") assert contract_call[1]["activation_state_filter"] == "active", contract_call @@ -257,7 +266,7 @@ def build_attention_queue(**kwargs: Any) -> dict[str, Any]: assert [name for name, _ in calls][:4] == [ "load_registry", "collect_global_registry_health", - "collect_history", + "collect_status_history", "check_contract", ], calls diff --git a/loopx/contract.py b/loopx/contract.py index ae1db7a0f8..aa855fdc2d 100644 --- a/loopx/contract.py +++ b/loopx/contract.py @@ -23,7 +23,12 @@ index_identity, ) from .control_plane.todos.active_state_editing import COMPLETED_WORK_ARCHIVE_HEADING -from .history import collect_history, load_registry +from .history import ( + RunHistoryAudit, + build_run_history_audit, + collect_history, + load_registry, +) from .paths import DEFAULT_RUNTIME_ROOT, rel_or_abs, resolve_runtime_root from .registry import inspect_registry, inspect_registry_boundary, registry_goals, resolve_state_file from .state_projection import state_projection_gap_warning @@ -877,6 +882,7 @@ def check_contract( goal_id_filter: str | None = None, activation_state_filter: GoalActivationState | str | None = None, include_public_boundary_scan: bool = True, + history_audit: RunHistoryAudit | None = None, ) -> dict[str, Any]: error_diagnostics: list[dict[str, Any]] = [] warnings: list[str] = [] @@ -953,34 +959,70 @@ def add_global_error(code: str, message: str) -> None: else: warnings.append(f"runtime root does not exist yet: {runtime_root}") - history = collect_history( + if history_audit is None: + history = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id_filter, + limit=limit, + activation_state_filter=activation_state_filter, + ) + history_audit = build_run_history_audit( + history, + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id_filter, + activation_state_filter=activation_state_filter, + include_runtime_goals=True, + ) + elif not history_audit.matches( registry_path=registry_path, runtime_root=runtime_root, goal_id=goal_id_filter, - limit=limit, activation_state_filter=activation_state_filter, + include_runtime_goals=True, + ): + raise ValueError("history audit scope does not match contract request") + checks.append( + f"run-history goals={history_audit.goal_count} runs={history_audit.run_count}" ) - checks.append(f"run-history goals={history.get('goal_count')} runs={history.get('run_count')}") - for item in history.get("goals") or []: - raw = int(item.get("raw_index_records") or 0) - unique = int(item.get("unique_runs") or 0) - if item.get("legacy_runtime_goal") and raw > unique: - checks.append(f"{item.get('id')}: legacy runtime goal has duplicate rows raw={raw} unique={unique}") + for item in history_audit.goals: + raw = item.raw_index_records + unique = item.unique_runs + if item.legacy_runtime_goal and raw > unique: + checks.append( + f"{item.goal_id}: legacy runtime goal has duplicate rows " + f"raw={raw} unique={unique}" + ) continue if raw > unique: - duplicate_summary = _index_duplicate_summary(Path(str(item.get("index_path") or ""))) + duplicate_summary = _index_duplicate_summary(item.index_path) if duplicate_summary.get("unexpected_duplicate_rows"): - warnings.append(_index_duplicate_warning(item.get("id"), raw, unique, duplicate_summary)) + warnings.append( + _index_duplicate_warning( + item.goal_id, + raw, + unique, + duplicate_summary, + ) + ) else: emitted_check = False if duplicate_summary.get("reward_overlay_rows"): emitted_check = True checks.append( - f"{item.get('id')}: reward overlay rows raw={raw} unique={unique} " + f"{item.goal_id}: reward overlay rows raw={raw} unique={unique} " f"overlays={duplicate_summary.get('reward_overlay_rows')}" ) if not emitted_check: - warnings.append(_index_duplicate_warning(item.get("id"), raw, unique, duplicate_summary)) + warnings.append( + _index_duplicate_warning( + item.goal_id, + raw, + unique, + duplicate_summary, + ) + ) if include_public_boundary_scan: boundary = scan_public_boundary(scan_roots, registry=registry) diff --git a/loopx/control_plane/status/collection.py b/loopx/control_plane/status/collection.py index 9be3e93fba..c07caea5e4 100644 --- a/loopx/control_plane/status/collection.py +++ b/loopx/control_plane/status/collection.py @@ -49,7 +49,7 @@ class StatusCollectionContext: load_registry: StatusCallback resolve_runtime_root: StatusCallback collect_global_registry_health: StatusCallback - collect_history: StatusCallback + collect_status_history: StatusCallback check_contract: StatusCallback build_attention_queue: StatusCallback build_runtime_summaries: StatusCallback @@ -102,15 +102,16 @@ def collect_status( current_registry=registry, ) include_runtime_goals = bool(global_registry.get("current_registry_is_global")) - history = context.collect_history( + history_collection = context.collect_status_history( registry_path=registry_path, runtime_root=runtime_root, goal_id=goal_filter, limit=control_plane_limit, - include_runtime_goals=include_runtime_goals, + status_include_runtime_goals=include_runtime_goals, activation_state_filter=activation_filter, agent_lane_id=agent_lane_id, ) + history = history_collection.status_history contract = context.check_contract( registry_path=registry_path, runtime_root_override=str(runtime_root), @@ -119,6 +120,7 @@ def collect_status( goal_id_filter=goal_filter, include_public_boundary_scan=include_public_boundary_scan, activation_state_filter=activation_filter, + history_audit=history_collection.contract_audit, ) contract = project_contract_health_for_goal(contract, goal_id=goal_filter) queue = context.build_attention_queue( diff --git a/loopx/history.py b/loopx/history.py index d9dd908604..f1cf0d962e 100644 --- a/loopx/history.py +++ b/loopx/history.py @@ -3,6 +3,7 @@ import json from contextlib import nullcontext from collections.abc import Callable +from dataclasses import dataclass from datetime import datetime, timezone from heapq import merge from itertools import islice @@ -71,6 +72,55 @@ ) +@dataclass(frozen=True, slots=True) +class RunIndexAudit: + goal_id: str + index_path: Path + raw_index_records: int + unique_runs: int + legacy_runtime_goal: bool + + +@dataclass(frozen=True, slots=True) +class RunHistoryAudit: + registry_path: Path + runtime_root: Path + goal_id: str | None + activation_state_filter: str | None + include_runtime_goals: bool + goal_count: int + run_count: int + goals: tuple[RunIndexAudit, ...] + + def matches( + self, + *, + registry_path: Path, + runtime_root: Path, + goal_id: str | None, + activation_state_filter: GoalActivationState | str | None, + include_runtime_goals: bool, + ) -> bool: + normalized_activation = ( + normalize_goal_activation_state(activation_state_filter).value + if activation_state_filter is not None + else None + ) + return ( + self.registry_path == registry_path.expanduser().resolve() + and self.runtime_root == runtime_root.expanduser().resolve() + and self.goal_id == (str(goal_id or "").strip() or None) + and self.activation_state_filter == normalized_activation + and self.include_runtime_goals is include_runtime_goals + ) + + +@dataclass(frozen=True, slots=True) +class StatusHistoryCollection: + status_history: dict[str, Any] + contract_audit: RunHistoryAudit + + def now_local() -> str: return now_local_iso() @@ -386,6 +436,116 @@ def collect_history( } +def build_run_history_audit( + history: dict[str, Any], + *, + registry_path: Path, + runtime_root: Path, + goal_id: str | None, + activation_state_filter: GoalActivationState | str | None, + include_runtime_goals: bool, +) -> RunHistoryAudit: + normalized_activation = ( + normalize_goal_activation_state(activation_state_filter).value + if activation_state_filter is not None + else None + ) + return RunHistoryAudit( + registry_path=registry_path.expanduser().resolve(), + runtime_root=runtime_root.expanduser().resolve(), + goal_id=str(goal_id or "").strip() or None, + activation_state_filter=normalized_activation, + include_runtime_goals=include_runtime_goals, + goal_count=int(history.get("goal_count") or 0), + run_count=int(history.get("run_count") or 0), + goals=tuple( + RunIndexAudit( + goal_id=str(item.get("id") or ""), + index_path=Path(str(item.get("index_path") or "")), + raw_index_records=int(item.get("raw_index_records") or 0), + unique_runs=int(item.get("unique_runs") or 0), + legacy_runtime_goal=bool(item.get("legacy_runtime_goal")), + ) + for item in history.get("goals") or [] + if isinstance(item, dict) + ), + ) + + +def history_for_registry_members( + history: dict[str, Any], + *, + limit: int, +) -> dict[str, Any]: + goals = [ + goal + for goal in history.get("goals") or [] + if isinstance(goal, dict) and goal.get("registry_member") is True + ] + recent_limit = max(0, limit) + recent_runs = list( + islice( + merge( + *( + goal.get("latest_runs", [])[:recent_limit] + for goal in goals + if isinstance(goal.get("latest_runs"), list) + ), + key=lambda item: _chronology_key(item.get("generated_at")), + reverse=True, + ), + recent_limit, + ) + ) + return { + **history, + "goal_count": len(goals), + "run_count": sum(int(goal.get("unique_runs") or 0) for goal in goals), + "goals": goals, + "runs": recent_runs, + } + + +def collect_status_history( + *, + registry_path: Path, + runtime_root: Path, + goal_id: str | None, + limit: int, + status_include_runtime_goals: bool, + activation_state_filter: GoalActivationState | str | None = None, + agent_lane_id: str | None = None, +) -> StatusHistoryCollection: + history = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + limit=limit, + include_runtime_goals=True, + activation_state_filter=activation_state_filter, + agent_lane_id=agent_lane_id, + ) + audit = build_run_history_audit( + history, + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + activation_state_filter=activation_state_filter, + include_runtime_goals=True, + ) + status_history = history + if ( + not status_include_runtime_goals + and not str(goal_id or "").strip() + and activation_state_filter is None + ): + status_history = history_for_registry_members(history, limit=limit) + return StatusHistoryCollection( + status_history=status_history, + contract_audit=audit, + ) + + def inspect_index_duplicates( *, registry_path: Path, diff --git a/loopx/status.py b/loopx/status.py index f7c12377bb..c42ba32971 100644 --- a/loopx/status.py +++ b/loopx/status.py @@ -49,7 +49,7 @@ build_goal_channel_notification_projection, ) from .handoff_budget import handoff_budget_contract -from .history import collect_history, load_registry +from .history import collect_status_history, load_registry from .history import STATUS_NEUTRAL_CLASSIFICATIONS as HISTORY_STATUS_NEUTRAL_CLASSIFICATIONS from .interface_budget import interface_budget_cadence_for_runs from .long_task_cadence import build_long_task_cadence_hint @@ -1220,7 +1220,7 @@ def build_status_collection_context() -> StatusCollectionContext: load_registry=load_registry, resolve_runtime_root=resolve_runtime_root, collect_global_registry_health=collect_global_registry_health, - collect_history=collect_history, + collect_status_history=collect_status_history, check_contract=check_contract, build_attention_queue=build_attention_queue, build_runtime_summaries=build_status_runtime_summaries, diff --git a/tests/control_plane/test_status_collection_material_capability_wiring.py b/tests/control_plane/test_status_collection_material_capability_wiring.py index 636d0bbde7..76d3f597c9 100644 --- a/tests/control_plane/test_status_collection_material_capability_wiring.py +++ b/tests/control_plane/test_status_collection_material_capability_wiring.py @@ -1,6 +1,7 @@ from __future__ import annotations from pathlib import Path +from types import SimpleNamespace from typing import Any import pytest @@ -36,11 +37,14 @@ def _context( "ok": True, "current_registry_is_global": True, }, - collect_history=lambda **_kwargs: { - "goal_count": 1, - "run_count": 0, - "goals": [], - }, + collect_status_history=lambda **_kwargs: SimpleNamespace( + status_history={ + "goal_count": 1, + "run_count": 0, + "goals": [], + }, + contract_audit=object(), + ), check_contract=lambda **_kwargs: { "ok": True, "summary": "ok", diff --git a/tests/control_plane/test_status_history_reuse.py b/tests/control_plane/test_status_history_reuse.py new file mode 100644 index 0000000000..138ed01bfb --- /dev/null +++ b/tests/control_plane/test_status_history_reuse.py @@ -0,0 +1,261 @@ +from __future__ import annotations + +from dataclasses import replace +import json +from pathlib import Path +from typing import Any + +import pytest + +import loopx.history as history_module +from loopx.contract import check_contract +from loopx.history import collect_history, collect_status_history +from loopx.status import collect_status + + +REGISTERED_GOAL_ID = "registered-goal" +SECOND_REGISTERED_GOAL_ID = "registered-goal-b" +RUNTIME_ONLY_GOAL_ID = "runtime-only-goal" + + +def _write_run_index( + runtime_root: Path, + goal_id: str, + runs: list[dict[str, Any]], +) -> None: + runs_dir = runtime_root / "goals" / goal_id / "runs" + runs_dir.mkdir(parents=True) + artifact = runs_dir / "artifact.json" + artifact.write_text("{}\n", encoding="utf-8") + rows = [] + for run in runs: + rows.append( + json.dumps( + { + **run, + "goal_id": goal_id, + "json_path": str(artifact), + "markdown_path": str(artifact), + } + ) + ) + (runs_dir / "index.jsonl").write_text("\n".join(rows) + "\n", encoding="utf-8") + + +def _write_fixture(tmp_path: Path) -> tuple[Path, Path]: + runtime_root = tmp_path / "runtime" + project_root = tmp_path / "project" + project_root.mkdir(parents=True) + goals = [] + for goal_id in (REGISTERED_GOAL_ID, SECOND_REGISTERED_GOAL_ID): + state_path = project_root / f"{goal_id}.md" + state_path.write_text( + "---\nstatus: active\n---\n\n" + f"# {goal_id}\n\n" + "## Agent Todo\n\n" + "- [ ] Continue work.\n" + f" \n", + encoding="utf-8", + ) + goals.append( + { + "id": goal_id, + "domain": "status-history-reuse", + "status": "active", + "repo": str(project_root), + "state_file": state_path.name, + "adapter": { + "kind": "test_v0", + "status": "connected-read-only", + }, + "authority_sources": [], + } + ) + registry_path = project_root / ".loopx" / "registry.json" + registry_path.parent.mkdir() + registry_path.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(runtime_root), + "goals": goals, + } + ) + + "\n", + encoding="utf-8", + ) + _write_run_index( + runtime_root, + REGISTERED_GOAL_ID, + [ + { + "classification": "registered-tied", + "generated_at": "2026-09-17T00:00:00Z", + "agent_id": "agent-a", + }, + { + "classification": "registered-malformed", + "generated_at": "not-a-time", + "agent_id": "agent-a", + }, + ], + ) + _write_run_index( + runtime_root, + SECOND_REGISTERED_GOAL_ID, + [ + { + "classification": "second-registered-tied", + "generated_at": "2026-09-17T00:00:00+00:00", + }, + { + "classification": "second-registered-older", + "generated_at": "2026-09-16T00:00:00Z", + }, + ], + ) + _write_run_index( + runtime_root, + RUNTIME_ONLY_GOAL_ID, + [ + { + "classification": "runtime-newer", + "generated_at": "2026-09-18T00:00:00Z", + } + ], + ) + return registry_path, runtime_root + + +@pytest.mark.parametrize("limit", [0, 1, 2]) +def test_status_history_projection_matches_direct_registry_history( + tmp_path: Path, + limit: int, +) -> None: + registry_path, runtime_root = _write_fixture(tmp_path) + + snapshot = collect_status_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=None, + limit=limit, + status_include_runtime_goals=False, + agent_lane_id="agent-a", + ) + direct = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=None, + limit=limit, + include_runtime_goals=False, + agent_lane_id="agent-a", + ) + + assert snapshot.status_history == direct + assert snapshot.contract_audit.goal_count == 3 + assert snapshot.contract_audit.run_count == 5 + assert [goal.goal_id for goal in snapshot.contract_audit.goals] == [ + REGISTERED_GOAL_ID, + SECOND_REGISTERED_GOAL_ID, + RUNTIME_ONLY_GOAL_ID, + ] + + +@pytest.mark.parametrize( + ("status_include_runtime_goals", "goal_id", "activation_state_filter"), + [ + (True, None, None), + (False, REGISTERED_GOAL_ID, None), + (False, None, "active"), + ], +) +def test_status_history_keeps_existing_nonproject_scopes( + tmp_path: Path, + status_include_runtime_goals: bool, + goal_id: str | None, + activation_state_filter: str | None, +) -> None: + registry_path, runtime_root = _write_fixture(tmp_path) + + snapshot = collect_status_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + limit=1, + status_include_runtime_goals=status_include_runtime_goals, + activation_state_filter=activation_state_filter, + ) + direct = collect_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + limit=1, + include_runtime_goals=status_include_runtime_goals, + activation_state_filter=activation_state_filter, + ) + + assert snapshot.status_history == direct + + +def test_status_reuses_one_history_scan_without_hiding_runtime_only_audit( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry_path, runtime_root = _write_fixture(tmp_path) + loaded_goal_ids: list[str] = [] + original_load_index = history_module.load_index + + def record_load_index(path: Path, **kwargs: Any): + loaded_goal_ids.append(path.parts[-3]) + return original_load_index(path, **kwargs) + + monkeypatch.setattr(history_module, "load_index", record_load_index) + + payload = collect_status( + registry_path=registry_path, + runtime_root_override=str(runtime_root), + scan_roots=[tmp_path], + limit=2, + include_public_boundary_scan=False, + agent_lane_id="agent-a", + ) + + assert payload["goal_count"] == 2 + assert [goal["id"] for goal in payload["run_history"]["goals"]] == [ + REGISTERED_GOAL_ID, + SECOND_REGISTERED_GOAL_ID, + ] + assert "run-history goals=3 runs=5" in payload["contract"]["checks"] + assert loaded_goal_ids == [ + REGISTERED_GOAL_ID, + SECOND_REGISTERED_GOAL_ID, + RUNTIME_ONLY_GOAL_ID, + ] + + +def test_contract_rejects_a_history_audit_from_another_runtime( + tmp_path: Path, +) -> None: + registry_path, runtime_root = _write_fixture(tmp_path) + snapshot = collect_status_history( + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=None, + limit=2, + status_include_runtime_goals=False, + ) + + mismatched = replace( + snapshot.contract_audit, + runtime_root=tmp_path / "other-runtime", + ) + with pytest.raises(ValueError, match="history audit scope"): + check_contract( + registry_path=registry_path, + runtime_root_override=str(runtime_root), + scan_roots=[tmp_path], + limit=2, + include_public_boundary_scan=False, + history_audit=mismatched, + )