From 9158ff63d62269f04cbd2b7b5ad02f60b0e06b9e Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Sat, 27 Jun 2026 11:55:46 -0400 Subject: [PATCH] session: add jsonl transcript persistence --- tartarus/session.py | 123 +++++++++++++++++++++++++++++++++++++++ tests/test_session.py | 132 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 255 insertions(+) create mode 100644 tartarus/session.py create mode 100644 tests/test_session.py diff --git a/tartarus/session.py b/tartarus/session.py new file mode 100644 index 0000000..c230af9 --- /dev/null +++ b/tartarus/session.py @@ -0,0 +1,123 @@ +"""Append-only JSONL persistence for conversation transcripts (PLAN.md §10). + +A session is the provider-native `messages` list the AgentLoop maintains. The loop +only ever appends to it, and only at whole-round-trip boundaries, so a session +file is always a valid transcript: one JSON message per line, in order. Modeled on +`FileAuditLog` (audit.py) — makedirs on write, OSError surfaced as a typed error. +""" + +import json +import os +import secrets +import time + +SESSION_SUFFIX = ".jsonl" + +# Bytes of random entropy appended to the timestamp in new_id(). +# 8 bytes makes accidental collisions astronomically unlikely while keeping the +# id short enough to type/see in a file listing. +_ID_RANDOM_BYTES = 4 + + +class SessionError(Exception): + """Raised when a session cannot be read or written.""" + + +class SessionStore: + def __init__(self, session_dir: str, session_id: str): + self._dir = os.path.abspath(session_dir) + self.session_id = session_id + self._path = os.path.join(self._dir, session_id + SESSION_SUFFIX) + # Number of messages already on disk; the next append writes the tail + # past this index so re-appends never duplicate. + self._flushed = 0 + + @property + def path(self) -> str: + return self._path + + @staticmethod + def new_id() -> str: + """A sortable, collision-resistant, typeable id. + + Timestamp prefix (down to microseconds) means lexical order equals + chronological order, so `latest` is a plain max. The random suffix + avoids extremely rare same-microsecond clashes. + """ + now = time.time() + secs = int(now) + usecs = int((now - secs) * 1_000_000) + timestamp = time.strftime("%Y%m%d-%H%M%S", time.localtime(secs)) + return f"{timestamp}-{usecs:06d}-{secrets.token_hex(_ID_RANDOM_BYTES)}" + + @staticmethod + def list_ids(session_dir: str) -> list[str]: + """Existing session ids, newest first (lexical order == chronological).""" + try: + names = os.listdir(session_dir) + except FileNotFoundError: + return [] + except OSError as error: + raise SessionError( + f"cannot list sessions in {session_dir}: {error}" + ) from error + ids = [n[: -len(SESSION_SUFFIX)] for n in names if n.endswith(SESSION_SUFFIX)] + return sorted(ids, reverse=True) + + @staticmethod + def latest(session_dir: str) -> str | None: + """The most recent session id, or None if there are none.""" + ids = SessionStore.list_ids(session_dir) + return ids[0] if ids else None + + @staticmethod + def resolve(session_dir: str, wanted: str) -> str: + """Resolve `wanted` to a session id: exact match, else unique prefix. + + Raises SessionError if nothing matches or a prefix is ambiguous. + """ + ids = SessionStore.list_ids(session_dir) + if wanted in ids: + return wanted + matches = [i for i in ids if i.startswith(wanted)] + if not matches: + raise SessionError(f"no session matching '{wanted}'") + if len(matches) > 1: + raise SessionError(f"'{wanted}' is ambiguous: matches {', '.join(matches)}") + return matches[0] + + def load(self) -> list[dict]: + """Read the persisted transcript, marking those messages as flushed.""" + try: + with open(self._path, encoding="utf-8") as session_file: + messages = [json.loads(line) for line in session_file if line.strip()] + except FileNotFoundError: + return [] + except (OSError, json.JSONDecodeError) as error: + raise SessionError(f"cannot read session {self._path}: {error}") from error + self._flushed = len(messages) + return messages + + def append(self, messages: list[dict]) -> None: + """Persist any messages added since the last flush.""" + tail = messages[self._flushed :] + if not tail: + return + if self._dir: + os.makedirs(self._dir, exist_ok=True) + try: + with open(self._path, "a", encoding="utf-8") as session_file: + for message in tail: + session_file.write(json.dumps(message) + "\n") + except OSError as error: + raise SessionError(f"cannot write session {self._path}: {error}") from error + self._flushed = len(messages) + + def first_user_message(self) -> str | None: + """First user-authored text in the session, for listing previews.""" + for message in self.load(): + if message.get("role") == "user": + content = message.get("content") + if isinstance(content, str): + return content + return None diff --git a/tests/test_session.py b/tests/test_session.py new file mode 100644 index 0000000..8a37b02 --- /dev/null +++ b/tests/test_session.py @@ -0,0 +1,132 @@ +import re + +import pytest + +from tartarus.session import SessionError, SessionStore + + +def test_new_id_is_sortable_and_unique(): + ids = sorted(SessionStore.new_id() for _ in range(50)) + # Shape: YYYYMMDD-HHMMSS-ffffff-xxxxxxxx + assert all(re.fullmatch(r"\d{8}-\d{6}-\d{6}-[0-9a-f]{8}", i) for i in ids) + assert len(set(ids)) == len(ids) + + +def test_append_then_load_round_trips(tmp_path): + store = SessionStore(str(tmp_path), "20260627-120000-aaaa") + messages = [ + {"role": "user", "content": "hi"}, + {"role": "assistant", "content": "hello"}, + ] + store.append(messages) + + reopened = SessionStore(str(tmp_path), "20260627-120000-aaaa") + assert reopened.load() == messages + + +def test_append_is_incremental_and_does_not_duplicate(tmp_path): + store = SessionStore(str(tmp_path), "s1") + store.append([{"role": "user", "content": "one"}]) + store.append( + [{"role": "user", "content": "one"}, {"role": "assistant", "content": "two"}] + ) + + assert SessionStore(str(tmp_path), "s1").load() == [ + {"role": "user", "content": "one"}, + {"role": "assistant", "content": "two"}, + ] + + +def test_load_after_resume_marks_existing_as_flushed(tmp_path): + SessionStore(str(tmp_path), "s1").append([{"role": "user", "content": "old"}]) + + resumed = SessionStore(str(tmp_path), "s1") + messages = resumed.load() + messages.append({"role": "assistant", "content": "new"}) + resumed.append(messages) + + # Only the new message was written; the old one is not duplicated. + assert SessionStore(str(tmp_path), "s1").load() == [ + {"role": "user", "content": "old"}, + {"role": "assistant", "content": "new"}, + ] + + +def test_load_missing_session_returns_empty(tmp_path): + assert SessionStore(str(tmp_path), "nope").load() == [] + + +def test_append_nothing_writes_no_file(tmp_path): + store = SessionStore(str(tmp_path), "empty") + store.append([]) + assert not (tmp_path / "empty.jsonl").exists() + + +def test_latest_and_list_ids_order_newest_first(tmp_path): + for session_id in ("20260101-000000-aaaa", "20260627-120000-bbbb"): + SessionStore(str(tmp_path), session_id).append( + [{"role": "user", "content": "x"}] + ) + + assert SessionStore.latest(str(tmp_path)) == "20260627-120000-bbbb" + assert SessionStore.list_ids(str(tmp_path)) == [ + "20260627-120000-bbbb", + "20260101-000000-aaaa", + ] + + +def test_latest_and_list_ids_on_missing_dir(tmp_path): + missing = str(tmp_path / "nope") + assert SessionStore.latest(missing) is None + assert SessionStore.list_ids(missing) == [] + + +def test_resolve_matches_exact_then_unique_prefix(tmp_path): + for session_id in ("20260627-120000-aaaa", "20260627-130000-bbbb"): + SessionStore(str(tmp_path), session_id).append( + [{"role": "user", "content": "x"}] + ) + + assert ( + SessionStore.resolve(str(tmp_path), "20260627-120000-aaaa") + == "20260627-120000-aaaa" + ) + # Unambiguous prefix resolves. + assert SessionStore.resolve(str(tmp_path), "20260627-13") == "20260627-130000-bbbb" + + +def test_resolve_rejects_missing_and_ambiguous(tmp_path): + for session_id in ("20260627-120000-aaaa", "20260627-130000-bbbb"): + SessionStore(str(tmp_path), session_id).append( + [{"role": "user", "content": "x"}] + ) + + with pytest.raises(SessionError, match="no session"): + SessionStore.resolve(str(tmp_path), "nope") + with pytest.raises(SessionError, match="ambiguous"): + SessionStore.resolve(str(tmp_path), "20260627") + + +def test_list_ids_surfaces_os_errors_as_session_error(tmp_path): + bad_path = str(tmp_path / "not_a_dir") + (tmp_path / "not_a_dir").write_text("i am a file, not a directory") + with pytest.raises(SessionError, match="cannot list sessions"): + SessionStore.list_ids(bad_path) + + +def test_new_id_order_is_chronological(): + ids = [SessionStore.new_id() for _ in range(10)] + assert ids == sorted(ids) + + +def test_first_user_message_preview(tmp_path): + store = SessionStore(str(tmp_path), "s1") + store.append( + [ + {"role": "user", "content": "summarize the repo"}, + {"role": "assistant", "content": "sure"}, + ] + ) + assert ( + SessionStore(str(tmp_path), "s1").first_user_message() == "summarize the repo" + ) -- 2.51.2