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
146 changes: 79 additions & 67 deletions loopx/history.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
from __future__ import annotations

import json
from contextlib import nullcontext
from collections.abc import Callable
from datetime import datetime, timezone
from heapq import merge
from itertools import islice
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 (
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 {
Expand Down
129 changes: 129 additions & 0 deletions tests/test_history_index_write_serialization.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
"""Refs GH-C07: per-goal history index writers share one lock.

`write_reserved_run_artifacts` appends to `<runs>/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"
Loading