From 9297750db342c5d2562f2e6a190e329b4ca226af Mon Sep 17 00:00:00 2001 From: YZJF <195568136+YZJF@users.noreply.github.com> Date: Wed, 16 Sep 2026 17:40:43 +0800 Subject: [PATCH] fix(history): serialize the per-goal index append and repair Refs GH-C07 The row asks for the global-registry write lock to reach the per-goal todo, refresh and history writers. Todo and refresh are already guarded: every todo mutation runs inside legacy_todo_write_transaction, which holds exclusive_cross_runtime_file_lock on the per-goal todo lock path, and state_refresh.py locks the state file. loopx/history.py had none. That left one real gap. index.jsonl has an appending writer and a rewriting one: write_reserved_run_artifacts appends a row, while repair_index_duplicates reads every line, decides which duplicates to drop, then replaces the file from that snapshot. An append landing between the read and the replace is lost, because the rewrite puts back a snapshot that never saw it. Both sides now take the same lock on the same index path. The repair holds it around the whole read-analyse-rewrite block, not just the final replace, since guarding only the write would still rewrite from an unlocked snapshot. A dry run keeps nullcontext and takes no lock, so a read-only preview does not block writers. No nested lock is introduced: the callers in promotion_gate.py, dreaming.py and cli_commands/history.py hold no lock of their own. Signed-off-by: YZJF <195568136+YZJF@users.noreply.github.com> --- loopx/history.py | 146 ++++++++++-------- .../test_history_index_write_serialization.py | 129 ++++++++++++++++ 2 files changed, 208 insertions(+), 67 deletions(-) create mode 100644 tests/test_history_index_write_serialization.py diff --git a/loopx/history.py b/loopx/history.py index 7d247d9ea1..3515b0bed5 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 0000000000..11c292aa46 --- /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"