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
102 changes: 83 additions & 19 deletions loopx/chat_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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)

Expand All @@ -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."""
Expand Down
226 changes: 225 additions & 1 deletion tests/test_chat_event_retention.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from datetime import datetime, timedelta, timezone
import gc
import inspect
import json
from pathlib import Path

Expand All @@ -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(
{
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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"]
Loading