diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 7762a179c9..14001fbfbf 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -79,10 +79,16 @@ def _read_json(path: Path) -> dict[str, Any]: return payload if isinstance(payload, dict) else {} -def _read_jsonl(path: Path) -> list[dict[str, Any]]: +def _read_jsonl( + path: Path, *, raise_on_error: bool = False, +) -> list[dict[str, Any]]: try: lines = path.read_bytes().split(b"\n") + except FileNotFoundError: + return [] except OSError: + if raise_on_error: + raise return [] rows: list[dict[str, Any]] = [] for line in lines: @@ -194,7 +200,7 @@ def __init__(self, runtime_root: Path) -> None: self._session_locks: WeakValueDictionary[str, threading.Lock] = WeakValueDictionary() self._event_lock = threading.RLock() self._event_cache = ChatEventCache(self._event_lock) - self._event_pending: dict[tuple[str, str], list[dict[str, Any]]] = {} + self._event_pending: dict[tuple[str, str], list[tuple[dict[str, Any], str]]] = {} self._event_flush_locks: WeakValueDictionary[tuple[str, str], threading.Lock] = WeakValueDictionary() self.sessions_root.mkdir(parents=True, exist_ok=True, mode=0o700) os.chmod(self.root, 0o700) @@ -1684,7 +1690,7 @@ def append_event( "payload": payload, } with self._event_lock: - self._event_pending.setdefault(key, []).append(event) + self._event_pending.setdefault(key, []).append((event, uuid.uuid4().hex)) if not buffered: self.flush_events(session_id, turn_id) return event @@ -1703,24 +1709,79 @@ def flush_events(self, session_id: str, turn_id: str) -> int: return flushed try: with exclusive_file_lock(path, agent_id="loopx-chat", operation="append_chat_events"): - rows = self._event_rows_locked(session_id, turn_id) - # Allocate after the last persisted sequence and append under the - # same file lock: row order stays strictly increasing by sequence. - # Compaction preserves this order but may leave sequence gaps. - sequence = int(rows[-1].get("sequence") or 0) if rows else 0 - for event in pending: - sequence += 1 - event["event_id"] = str(sequence) - event["sequence"] = sequence - _append_jsonl_rows(path, pending) - if any(row["kind"] in TERMINAL_EVENT_KINDS for row in pending): - self._event_cache.drop(key) - else: - self._event_cache.put(key, self._event_revision(path), [*rows, *pending]) + append_started = False + try: + uncertain = any("sequence" in event for event, _ in pending) + rows = ( + _read_jsonl(path, raise_on_error=True) + if uncertain + else self._event_rows_locked(session_id, turn_id) + ) + # A failed append may already have persisted all or part of + # its batch. Match only this store's stable append ids; two + # otherwise identical events are still separate writes. + durable = ( + { + row["_append_id"]: row + for row in rows + if isinstance(row.get("_append_id"), str) + } + if uncertain + else {} + ) + # Allocate after the last persisted sequence under the same + # file lock. Compaction may leave gaps; another writer may + # have appended after a partial failed batch. + sequence = int(rows[-1].get("sequence") or 0) if rows else 0 + to_append: list[dict[str, Any]] = [] + for event, append_id in pending: + committed = durable.get(append_id) + if committed is not None: + event["event_id"] = committed["event_id"] + event["sequence"] = committed["sequence"] + continue + sequence += 1 + event["event_id"] = str(sequence) + event["sequence"] = sequence + to_append.append({**event, "_append_id": append_id}) + append_started = True + _append_jsonl_rows(path, to_append) + if any(event["kind"] in TERMINAL_EVENT_KINDS for event, _ in pending): + self._event_cache.drop(key) + else: + self._event_cache.put(key, self._event_revision(path), [*rows, *to_append]) + except Exception: + if not append_started: + raise + # Confirm visible writes before releasing the file lock: + # another store may compact old events before our retry. + try: + written = { + row["_append_id"]: row + for row in _read_jsonl(path, raise_on_error=True) + if isinstance(row.get("_append_id"), str) + } + except OSError: + pass # Keep all uncertain events for a later retry. + else: + remaining = [] + for event, append_id in pending: + committed = written.get(append_id) + if committed is None: + remaining.append((event, append_id)) + else: + event["event_id"] = committed["event_id"] + event["sequence"] = committed["sequence"] + pending = remaining + raise except Exception: + self._event_cache.drop(key) with self._event_lock: later = self._event_pending.get(key, []) - self._event_pending[key] = [*pending, *later] + if pending or later: + self._event_pending[key] = [*pending, *later] + else: + self._event_pending.pop(key, None) raise flushed += len(pending) @@ -1743,7 +1804,10 @@ def events_after(self, session_id: str, turn_id: str, event_id: str | None) -> l result = rows[start:] if rows and rows[-1].get("kind") in TERMINAL_EVENT_KINDS: self._event_cache.retain_terminal(key, rows) - return result + return [ + {field: value for field, value in row.items() if field != "_append_id"} + for row in result + ] def compact_completed_events(self, *, older_than_hours: float = 24.0) -> int: """Drop replay-only deltas after the durable final message is old enough.""" diff --git a/tests/test_chat_event_retention.py b/tests/test_chat_event_retention.py index 5235d12d0f..0c93238ff3 100644 --- a/tests/test_chat_event_retention.py +++ b/tests/test_chat_event_retention.py @@ -2,6 +2,7 @@ from datetime import datetime, timedelta, timezone import gc +import inspect import json from pathlib import Path @@ -13,7 +14,7 @@ def _write_completed_turn(root: Path, *, with_events: bool = True) -> tuple[Path, Path]: turn_path = root / "chat" / "sessions" / "session" / "turns" / "turn.json" - turn_path.parent.mkdir(parents=True) + turn_path.parent.mkdir(parents=True, exist_ok=True) turn_path.write_text( json.dumps( { @@ -64,6 +65,228 @@ def counted(path): assert len(reads) == 3 # independent writer plus invalidated reader +def test_event_flush_retry_after_durable_fsync_keeps_distinct_identical_events( + tmp_path: Path, monkeypatch, +) -> None: + store = ChatSessionStore(tmp_path) + key = ("session", "turn") + monkeypatch.setattr(chat_store, "utc_now", lambda: "2026-10-05T00:00:00Z") + first = store.append_event( + *key, kind="assistant.delta", payload={"text": "same"}, buffered=True + ) + second = store.append_event( + *key, kind="assistant.delta", payload={"text": "same"}, buffered=True + ) + actual_fsync = chat_store.os.fsync + failed = False + + def uncertain_fsync(fd: int) -> None: + nonlocal failed + actual_fsync(fd) + caller_name = inspect.currentframe().f_back.f_code.co_name + if not failed and caller_name == "_append_jsonl_rows": + failed = True + raise OSError("fsync completed, but its result was lost") + + monkeypatch.setattr(chat_store.os, "fsync", uncertain_fsync) + with pytest.raises(OSError, match="result was lost"): + store.flush_events(*key) + + event_path = store._event_path(*key) + assert [row["sequence"] for row in chat_store._read_jsonl(event_path)] == [1, 2] + assert store.flush_events(*key) == 0 + replay = ChatSessionStore(tmp_path).events_after(*key, None) + assert [row["sequence"] for row in replay] == [1, 2] + assert [row["payload"] for row in replay] == [{"text": "same"}, {"text": "same"}] + assert [first["sequence"], second["sequence"]] == [1, 2] + assert store.events_after(*key, None) == replay + assert all("_append_id" not in row for row in replay) + _write_completed_turn(tmp_path, with_events=False) + + from http.client import HTTPConnection + from threading import Thread + from loopx.chat_server import ChatHTTPServer, ChatRequestHandler + + server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler) + server.verbose = False + server.chat_store = store + thread = Thread(target=server.serve_forever, daemon=True) + thread.start() + connection = HTTPConnection("127.0.0.1", server.server_address[1], timeout=5) + try: + connection.request("GET", "/api/chat/sessions/session/turns/turn/events") + response = connection.getresponse() + body = response.read().decode("utf-8") + assert response.status == 200 + assert response.getheader("Content-Type") == "text/event-stream; charset=utf-8" + assert [line for line in body.splitlines() if line.startswith("id: ")] == [ + "id: 1", "id: 2", + ] + data = [ + json.loads(line[6:]) + for line in body.splitlines() + if line.startswith("data: ") + ] + assert [row["sequence"] for row in data] == [1, 2] + assert all("_append_id" not in row for row in data) + finally: + connection.close() + server.shutdown() + thread.join(timeout=5) + server.server_close() + + +def test_failed_flush_does_not_resurrect_events_after_other_store_compacts( + tmp_path: Path, monkeypatch, +) -> None: + store = ChatSessionStore(tmp_path) + _write_completed_turn(tmp_path, with_events=False) + key = ("session", "turn") + delta = store.append_event( + *key, kind="assistant.delta", payload={"text": "late"}, buffered=True + ) + terminal = store.append_event(*key, kind="turn.completed", payload={}, buffered=True) + append_rows = chat_store._append_jsonl_rows + failed = False + + def durable_then_lose_result(path: Path, rows: list[dict]) -> None: + nonlocal failed + append_rows(path, rows) + if not failed: + failed = True + raise OSError("durable write result lost") + + monkeypatch.setattr(chat_store, "_append_jsonl_rows", durable_then_lose_result) + with pytest.raises(OSError, match="result lost"): + store.flush_events(*key) + path = store._event_path(*key) + assert [(row["sequence"], row["kind"]) for row in chat_store._read_jsonl(path)] == [ + (1, "assistant.delta"), (2, "turn.completed"), + ] + + compacting_store = ChatSessionStore(tmp_path) + assert [row["kind"] for row in compacting_store.events_after(*key, None)] == [ + "turn.completed", + ] + assert store.flush_events(*key) == 0 + assert [(row["sequence"], row["kind"]) for row in chat_store._read_jsonl(path)] == [ + (2, "turn.completed"), + ] + assert [delta["sequence"], terminal["sequence"]] == [1, 2] + + +def test_failed_flush_keeps_uncertain_events_if_locked_readback_fails( + tmp_path: Path, monkeypatch, +) -> None: + store = ChatSessionStore(tmp_path) + key = ("session", "turn") + event = store.append_event(*key, kind="assistant.delta", payload={}, buffered=True) + event_path = store._event_path(*key) + append_rows = chat_store._append_jsonl_rows + read_bytes = Path.read_bytes + readback_failures = 0 + result_lost = False + nonempty_appends = 0 + + def lose_result(path: Path, rows: list[dict]) -> None: + nonlocal readback_failures, result_lost, nonempty_appends + if rows: + nonempty_appends += 1 + append_rows(path, rows) + if not result_lost: + result_lost = True + readback_failures = 2 + raise OSError("append result lost") + + def fail_two_reads(path: Path) -> bytes: + nonlocal readback_failures + if path == event_path and readback_failures: + readback_failures -= 1 + raise PermissionError("readback temporarily unavailable") + return read_bytes(path) + + monkeypatch.setattr(chat_store, "_append_jsonl_rows", lose_result) + monkeypatch.setattr(Path, "read_bytes", fail_two_reads) + with pytest.raises(OSError, match="append result lost"): + store.flush_events(*key) + persisted_size = event_path.stat().st_size + with pytest.raises(PermissionError, match="readback temporarily unavailable"): + store.flush_events(*key) + assert event_path.stat().st_size == persisted_size + assert nonempty_appends == 1 + assert store.flush_events(*key) == 1 + assert [row["sequence"] for row in chat_store._read_jsonl(event_path)] == [1] + assert [row["sequence"] for row in store.events_after(*key, None)] == [1] + assert event["sequence"] == 1 + + +def test_event_flush_retry_after_partial_batch_and_concurrent_writer( + tmp_path: Path, monkeypatch, +) -> None: + store = ChatSessionStore(tmp_path) + key = ("session", "turn") + first = store.append_event( + *key, kind="assistant.delta", payload={"text": "first"}, buffered=True + ) + second = store.append_event( + *key, kind="assistant.delta", payload={"text": "second"}, buffered=True + ) + append_rows = chat_store._append_jsonl_rows + failed = False + + def partially_durable(path: Path, rows: list[dict]) -> None: + nonlocal failed + if not failed: + failed = True + append_rows(path, rows[:1]) + raise OSError("batch interrupted after one durable row") + append_rows(path, rows) + + monkeypatch.setattr(chat_store, "_append_jsonl_rows", partially_durable) + with pytest.raises(OSError, match="one durable row"): + store.flush_events(*key) + assert [ + row["sequence"] for row in chat_store._read_jsonl(store._event_path(*key)) + ] == [1] + + other = ChatSessionStore(tmp_path) + other.append_event(*key, kind="assistant.delta", payload={"text": "other"}) + assert store.flush_events(*key) == 1 + replay = store.events_after(*key, None) + assert [(row["sequence"], row["payload"]["text"]) for row in replay] == [ + (1, "first"), (2, "other"), (3, "second"), + ] + assert [first["sequence"], second["sequence"]] == [1, 3] + assert [row["sequence"] for row in store.events_after(*key, "1")] == [2, 3] + + +def test_event_flush_retry_after_no_write_preserves_once_only_readback( + tmp_path: Path, monkeypatch, +) -> None: + store = ChatSessionStore(tmp_path) + key = ("session", "turn") + event = store.append_event(*key, kind="turn.completed", payload={}, buffered=True) + append_rows = chat_store._append_jsonl_rows + failed = False + + def fail_before_write(path: Path, rows: list[dict]) -> None: + nonlocal failed + if not failed: + failed = True + raise OSError("write never started") + append_rows(path, rows) + + monkeypatch.setattr(chat_store, "_append_jsonl_rows", fail_before_write) + with pytest.raises(OSError, match="write never started"): + store.flush_events(*key) + assert chat_store._read_jsonl(store._event_path(*key)) == [] + assert store.flush_events(*key) == 1 + assert [ + row["sequence"] for row in ChatSessionStore(tmp_path).events_after(*key, None) + ] == [1] + assert event["sequence"] == 1 + + def test_completed_replay_evicts_least_recently_used_log(tmp_path: Path) -> None: store = ChatSessionStore(tmp_path) for i in range(8): @@ -204,3 +427,4 @@ def test_compaction_marker_does_not_make_legacy_terminal_turns_unreadable( assert persisted["event_compaction_revision"] == list( store._event_revision(event_path) or () ) + assert [row["event_id"] for row in store.events_after("session", "turn", None)] == ["2"]