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
17 changes: 17 additions & 0 deletions docs/status-data-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -1833,6 +1833,23 @@ It mirrors the compact run index, but strips local artifact paths. UIs should
show artifact availability with `json_exists` and `markdown_exists` instead of
linking directly to local files.

Artifact flags are fresh observations for every returned Run, not facts trusted
from the persisted index. History collection first reads and deduplicates the
complete index and computes quota and semantic history; only then does it check
files for retained recent, latest-status, agent-lane and semantic-evidence rows.
`--limit 0` does not discard old vision, owner correction or active blocked-retry
evidence. Creating or deleting an artifact is visible on the next read even when
the index bytes are unchanged. Full `load_index`/`load_index_snapshot` readers
still observe every row by default. This is a caller filesystem optimization,
independent of the Todo authority provider, not an index cache or a limit on
decision history.

文件存在性是每次读取对返回记录的实时观测,不信任索引中保存的布尔值。先用完整索引
完成去重、额度和语义历史归约,再检查近期记录、最新状态、Agent 窗口及保留语义证据
的文件;即使 `--limit 0`,旧愿景、用户纠偏和有效阻塞重试也不会因此丢失。索引不变时,
文件创建/删除仍在下次读取生效;底层完整索引接口默认保留逐条检查行为。这不改变
File/SQLite 等 Todo provider 的权威,也不限制决策历史。

On the `status`, `quota should-run`, and `history` read paths, relative
`common_runtime_root` values, relative `--runtime-root` overrides, and relative
run-index artifact paths are resolved against the project root that owns the
Expand Down
19 changes: 18 additions & 1 deletion loopx/control_plane/runtime/run_context_retention.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
from __future__ import annotations

from collections.abc import Callable
from collections.abc import Callable, Iterator
from typing import Any

