diff --git a/loopx/history.py b/loopx/history.py index 7d247d9ea..3515b0bed 100644 --- a/loopx/history.py +++ b/loopx/history.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +from contextlib import nullcontext from collections.abc import Callable from datetime import datetime, timezone from heapq import merge @@ -8,6 +9,7 @@ from pathlib import Path from typing import Any +from .file_lock import exclusive_file_lock from .authority import goal_authority_registry_summary from .control_plane import compact_control_plane_policy from .control_plane.goals.activation import ( @@ -121,22 +123,24 @@ def write_reserved_run_artifacts( # the durable run record and its index row before append. Fail closed here so # malformed or negative usage never enters run history. ingest_usage_into_run_record(record, index_record=index_record) - json_path, markdown_path = reserve_unique_run_paths(runs_dir, generated_at) + # GH-C07: one lock per goal history index, shared with the repair path. index_path = runs_dir / "index.jsonl" - index_record["json_path"] = str(json_path) - index_record["markdown_path"] = str(markdown_path) - payload["json_path"] = str(json_path) - payload["markdown_path"] = str(markdown_path) - payload["index_path"] = str(index_path) - if isinstance(record.get("usage"), dict): - payload["usage"] = dict(record["usage"]) - json_path.write_text( - json.dumps(record, ensure_ascii=False, indent=2, allow_nan=False) + "\n", - encoding="utf-8", - ) - markdown_path.write_text(render_markdown(payload) + "\n", encoding="utf-8") - with index_path.open("a", encoding="utf-8") as f: - f.write(json.dumps(index_record, ensure_ascii=False, allow_nan=False) + "\n") + with exclusive_file_lock(index_path, operation="history_run_append"): + json_path, markdown_path = reserve_unique_run_paths(runs_dir, generated_at) + index_record["json_path"] = str(json_path) + index_record["markdown_path"] = str(markdown_path) + payload["json_path"] = str(json_path) + payload["markdown_path"] = str(markdown_path) + payload["index_path"] = str(index_path) + if isinstance(record.get("usage"), dict): + payload["usage"] = dict(record["usage"]) + json_path.write_text( + json.dumps(record, ensure_ascii=False, indent=2, allow_nan=False) + "\n", + encoding="utf-8", + ) + markdown_path.write_text(render_markdown(payload) + "\n", encoding="utf-8") + with index_path.open("a", encoding="utf-8") as f: + f.write(json.dumps(index_record, ensure_ascii=False, allow_nan=False) + "\n") def validate_goal_id_path_segment(goal_id: str) -> str: @@ -507,60 +511,68 @@ def repair_index_duplicates( if not index_path.exists(): continue - raw_lines = index_path.read_text(encoding="utf-8").splitlines() - grouped: dict[tuple[str, str, str], list[tuple[int, dict[str, Any]]]] = {} - for line_number, line in enumerate(raw_lines, start=1): - if not line.strip(): - continue - raw_index_records += 1 - try: - item = json.loads(line) - except json.JSONDecodeError: - continue - if not isinstance(item, dict): - continue - grouped.setdefault(index_identity(item), []).append((line_number, item)) - - remove_lines: set[int] = set() - for records in grouped.values(): - if len(records) <= 1: - continue - decision = duplicate_repair_decision(records) - removed_lines = list(decision.get("removed_line_numbers") or []) - if decision.get("action") == "preserve_reward_overlay": - preserved_reward_overlay_rows += len(records) - 1 - elif decision.get("repairable"): - remove_lines.update(int(line_number) for line_number in removed_lines) - removed_row_count += len(removed_lines) - else: - unrepaired_group_count += 1 + # GH-C07: read and rewrite the index under the same lock the append + # path takes. A dry run only reports, so it must not block writers. + lock = ( + exclusive_file_lock(index_path, operation="history_index_repair") + if execute + else nullcontext() + ) + with lock: + raw_lines = index_path.read_text(encoding="utf-8").splitlines() + grouped: dict[tuple[str, str, str], list[tuple[int, dict[str, Any]]]] = {} + for line_number, line in enumerate(raw_lines, start=1): + if not line.strip(): + continue + raw_index_records += 1 + try: + item = json.loads(line) + except json.JSONDecodeError: + continue + if not isinstance(item, dict): + continue + grouped.setdefault(index_identity(item), []).append((line_number, item)) - first_record = records[0][1] - groups.append( - { - "goal_id": current_goal_id, - "index_path": str(index_path), - "generated_at": first_record.get("generated_at"), - "json_path": first_record.get("json_path"), - "markdown_path": first_record.get("markdown_path"), - "action": decision.get("action"), - "repairable": decision.get("repairable"), - "line_numbers": decision.get("line_numbers"), - "kept_line_numbers": decision.get("kept_line_numbers"), - "removed_line_numbers": removed_lines, - "reason": decision.get("reason"), - } - ) + remove_lines: set[int] = set() + for records in grouped.values(): + if len(records) <= 1: + continue + decision = duplicate_repair_decision(records) + removed_lines = list(decision.get("removed_line_numbers") or []) + if decision.get("action") == "preserve_reward_overlay": + preserved_reward_overlay_rows += len(records) - 1 + elif decision.get("repairable"): + remove_lines.update(int(line_number) for line_number in removed_lines) + removed_row_count += len(removed_lines) + else: + unrepaired_group_count += 1 + + first_record = records[0][1] + groups.append( + { + "goal_id": current_goal_id, + "index_path": str(index_path), + "generated_at": first_record.get("generated_at"), + "json_path": first_record.get("json_path"), + "markdown_path": first_record.get("markdown_path"), + "action": decision.get("action"), + "repairable": decision.get("repairable"), + "line_numbers": decision.get("line_numbers"), + "kept_line_numbers": decision.get("kept_line_numbers"), + "removed_line_numbers": removed_lines, + "reason": decision.get("reason"), + } + ) - if execute and remove_lines: - rewritten = [ - line - for line_number, line in enumerate(raw_lines, start=1) - if line_number not in remove_lines - ] - tmp_path = index_path.with_suffix(index_path.suffix + ".tmp") - tmp_path.write_text("".join(line + "\n" for line in rewritten), encoding="utf-8") - tmp_path.replace(index_path) + if execute and remove_lines: + rewritten = [ + line + for line_number, line in enumerate(raw_lines, start=1) + if line_number not in remove_lines + ] + tmp_path = index_path.with_suffix(index_path.suffix + ".tmp") + tmp_path.write_text("".join(line + "\n" for line in rewritten), encoding="utf-8") + tmp_path.replace(index_path) limited_groups = groups[: max(0, limit)] return { diff --git a/tests/test_history_index_write_serialization.py b/tests/test_history_index_write_serialization.py new file mode 100644 index 000000000..11c292aa4 --- /dev/null +++ b/tests/test_history_index_write_serialization.py @@ -0,0 +1,129 @@ +"""Refs GH-C07: per-goal history index writers share one lock. + +`write_reserved_run_artifacts` appends to `/index.jsonl` while +`repair_index_duplicates` rewrites that same file from a snapshot. If only one +side locks, an append that lands after the repair reads the index is still lost, +so the durable invariant is that **both** sides take the lock for the **same** +path. These cases pin that, plus the rule that a dry-run repair stays read-only +and does not block writers. +""" + +from __future__ import annotations + +import json +from contextlib import contextmanager +from pathlib import Path +from typing import Any, Iterator + +import pytest + +from loopx import history + +GOAL_ID = "goal-history-lock" + + +@pytest.fixture +def record_lock_calls(monkeypatch: pytest.MonkeyPatch) -> list[Path]: + """Capture every lock path taken through the history module.""" + taken: list[Path] = [] + real_lock = history.exclusive_file_lock + + @contextmanager + def recording_lock(path: Path, **_kwargs: Any) -> Iterator[Path]: + taken.append(path) + with real_lock(path) as locked: + yield locked + + monkeypatch.setattr(history, "exclusive_file_lock", recording_lock) + return taken + + +def _write_artifacts(runs_dir: Path) -> Path: + runs_dir.mkdir(parents=True, exist_ok=True) + generated_at = "2026-09-16T00:00:00+00:00" + history.write_reserved_run_artifacts( + runs_dir=runs_dir, + generated_at=generated_at, + record={"goal_id": GOAL_ID, "generated_at": generated_at}, + index_record={"goal_id": GOAL_ID, "generated_at": generated_at}, + payload={"goal_id": GOAL_ID, "generated_at": generated_at}, + render_markdown=lambda payload: "# record\n", + ) + return runs_dir / "index.jsonl" + + +def test_append_takes_the_goal_index_lock(tmp_path: Path, record_lock_calls: list[Path]) -> None: + index_path = _write_artifacts(tmp_path / "runs") + + assert record_lock_calls == [index_path], record_lock_calls + assert index_path.exists() + + +def test_append_is_visible_after_the_lock_released( + tmp_path: Path, record_lock_calls: list[Path] +) -> None: + index_path = _write_artifacts(tmp_path / "runs") + + rows = [ + json.loads(line) + for line in index_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + assert [row["goal_id"] for row in rows] == [GOAL_ID] + + +def _history_fixture(root: Path) -> tuple[Path, Path]: + runtime_root = root / "runtime" + index_path = runtime_root / "goals" / GOAL_ID / "runs" / "index.jsonl" + index_path.parent.mkdir(parents=True, exist_ok=True) + row = {"goal_id": GOAL_ID, "generated_at": "2026-09-16T00:00:00+00:00", "kind": "run"} + index_path.write_text( + "".join(json.dumps(row) + "\n" for _ in range(2)), encoding="utf-8" + ) + registry_path = root / "registry.json" + registry_path.write_text( + json.dumps( + { + "schema_version": "0.1", + "common_runtime_root": str(runtime_root), + "projects": [], + "goals": [], + } + ), + encoding="utf-8", + ) + return registry_path, runtime_root + + +def test_repair_locks_the_same_index_path_on_execute( + tmp_path: Path, record_lock_calls: list[Path] +) -> None: + registry_path, runtime_root = _history_fixture(tmp_path) + index_path = runtime_root / "goals" / GOAL_ID / "runs" / "index.jsonl" + + history.repair_index_duplicates( + registry_path=registry_path, + runtime_root_override=str(runtime_root), + goal_id=GOAL_ID, + limit=10, + execute=True, + ) + + assert record_lock_calls == [index_path], record_lock_calls + + +def test_repair_dry_run_does_not_take_the_write_lock( + tmp_path: Path, record_lock_calls: list[Path] +) -> None: + registry_path, runtime_root = _history_fixture(tmp_path) + + result = history.repair_index_duplicates( + registry_path=registry_path, + runtime_root_override=str(runtime_root), + goal_id=GOAL_ID, + limit=10, + execute=False, + ) + + assert result["dry_run"] is True + assert record_lock_calls == [], "a read-only repair must not block writers"