diff --git a/docs/status-data-contract.md b/docs/status-data-contract.md index 568a4f2ca6..22f291deb4 100644 --- a/docs/status-data-contract.md +++ b/docs/status-data-contract.md @@ -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 diff --git a/loopx/control_plane/runtime/run_context_retention.py b/loopx/control_plane/runtime/run_context_retention.py index 1b0149aa49..a62957fa7b 100644 --- a/loopx/control_plane/runtime/run_context_retention.py +++ b/loopx/control_plane/runtime/run_context_retention.py @@ -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 @@ -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, *, diff --git a/loopx/history.py b/loopx/history.py index 7c44802ab5..6dba847635 100644 --- a/loopx/history.py +++ b/loopx/history.py @@ -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 @@ -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 ( @@ -229,11 +230,30 @@ 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: @@ -241,20 +261,9 @@ def load_index_snapshot( 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) @@ -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, @@ -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 @@ -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): diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 09c271ee70..346264a947 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1735,7 +1735,7 @@ }, { "site": "loopx/history.py::.collect_history::codec_read:load_registry#1", - "line": 329, + "line": 341, "column": 20, "kind": "codec_read", "api": "load_registry", @@ -1743,7 +1743,7 @@ }, { "site": "loopx/history.py::.inspect_index_duplicates::codec_read:load_registry#1", - "line": 573, + "line": 597, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -1751,7 +1751,7 @@ }, { "site": "loopx/history.py::.rebuild_index_artifact_collisions::codec_read:load_registry#1", - "line": 787, + "line": 811, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -1759,7 +1759,7 @@ }, { "site": "loopx/history.py::.repair_index_duplicates::codec_read:load_registry#1", - "line": 677, + "line": 701, "column": 16, "kind": "codec_read", "api": "load_registry", diff --git a/tests/test_history_artifact_observation.py b/tests/test_history_artifact_observation.py new file mode 100644 index 0000000000..4f67c8d6d3 --- /dev/null +++ b/tests/test_history_artifact_observation.py @@ -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()