from ..quota.blocked_retry import active_turn_retry_for_run
Expand Down Expand Up @@ -232,6 +232,23 @@ def goal_semantic_history_from_runs(
return semantic_history


def iter_goal_semantic_history_runs(value: dict[str, Any]) -> Iterator[dict[str, Any]]:
"""Enumerate retained Run references without reselecting or limiting history.

The filesystem adapter uses this shape-owned traversal to observe artifacts
after semantic reduction, including evidence outside the recent-run window.
"""
for context in value.get("agents", []):
for field in SEMANTIC_CONTEXT_RUN_FIELDS:
run = context.get(field)
if isinstance(run, dict):
yield run
yield from value.get("active_blocked_retry_runs", [])
correction = value.get("latest_owner_correction_run")
if isinstance(correction, dict):
yield correction


def compact_goal_semantic_history(
value: Any,
*,
Expand Down
56 changes: 40 additions & 16 deletions loopx/history.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,10 @@
import hashlib
import json
from contextlib import nullcontext
from collections.abc import Callable
from collections.abc import Callable, Iterable
from dataclasses import dataclass
from heapq import merge
from itertools import islice
from itertools import chain, islice
from pathlib import Path
from typing import Any

Expand Down Expand Up @@ -35,6 +35,7 @@
)
from .control_plane.runtime.run_context_retention import (
goal_semantic_history_from_runs,
iter_goal_semantic_history_runs,
latest_runs_with_agent_context,
)
from .control_plane.runtime.run_index_duplicates import (
Expand Down Expand Up @@ -229,32 +230,40 @@ def _indexed_artifact_exists(value: Any, *, artifact_root: Path | None) -> bool:
return path.exists()


def _observe_run_artifacts(
records: Iterable[dict[str, Any]], *, artifact_root: Path | None,
) -> None:
"""Fresh request-local observations, shared by full and selected reads."""
observed: dict[str, bool] = {}
for record in records:
for source, target in (("json_path", "json_exists"), ("markdown_path", "markdown_exists")):
path = str(record.get(source) or "").strip()
if path not in observed:
observed[path] = _indexed_artifact_exists(path, artifact_root=artifact_root)
record[target] = observed[path]


def load_index_snapshot(
path: Path,
*,
artifact_root: Path | None = None,
include_artifact_status: bool = True,
) -> RunIndexSnapshot:
"""Decode the complete index; derived artifact status is optional internally.

Full readers keep fresh status by default. History collection defers only
filesystem observation until its complete semantic selection is known.
"""
try:
stream = path.open("rb")
except FileNotFoundError:
return RunIndexSnapshot(records=[], raw_count=0, digest=None)

records: list[dict[str, Any]] = []
positions: dict[tuple[str, str, str], int] = {}
artifact_exists: dict[tuple[str, str], bool] = {}
raw_count = 0
digest = hashlib.sha256()

def artifact_is_present(value: Any) -> bool:
text = str(value or "").strip()
cache_key = (text, str(artifact_root or ""))
if cache_key not in artifact_exists:
artifact_exists[cache_key] = _indexed_artifact_exists(
value,
artifact_root=artifact_root,
)
return artifact_exists[cache_key]

with stream:
for encoded_line in stream:
digest.update(encoded_line)
Expand All @@ -275,13 +284,16 @@ def artifact_is_present(value: Any) -> bool:
str(item.get("markdown_path") or ""),
)
item = dict(item)
item["json_exists"] = artifact_is_present(item.get("json_path"))
item["markdown_exists"] = artifact_is_present(item.get("markdown_path"))
# These are observations, never persisted index authority.
item.pop("json_exists", None)
item.pop("markdown_exists", None)
if key in positions:
records[positions[key]].update(item)
else:
positions[key] = len(records)
records.append(item)
if include_artifact_status:
_observe_run_artifacts(records, artifact_root=artifact_root)
return RunIndexSnapshot(
records=records,
raw_count=raw_count,
Expand Down Expand Up @@ -365,7 +377,7 @@ def collect_history(
index_path = runtime_root / "goals" / current_goal_id / "runs" / "index.jsonl"
index_snapshot = load_index_snapshot(
index_path,
artifact_root=registry_project_root(registry_path),
include_artifact_status=False,
)
runs = index_snapshot.records
raw_count = index_snapshot.raw_count
Expand Down Expand Up @@ -433,6 +445,18 @@ def collect_history(
),
"semantic_history": goal_semantic_history_from_runs(runs),
}
# Reduce the full index first. Old semantic evidence and lane windows
# may survive outside --limit; every returned Run still needs fresh
# flags. These references also cover the global recent_runs window.
status_run = goal_record["latest_status_run"]
_observe_run_artifacts(
chain(
goal_record["latest_runs"],
[status_run] if status_run is not None else [],
iter_goal_semantic_history_runs(goal_record["semantic_history"]),
),
artifact_root=registry_project_root(registry_path),
)
if registry_member:
for field in REGISTRY_STATUS_FIELDS:
if meta.get(field):
Expand Down
8 changes: 4 additions & 4 deletions loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -1735,31 +1735,31 @@
},
{
"site": "loopx/history.py::<module>.collect_history::codec_read:load_registry#1",
"line": 329,
"line": 341,
"column": 20,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/history.py::<module>.inspect_index_duplicates::codec_read:load_registry#1",
"line": 573,
"line": 597,
"column": 16,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/history.py::<module>.rebuild_index_artifact_collisions::codec_read:load_registry#1",
"line": 787,
"line": 811,
"column": 16,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/history.py::<module>.repair_index_duplicates::codec_read:load_registry#1",
"line": 677,
"line": 701,
"column": 16,
"kind": "codec_read",
"api": "load_registry",
Expand Down
179 changes: 179 additions & 0 deletions tests/test_history_artifact_observation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,179 @@
"""History retains control evidence without probing every discarded artifact."""
from __future__ import annotations

import hashlib
import json
import subprocess
import sys
from datetime import datetime, timedelta, timezone

import pytest

import loopx.history as history_module
from loopx.history import collect_history, load_index_snapshot


@pytest.fixture
def history_case(tmp_path, monkeypatch):
registry = tmp_path / "registry.json"
registry.write_text("{}\n")
runtime = tmp_path / "runtime"
index = runtime / "goals" / "example" / "runs" / "index.jsonl"
index.parent.mkdir(parents=True)
artifacts = tmp_path / "artifacts"
artifacts.mkdir()
start = datetime(2026, 9, 1, tzinfo=timezone.utc)
now = start + timedelta(seconds=1001)
monkeypatch.setattr(
"loopx.control_plane.runtime.run_context_retention.now_utc_iso",
lambda: now.isoformat(),
)
rows = []
for n in range(1000):
rows.append({
"generated_at": (start + timedelta(seconds=n)).isoformat(),
"run_id": f"run-{n}", "classification": "progress",
"agent_id": "worker", "json_path": f"artifacts/{n}.json",
"markdown_path": f"artifacts/{n}.md",
# Persisted booleans are deliberately wrong: each read observes files.
"json_exists": True, "markdown_exists": True,
})
rows[0]["agent_vision"] = {"agent_id": "worker", "revision": "retained"}
rows[1]["human_reward"] = {"decision": "revise", "reason_summary": "Keep evidence"}
rows[990].update({
"todo_id": "todo_blocked_work", "turn_instance_id": "turn-1",
"delivery_outcome": "outcome_gap",
"progress_observation": {"result_class": "blocked"},
"blocked_retry": {
"schema_version": "quota_blocked_retry_v0", "source": "turn_settlement",
"todo_id": "todo_blocked_work", "observed_at": rows[990]["generated_at"],
"due_at": (start + timedelta(seconds=1290)).isoformat(),
"resume_when": f"resume_at:{(start + timedelta(seconds=1290)).isoformat()}",
},
})
rows[995]["agent_id"] = "quiet-worker"
# Recent neutral records must not hide the latest status-relevant run.
rows[999]["classification"] = "quota_slot_spent"
rows[998]["classification"] = "quota_slot_spent"
# Last duplicate replaces only supplied fields, without another unique Run.
duplicate = {**rows[0], "annotation": "merged duplicate"}
index.write_text("\n".join(json.dumps(row) for row in [*rows, duplicate]) + "\nnot-json\n[]\n\n")
for n in (0, 1, 990, 995, 997, 999):
(artifacts / f"{n}.json").write_text("{}\n")
return registry, runtime, index, artifacts


def _collect(case, *, limit=1, agent_lane_id=None):
registry, runtime, _, _ = case
return collect_history(registry_path=registry, runtime_root=runtime,
goal_id="example", limit=limit, agent_lane_id=agent_lane_id)


@pytest.mark.parametrize("limit", [0, 1, 3])
def test_history_observes_retained_evidence_after_full_reduction(history_case, monkeypatch, limit):
_, _, index, _ = history_case
calls = []
original = history_module._indexed_artifact_exists

def observe(value, *, artifact_root):
calls.append(value)
return original(value, artifact_root=artifact_root)

monkeypatch.setattr(history_module, "_indexed_artifact_exists", observe)
result = _collect(history_case, limit=limit)
goal = result["goals"][0]
assert result["run_count"] == goal["unique_runs"] == 1000
assert goal["raw_index_records"] == 1003
assert goal["index_digest"] == f"sha256:{hashlib.sha256(index.read_bytes()).hexdigest()}"
assert len(result["runs"]) == len(goal["latest_runs"]) == limit
semantic = goal["semantic_history"]
vision = semantic["agents"][0]["latest_agent_vision_run"]
correction = semantic["latest_owner_correction_run"]
retry = semantic["active_blocked_retry_runs"][0]
assert vision["run_id"] == "run-0"
assert vision["annotation"] == "merged duplicate"
assert correction["run_id"] == "run-1"
assert retry["run_id"] == "run-990"
assert goal["latest_status_run"]["run_id"] == "run-997"
for row in [vision, correction, retry, goal["latest_status_run"]]:
assert row["json_exists"] is True
assert row["markdown_exists"] is False
# Counts depend on retained output, never the 1,000-row historical population.
expected = {0, 1, 990, 997, *range(1000 - limit, 1000)}
assert set(calls) == {f"artifacts/{n}.{ext}" for n in expected for ext in ("json", "md")}
assert len(calls) == len(set(calls))


def test_lane_window_artifacts_remain_observed(history_case):
result = _collect(history_case, agent_lane_id="quiet-worker")
runs = result["goals"][0]["latest_runs"]
assert [run["run_id"] for run in runs] == ["run-999", "run-995"]
assert all(run["json_exists"] and not run["markdown_exists"] for run in runs)


def test_unchanged_index_does_not_cache_artifact_existence(history_case):
_, _, index, artifacts = history_case
before = index.read_bytes()
initial = _collect(history_case)
assert initial["runs"][0]["json_exists"] is True
assert initial["runs"][0]["markdown_exists"] is False
(artifacts / "999.json").unlink()
(artifacts / "999.md").write_text("# Delivered\n")
(artifacts / "0.json").unlink()
reread = _collect(history_case)
assert reread["runs"][0]["json_exists"] is False
assert reread["runs"][0]["markdown_exists"] is True
assert reread["goals"][0]["semantic_history"]["agents"][0]["latest_agent_vision_run"]["json_exists"] is False
assert index.read_bytes() == before
assert reread["goals"][0]["index_digest"] == initial["goals"][0]["index_digest"]


def test_full_index_reader_still_observes_every_row(history_case):
registry, _, index, artifacts = history_case
(artifacts / "500.md").write_text("# Older artifact\n")
(artifacts / "501.json").symlink_to(artifacts / "absent.json")
snapshot = load_index_snapshot(index, artifact_root=registry.parent)
assert len(snapshot.records) == 1000
assert snapshot.records[500]["markdown_exists"] is True
assert snapshot.records[501]["json_exists"] is False
assert snapshot.records[502]["json_exists"] is False


@pytest.mark.parametrize("provider", ["file", "sqlite"])
def test_public_history_and_status_keep_artifacts_with_canonical_todos(history_case, monkeypatch, tmp_path, provider):
from tests.control_plane.canonical_authority_fixture import (
initialize_canonical_authority, isolate_sqlite_runtime,
)
from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection
from loopx.control_plane.effect_runtime import restart_effect_runtime
from loopx.control_plane.testing.canary_harness import write_fixture_registry

isolate_sqlite_runtime(tmp_path, monkeypatch)
registry, runtime, _, _ = history_case
state = tmp_path / "state.md"
state.write_text("---\nstatus: active\n---\n# Example\n\n## Agent Todo\n")
write_fixture_registry(project=tmp_path, registry_path=registry, runtime_root=runtime,
goal_id="example", domain="engineering", state_file=str(state),
adapter_kind="generic_project_goal_v0")
initialize_canonical_authority(
runtime, "example", build_todo_runtime_shadow_projection(goal_id="example", todos=[]),
state_path=state, provider=provider,
)
try:
for command in ("history", "status"):
result = subprocess.run(
[sys.executable, "-m", "loopx.entrypoint", "--registry", str(registry),
"--runtime-root", str(runtime), "--format", "json", command,
"--goal-id", "example", "--limit", "1"],
capture_output=True, text=True, timeout=60,
)
assert result.returncode == 0, result.stdout + result.stderr
payload = json.loads(result.stdout)
history = payload if command == "history" else payload["run_history"]
goal = history["goals"][0]
assert goal["unique_runs"] == 1000
assert goal["latest_status_run"]["run_id"] == "run-997"
assert goal["latest_status_run"]["json_exists"] is True
assert goal["latest_status_run"]["markdown_exists"] is False
finally:
restart_effect_runtime()
Loading