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

1"""Append-only JSONL store for chat sessions. 

2 

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""" 

12 

13from __future__ import annotations 

14 

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 

24 

25from filelock import FileLock 

26 

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) 

34 

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})" 

51 

52 

53def sessions_enabled() -> bool: 

54 """True when session persistence is on for the human surfaces (default on).""" 

55 return cfg.sessions_enabled 

56 

57 

58def agent_sessions_enabled() -> bool: 

59 """True when agent (MCP) sessions are on (default off). 

60 

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 

65 

66 

67class SessionEventType(StrEnum): 

68 """The tag on every line of a session log.""" 

69 

70 META = "meta" 

71 TITLE = "title" 

72 MESSAGE = "message" 

73 SUMMARY = "summary" 

74 ORIGIN = "origin" 

75 

76 

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.""" 

81 

82 TUI = "tui" 

83 MCP = "mcp" 

84 HTTP = "http" 

85 CLI = "cli" 

86 

87 

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) 

94 

95 

96class MessageRole(StrEnum): 

97 """Author of a chat message.""" 

98 

99 USER = "user" 

100 ASSISTANT = "assistant" 

101 

102 

103class TitleSource(StrEnum): 

104 """Where a session title came from.""" 

105 

106 AUTO = "auto" 

107 CUSTOM = "custom" 

108 

109 

110@dataclass(frozen=True) 

111class SessionMessage: 

112 """One turn in a session. ``ts`` is stamped by the store on write.""" 

113 

114 role: MessageRole 

115 content: str 

116 sources: tuple[str, ...] = () 

117 ts: str = "" 

118 

119 

120@dataclass(frozen=True) 

121class SessionMeta: 

122 """Session metadata, reconstructed from the log without its message bodies.""" 

123 

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.""" 

136 

137 

138@dataclass(frozen=True) 

139class Session: 

140 """A session's metadata plus its full transcript.""" 

141 

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 = "" 

150 

151 

152class SessionNotFoundError(Exception): 

153 """Raised when a session id has no backing file.""" 

154 

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 

158 

159 

160def _may_append(surface: SessionOrigin, owner: SessionOrigin) -> bool: 

161 """Whether *surface* may append to a session owned by *owner*. 

162 

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 

171 

172 

173class SessionOwnershipError(Exception): 

174 """Raised when a surface appends to a session another surface owns.""" 

175 

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 

184 

185 

186class SessionForkRangeError(ValueError): 

187 """Raised when a fork point lies outside the source session's transcript.""" 

188 

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 

197 

198 

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 

208 

209 

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 

217 

218 

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 

227 

228 

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 } 

238 

239 

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 ) 

248 

249 

250class SessionStore: 

251 """Reads and appends session logs under ``cfg.data_dir/sessions``. 

252 

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 """ 

257 

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]] = {} 

262 

263 @property 

264 def _dir(self) -> Path: 

265 return cfg.data_dir / SESSIONS_DIRNAME 

266 

267 def _path(self, session_id: str) -> Path: 

268 return self._dir / f"{session_id}.jsonl" 

269 

270 def _now(self) -> str: 

271 return self._clock().isoformat() 

272 

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 

278 

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()) 

292 

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 

306 

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 } 

326 

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 

334 

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. 

343 

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 

355 

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 

371 

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) 

375 

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. 

380 

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())) 

391 

392 def transfer(self, session_id: str, origin: SessionOrigin) -> None: 

393 """Append an origin event handing the session to *origin*; newest wins. 

394 

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 ) 

403 

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 ) 

410 

411 def set_summary(self, session_id: str, summary: str) -> None: 

412 """Append a summary event; the newest summary wins on read. 

413 

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 ) 

423 

424 def delete(self, session_id: str) -> None: 

425 """Remove a session's file.""" 

426 self._require(session_id).unlink() 

427 

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) 

434 

435 def list(self, origins: frozenset[SessionOrigin] | None = None) -> list[SessionMeta]: 

436 """Sessions' metadata, newest first; *origins* narrows to those surfaces. 

437 

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) 

455 

456 def _meta_for(self, path: Path) -> SessionMeta | None: 

457 """Meta for one session file, reusing the last fold when it is unchanged. 

458 

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. 

462 

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 

480 

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. 

485 

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