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
23 changes: 16 additions & 7 deletions examples/control_plane/status-collection-readmodel-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from pathlib import Path
import sys
import tempfile
from types import SimpleNamespace
from typing import Any


Expand Down Expand Up @@ -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}))
Expand All @@ -136,23 +138,28 @@ 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,
"goals": [],
"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",
{
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
68 changes: 55 additions & 13 deletions loopx/contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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] = []
Expand Down Expand Up @@ -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)
Expand Down
8 changes: 5 additions & 3 deletions loopx/control_plane/status/collection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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),
Expand All @@ -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(
Expand Down
160 changes: 160 additions & 0 deletions loopx/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()

Expand Down Expand Up @@ -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,
Expand Down
Loading