Coverage for src/lilbee/sessions/store.py: 100%
251 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""Append-only JSONL store for chat sessions.
3Each session is one ``<id>.jsonl`` file under ``cfg.data_dir/sessions``. The file
4is a strictly append-only event log: one JSON object per line, appended and
5fsynced, never rewritten. Event types are ``meta`` (first line), ``title``
6(newest wins, so rename appends rather than rewrites), ``message``, and
7``summary`` (newest wins; compaction's condensed view of the turns that no
8longer fit the prompt). The only corruption an append log can suffer is a torn
9final line, which the reader skips. A fork's file is written whole once, then
10only appended to.
11"""
13from __future__ import annotations
15import json
16import os
17from collections.abc import Callable, Iterator
18from dataclasses import dataclass
19from datetime import UTC, datetime
20from enum import StrEnum
21from pathlib import Path
22from typing import Any
23from uuid import uuid4
25from filelock import FileLock
27from lilbee.core.config import cfg
28from lilbee.core.security import (
29 OWNER_ONLY_MODE,
30 ensure_private_dir,
31 private_opener,
32 write_private_text,
33)
35SESSIONS_DIRNAME = "sessions"
36SESSIONS_DISABLED_HINT = (
37 "Sessions are off. Turn them on with /set sessions_enabled true in the TUI, "
38 "settings_set via MCP, or sessions_enabled = true in config.toml."
39)
40AGENT_SESSIONS_DISABLED_HINT = (
41 "Agent sessions are off. Turn them on with settings_set mcp_sessions_enabled "
42 "true, /set mcp_sessions_enabled true in the TUI, or mcp_sessions_enabled = "
43 "true in config.toml."
44)
45# Bounds a wedged lock holder; a healthy append holds the lock for milliseconds.
46_APPEND_LOCK_TIMEOUT_S = 10
47UNTITLED_SESSION_TITLE = "Untitled chat"
48TITLE_MAX_LEN = 60
49TITLE_ELLIPSIS = "…"
50FORK_TITLE_SUFFIX = " (fork {number})"
53def sessions_enabled() -> bool:
54 """True when session persistence is on for the human surfaces (default on)."""
55 return cfg.sessions_enabled
58def agent_sessions_enabled() -> bool:
59 """True when agent (MCP) sessions are on (default off).
61 Independent of ``sessions_enabled``: the two govern separate domains, so an
62 agent's working state can be off while a human's conversations are saved.
63 """
64 return cfg.mcp_sessions_enabled
67class SessionEventType(StrEnum):
68 """The tag on every line of a session log."""
70 META = "meta"
71 TITLE = "title"
72 MESSAGE = "message"
73 SUMMARY = "summary"
74 ORIGIN = "origin"
77class SessionOrigin(StrEnum):
78 """The surface a session belongs to: whoever created it, or was last
79 explicitly transferred to. Appends from any other surface are refused, so
80 an agent cannot splice its turns into a conversation a human owns."""
82 TUI = "tui"
83 MCP = "mcp"
84 HTTP = "http"
85 CLI = "cli"
88# The surfaces a human drives directly. Their sessions are one conversation
89# space (start in Obsidian, continue in the TUI); agent sessions are working
90# state and stay out of it unless asked for.
91HUMAN_ORIGINS: frozenset[SessionOrigin] = frozenset(
92 {SessionOrigin.TUI, SessionOrigin.HTTP, SessionOrigin.CLI}
93)
96class MessageRole(StrEnum):
97 """Author of a chat message."""
99 USER = "user"
100 ASSISTANT = "assistant"
103class TitleSource(StrEnum):
104 """Where a session title came from."""
106 AUTO = "auto"
107 CUSTOM = "custom"
110@dataclass(frozen=True)
111class SessionMessage:
112 """One turn in a session. ``ts`` is stamped by the store on write."""
114 role: MessageRole
115 content: str
116 sources: tuple[str, ...] = ()
117 ts: str = ""
120@dataclass(frozen=True)
121class SessionMeta:
122 """Session metadata, reconstructed from the log without its message bodies."""
124 id: str
125 title: str
126 created_at: str
127 updated_at: str
128 model_ref: str
129 scope: str
130 message_count: int
131 origin: SessionOrigin = SessionOrigin.TUI
132 """Owning surface. Files written before ownership existed carry no origin;
133 the only writer then was the TUI, so that is the fallback."""
134 forked_from: str = ""
135 """Id of the session this one was forked from; empty when it is not a fork."""
138@dataclass(frozen=True)
139class Session:
140 """A session's metadata plus its full transcript."""
142 meta: SessionMeta
143 messages: tuple[SessionMessage, ...]
144 # Rolling summary of the turns compaction has folded away, empty until the
145 # conversation first outgrows the prompt budget. It lives here rather than on
146 # the meta because only replaying a session needs it: listing does not, and
147 # carrying a paragraph per session would bloat the drawer's hot path and
148 # every HTTP/MCP list payload.
149 summary: str = ""
152class SessionNotFoundError(Exception):
153 """Raised when a session id has no backing file."""
155 def __init__(self, session_id: str) -> None:
156 super().__init__(f"No session with id {session_id!r}")
157 self.session_id = session_id
160def _may_append(surface: SessionOrigin, owner: SessionOrigin) -> bool:
161 """Whether *surface* may append to a session owned by *owner*.
163 The human surfaces are one conversation space (the same person in the TUI,
164 Obsidian, or the shell), so they append to each other's sessions freely.
165 Agent sessions are working state: only the agent surface appends to them,
166 and it appends to nothing else without an explicit claim.
167 """
168 if surface is owner:
169 return True
170 return surface in HUMAN_ORIGINS and owner in HUMAN_ORIGINS
173class SessionOwnershipError(Exception):
174 """Raised when a surface appends to a session another surface owns."""
176 def __init__(self, session_id: str, owner: SessionOrigin, surface: SessionOrigin) -> None:
177 super().__init__(
178 f"Session {session_id!r} belongs to the {owner.value} surface; "
179 f"claim it before appending from {surface.value}."
180 )
181 self.session_id = session_id
182 self.owner = owner
183 self.surface = surface
186class SessionForkRangeError(ValueError):
187 """Raised when a fork point lies outside the source session's transcript."""
189 def __init__(self, session_id: str, message_count: int, available: int) -> None:
190 super().__init__(
191 f"Cannot fork session {session_id!r} after {message_count} messages: "
192 f"choose a count from 0 to {available}."
193 )
194 self.session_id = session_id
195 self.message_count = message_count
196 self.available = available
199def derive_title(text: str) -> str:
200 """Title a session from its first user message: first line, truncated."""
201 stripped = text.strip()
202 if not stripped:
203 return UNTITLED_SESSION_TITLE
204 first = stripped.splitlines()[0]
205 if len(first) > TITLE_MAX_LEN:
206 return first[:TITLE_MAX_LEN] + TITLE_ELLIPSIS
207 return first
210def _fork_title(source_title: str, number: int) -> str:
211 """``<source title> (fork N)``, the source part clipped so the whole fits TITLE_MAX_LEN."""
212 suffix = FORK_TITLE_SUFFIX.format(number=number)
213 room = TITLE_MAX_LEN - len(suffix)
214 if len(source_title) <= room:
215 return source_title + suffix
216 return source_title[: room - len(TITLE_ELLIPSIS)] + TITLE_ELLIPSIS + suffix
219def _fork_count(source: Session, message_count: int | None) -> int:
220 """The number of leading messages to copy, checked against *source*'s transcript."""
221 available = len(source.messages)
222 if message_count is None:
223 return available
224 if not 0 <= message_count <= available:
225 raise SessionForkRangeError(source.meta.id, message_count, available)
226 return message_count
229def _message_event(message: SessionMessage, ts: str) -> dict[str, Any]:
230 """The ``message`` event line for *message*, stamped *ts*."""
231 return {
232 "type": SessionEventType.MESSAGE,
233 "role": message.role,
234 "content": message.content,
235 "sources": list(message.sources),
236 "ts": ts,
237 }
240def _message_from_event(event: dict[str, Any], ts: str) -> SessionMessage:
241 """Reconstruct one message from its ``message`` event line."""
242 return SessionMessage(
243 role=MessageRole(event["role"]),
244 content=event["content"],
245 sources=tuple(event.get("sources", [])),
246 ts=ts,
247 )
250class SessionStore:
251 """Reads and appends session logs under ``cfg.data_dir/sessions``.
253 The directory is resolved late-bound from ``cfg`` on every call, so the store
254 follows a reconfigured data dir (and test isolation) without reconstruction.
255 ``clock`` is injectable for deterministic tests.
256 """
258 def __init__(self, clock: Callable[[], datetime] | None = None) -> None:
259 self._clock = clock or (lambda: datetime.now(UTC))
260 # path -> (size, mtime, meta) from the last fold of that file; see _meta_for.
261 self._meta_cache: dict[Path, tuple[int, float, SessionMeta]] = {}
263 @property
264 def _dir(self) -> Path:
265 return cfg.data_dir / SESSIONS_DIRNAME
267 def _path(self, session_id: str) -> Path:
268 return self._dir / f"{session_id}.jsonl"
270 def _now(self) -> str:
271 return self._clock().isoformat()
273 def _require(self, session_id: str) -> Path:
274 path = self._path(session_id)
275 if not path.exists():
276 raise SessionNotFoundError(session_id)
277 return path
279 @staticmethod
280 def _write_event(path: Path, event: dict[str, Any]) -> None:
281 # Per-session lock: two writers on one id (a second process, another
282 # surface) serialize instead of interleaving lines. Appends take
283 # milliseconds, so a blocked writer waits, never fails, under any
284 # realistic contention; the timeout only bounds a wedged holder.
285 with (
286 FileLock(str(path) + ".lock", timeout=_APPEND_LOCK_TIMEOUT_S, mode=OWNER_ONLY_MODE),
287 open(path, "a", encoding="utf-8", opener=private_opener) as fh,
288 ):
289 fh.write(json.dumps(event) + "\n")
290 fh.flush()
291 os.fsync(fh.fileno())
293 @staticmethod
294 def _iter_events(path: Path) -> Iterator[dict[str, Any]]:
295 # A torn multi-byte character decodes here, outside the try that skips
296 # torn lines; replacing makes it a JSON failure that try can catch.
297 with path.open(encoding="utf-8", errors="replace") as fh:
298 for raw in fh:
299 line = raw.strip()
300 if not line:
301 continue
302 try:
303 yield json.loads(line)
304 except json.JSONDecodeError:
305 continue # torn final line; skip it
307 @staticmethod
308 def _meta_event(
309 session_id: str,
310 model_ref: str,
311 scope: str,
312 origin: SessionOrigin,
313 now: str,
314 forked_from: str,
315 ) -> dict[str, Any]:
316 return {
317 "type": SessionEventType.META,
318 "id": session_id,
319 "created_at": now,
320 "model_ref": model_ref,
321 "scope": scope,
322 "origin": origin,
323 "forked_from": forked_from,
324 "ts": now,
325 }
327 def create(self, model_ref: str, scope: str, origin: SessionOrigin = SessionOrigin.TUI) -> str:
328 """Start a new session owned by *origin* and return its id."""
329 session_id = uuid4().hex
330 ensure_private_dir(self._dir)
331 meta = self._meta_event(session_id, model_ref, scope, origin, self._now(), forked_from="")
332 self._write_event(self._path(session_id), meta)
333 return session_id
335 def fork(
336 self,
337 session_id: str,
338 *,
339 message_count: int | None = None,
340 origin: SessionOrigin = SessionOrigin.TUI,
341 ) -> str:
342 """Copy the first *message_count* messages (all when None) into a new session.
344 The fork is owned by *origin*, which must be allowed to append to the
345 source. The source is read once and never written.
346 """
347 source = self.get(session_id)
348 if not _may_append(origin, source.meta.origin):
349 raise SessionOwnershipError(session_id, source.meta.origin, origin)
350 count = _fork_count(source, message_count)
351 fork_id = uuid4().hex
352 events = self._fork_events(fork_id, source, count, origin)
353 write_private_text(self._path(fork_id), "".join(json.dumps(e) + "\n" for e in events))
354 return fork_id
356 def _fork_events(
357 self, fork_id: str, source: Session, count: int, origin: SessionOrigin
358 ) -> list[dict[str, Any]]:
359 """A fork's whole log; the title event is last so ``updated_at`` is the fork time."""
360 now = self._now()
361 meta = source.meta
362 events = [self._meta_event(fork_id, meta.model_ref, meta.scope, origin, now, meta.id)]
363 events += [_message_event(message, message.ts) for message in source.messages[:count]]
364 if count == len(source.messages) and source.summary:
365 events.append({"type": SessionEventType.SUMMARY, "summary": source.summary, "ts": now})
366 title = _fork_title(meta.title, self._fork_number(meta.id))
367 events.append(
368 {"type": SessionEventType.TITLE, "title": title, "source": TitleSource.AUTO, "ts": now}
369 )
370 return events
372 def _fork_number(self, source_id: str) -> int:
373 """1 plus the number of existing forks of *source_id*."""
374 return 1 + sum(1 for meta in self.list() if meta.forked_from == source_id)
376 def add_message(
377 self, session_id: str, message: SessionMessage, *, surface: SessionOrigin | None = None
378 ) -> None:
379 """Append one message event to an existing session.
381 With *surface* given, the append is refused unless that surface owns the
382 session (see ``transfer``); without it the caller is a library embedder
383 that manages its own store and ownership does not apply.
384 """
385 path = self._require(session_id)
386 if surface is not None:
387 meta = self._meta_for(path)
388 if meta is not None and not _may_append(surface, meta.origin):
389 raise SessionOwnershipError(session_id, meta.origin, surface)
390 self._write_event(path, _message_event(message, self._now()))
392 def transfer(self, session_id: str, origin: SessionOrigin) -> None:
393 """Append an origin event handing the session to *origin*; newest wins.
395 This is the explicit bridge between the human and agent domains: an
396 agent claims a session whose id the user handed it, and POST /claim
397 brings one back. Never implicit in an append.
398 """
399 self._write_event(
400 self._require(session_id),
401 {"type": SessionEventType.ORIGIN, "origin": origin, "ts": self._now()},
402 )
404 def set_title(self, session_id: str, title: str, source: TitleSource) -> None:
405 """Append a title event; the newest title wins on read."""
406 self._write_event(
407 self._require(session_id),
408 {"type": SessionEventType.TITLE, "title": title, "source": source, "ts": self._now()},
409 )
411 def set_summary(self, session_id: str, summary: str) -> None:
412 """Append a summary event; the newest summary wins on read.
414 Compaction folds the oldest turns into a summary once they no longer fit
415 the prompt. The messages themselves stay in the log untouched: the
416 transcript the user scrolls is always complete, and only what is fed to
417 the model is condensed.
418 """
419 self._write_event(
420 self._require(session_id),
421 {"type": SessionEventType.SUMMARY, "summary": summary, "ts": self._now()},
422 )
424 def delete(self, session_id: str) -> None:
425 """Remove a session's file."""
426 self._require(session_id).unlink()
428 def get(self, session_id: str) -> Session:
429 """Replay a session's log into its reconstructed view."""
430 meta, messages, summary = self._replay(
431 session_id, self._require(session_id), collect_messages=True
432 )
433 return Session(meta=meta, messages=messages, summary=summary)
435 def list(self, origins: frozenset[SessionOrigin] | None = None) -> list[SessionMeta]:
436 """Sessions' metadata, newest first; *origins* narrows to those surfaces.
438 Listing replays every event of every session, so it is the one hot path
439 here: the drawer runs it on open. Messages are not materialised (only
440 counted), and each file's meta is memoised against its size and mtime so
441 reopening a vault that has not changed costs one stat() per session.
442 """
443 if not self._dir.exists():
444 return []
445 ensure_private_dir(self._dir)
446 paths = list(self._dir.glob("*.jsonl"))
447 metas = [meta for meta in (self._meta_for(path) for path in paths) if meta is not None]
448 if origins is not None:
449 metas = [meta for meta in metas if meta.origin in origins]
450 # Drop cache entries for sessions that no longer exist, so a long-lived
451 # store does not pin the meta of every session ever deleted.
452 live = {path for path in paths}
453 self._meta_cache = {p: v for p, v in self._meta_cache.items() if p in live}
454 return sorted(metas, key=lambda meta: (meta.updated_at, meta.id), reverse=True)
456 def _meta_for(self, path: Path) -> SessionMeta | None:
457 """Meta for one session file, reusing the last fold when it is unchanged.
459 The log is append-only, so any new event grows the file: size plus mtime
460 is enough to notice a change. A file that grows between the stat and the
461 read is simply re-folded on the next list(), never served stale.
463 Returns None when the file goes away underneath us, which is routine: the
464 CLI or another surface can delete a session while the drawer is listing.
465 Reading it instead would raise straight out of list().
466 """
467 try:
468 stat = path.stat()
469 except OSError:
470 return None
471 cached = self._meta_cache.get(path)
472 if cached is not None and cached[0] == stat.st_size and cached[1] == stat.st_mtime:
473 return cached[2]
474 try:
475 meta = self._replay(path.stem, path, collect_messages=False)[0]
476 except OSError:
477 return None
478 self._meta_cache[path] = (stat.st_size, stat.st_mtime, meta)
479 return meta
481 def _replay(
482 self, session_id: str, path: Path, *, collect_messages: bool
483 ) -> tuple[SessionMeta, tuple[SessionMessage, ...], str]:
484 """Fold a session's event log into its meta, messages and summary.
486 ``collect_messages=False`` is for listing, which needs only the count:
487 building a SessionMessage per message across a whole vault is pure waste.
488 """
489 created_at = ""
490 model_ref = ""
491 scope = ""
492 title = UNTITLED_SESSION_TITLE
493 updated_at = ""
494 summary = ""
495 origin = SessionOrigin.TUI
496 forked_from = ""
497 message_count = 0
498 messages: list[SessionMessage] = []
499 for event in self._iter_events(path):
500 ts = event.get("ts", "")
501 updated_at = ts
502 event_type = event.get("type")
503 if event_type == SessionEventType.META:
504 created_at = event["created_at"]
505 model_ref = event["model_ref"]
506 scope = event["scope"]
507 origin = SessionOrigin(event.get("origin", SessionOrigin.TUI))
508 forked_from = event.get("forked_from", "")
509 elif event_type == SessionEventType.ORIGIN:
510 origin = SessionOrigin(event["origin"])
511 elif event_type == SessionEventType.TITLE:
512 title = event["title"]
513 elif event_type == SessionEventType.SUMMARY:
514 summary = event["summary"]
515 elif event_type == SessionEventType.MESSAGE:
516 message_count += 1
517 if collect_messages:
518 messages.append(_message_from_event(event, ts))
519 meta = SessionMeta(
520 id=session_id,
521 title=title,
522 created_at=created_at,
523 updated_at=updated_at,
524 model_ref=model_ref,
525 scope=scope,
526 message_count=message_count,
527 origin=origin,
528 forked_from=forked_from,
529 )
530 return meta, tuple(messages), summary