1use std::collections::BTreeMap;
11use std::path::{Path, PathBuf};
12use std::sync::{Arc, Mutex};
13use std::time::Duration;
14
15use async_trait::async_trait;
16use harn_sqlite::{
17 initialize_file, initialize_transient, require_snapshot_initialized, sqlite_contention,
18 SchemaVersion, SqliteContention,
19};
20use rusqlite::{
21 params, Connection, OpenFlags, OptionalExtension, Transaction, TransactionBehavior,
22};
23use uuid::Uuid;
24
25use super::event::{
26 now_ms_and_rfc3339, AppendEvent, EventId, EventSignature, SessionEventKind, StoredEvent,
27};
28use super::redaction::{
29 prepare_append_event, prepare_stored_events_for_persistence, redact_stored_events,
30};
31use super::search::{
32 combined_score, fts_literal_query, ranks, redacted_search_document,
33 redacted_search_document_parts, snippet, vector_blob, vector_from_blob, SearchHit, SearchMode,
34 SearchQuery, SearchResponse,
35};
36use super::signing::{
37 chain_root_fold, chain_root_hash, chain_root_init, compute_record_hash, re_anchor_events,
38 verify_session_chain,
39};
40use super::store::{
41 CreateSession, EventPage, ForkResult, ImportResult, ImportSession, ListFilter, ListOrder,
42 ListSortKey, ReadRange, SessionId, SessionImporter, SessionMeta, SessionStatus, SessionStore,
43 SessionType, Snapshot, SnapshotId, StoreContention, StoreError, StoreHooks, StoreResult,
44 TruncateResult, UpdateSession, VerifyReport, MAX_READ_BATCH,
45};
46
47const SCHEMA_VERSION: i64 = 5;
52const SCHEMA_NAME: &str = "session_store";
53const DEFAULT_BUSY_TIMEOUT: Duration = Duration::from_secs(5);
54const SQLITE_SCHEMA: SchemaVersion = SchemaVersion::new(SCHEMA_NAME, SCHEMA_VERSION);
55
56#[derive(Clone)]
57pub struct SqliteSessionStore {
58 conn: Arc<Mutex<Connection>>,
59 hooks: Arc<StoreHooks>,
60 path: PathBuf,
61 _snapshot: Option<Arc<tempfile::TempDir>>,
62}
63
64impl SqliteSessionStore {
65 pub fn open(path: impl AsRef<Path>) -> StoreResult<Self> {
66 Self::open_with_hooks(path, StoreHooks::default())
67 }
68
69 pub fn open_in_memory() -> StoreResult<Self> {
70 Self::open_in_memory_with_hooks(StoreHooks::default())
71 }
72
73 pub fn open_in_memory_with_hooks(hooks: StoreHooks) -> StoreResult<Self> {
75 let conn =
76 Connection::open_in_memory().map_err(|error| StoreError::Backend(error.to_string()))?;
77 Self::initialize(conn, PathBuf::from(":memory:"), hooks)
78 }
79
80 pub fn open_read_only(path: impl AsRef<Path>) -> StoreResult<Self> {
89 Self::open_read_only_with_hooks(path, StoreHooks::default())
90 }
91
92 pub fn open_read_only_with_hooks(
94 path: impl AsRef<Path>,
95 hooks: StoreHooks,
96 ) -> StoreResult<Self> {
97 let path = path.as_ref().to_path_buf();
98 let (snapshot, snapshot_path) = copy_sqlite_snapshot(&path)?;
99 let conn = Connection::open_with_flags(&snapshot_path, OpenFlags::SQLITE_OPEN_READ_ONLY)
100 .map_err(|error| StoreError::Backend(error.to_string()))?;
101 conn.pragma_update(None, "foreign_keys", "ON")
102 .map_err(|error| StoreError::Backend(error.to_string()))?;
103 require_snapshot_initialized::<StoreError>(&conn, SQLITE_SCHEMA)
104 .map_err(map_initialization)?;
105 Ok(Self {
106 conn: Arc::new(Mutex::new(conn)),
107 hooks: Arc::new(hooks),
108 path,
109 _snapshot: Some(Arc::new(snapshot)),
110 })
111 }
112
113 pub fn open_with_hooks(path: impl AsRef<Path>, hooks: StoreHooks) -> StoreResult<Self> {
114 let path = path.as_ref().to_path_buf();
115 if let Some(parent) = path.parent() {
116 if !parent.as_os_str().is_empty() {
117 std::fs::create_dir_all(parent)
118 .map_err(|error| StoreError::Backend(error.to_string()))?;
119 }
120 }
121 let conn =
122 Connection::open(&path).map_err(|error| StoreError::Backend(error.to_string()))?;
123 Self::initialize(conn, path, hooks)
124 }
125
126 fn initialize(mut conn: Connection, path: PathBuf, hooks: StoreHooks) -> StoreResult<Self> {
127 conn.pragma_update(None, "foreign_keys", "ON")
128 .map_err(|error| StoreError::Backend(error.to_string()))?;
129 let initialization = if path == Path::new(":memory:") {
130 initialize_transient(
131 &conn,
132 DEFAULT_BUSY_TIMEOUT,
133 SQLITE_SCHEMA,
134 initialize_session_schema,
135 )
136 } else {
137 initialize_file(
138 &conn,
139 DEFAULT_BUSY_TIMEOUT,
140 SQLITE_SCHEMA,
141 initialize_session_schema,
142 )
143 };
144 initialization.map_err(map_initialization)?;
145 rebuild_search_index_if_needed(&mut conn, &hooks)?;
146 Ok(Self {
147 conn: Arc::new(Mutex::new(conn)),
148 hooks: Arc::new(hooks),
149 path,
150 _snapshot: None,
151 })
152 }
153
154 pub fn path(&self) -> &Path {
155 &self.path
156 }
157
158 fn lock(&self) -> std::sync::MutexGuard<'_, Connection> {
159 self.conn.lock().unwrap_or_else(|e| e.into_inner())
160 }
161}
162
163#[derive(PartialEq, Eq)]
164struct SnapshotSource {
165 path: PathBuf,
166 len: u64,
167 modified: Option<std::time::SystemTime>,
168}
169
170fn copy_sqlite_snapshot(path: &Path) -> StoreResult<(tempfile::TempDir, PathBuf)> {
171 const MAX_ATTEMPTS: usize = 3;
172 for attempt in 0..MAX_ATTEMPTS {
173 let before = snapshot_sources(path)?;
174 let snapshot = tempfile::Builder::new()
175 .prefix("harn-session-store-inspect-")
176 .tempdir()
177 .map_err(|error| StoreError::Backend(error.to_string()))?;
178 let snapshot_path = snapshot.path().join("session-store.sqlite");
179 let mut copy_failed = false;
180 for source in &before {
181 let destination = if source.path == path {
182 snapshot_path.clone()
183 } else if source.path == sqlite_sidecar_path(path, "-wal") {
184 sqlite_sidecar_path(&snapshot_path, "-wal")
185 } else {
186 sqlite_sidecar_path(&snapshot_path, "-shm")
187 };
188 if std::fs::copy(&source.path, destination).is_err() {
189 copy_failed = true;
190 break;
191 }
192 }
193 if !copy_failed && snapshot_sources(path)? == before {
194 return Ok((snapshot, snapshot_path));
195 }
196 if attempt + 1 == MAX_ATTEMPTS {
197 return Err(StoreError::Backend(
198 "session store changed while creating a read-only snapshot".to_string(),
199 ));
200 }
201 }
202 unreachable!("snapshot attempt loop returns on success or exhaustion")
203}
204
205fn snapshot_sources(path: &Path) -> StoreResult<Vec<SnapshotSource>> {
206 let mut paths = vec![path.to_path_buf()];
207 for suffix in ["-wal", "-shm"] {
208 let sidecar = sqlite_sidecar_path(path, suffix);
209 if sidecar.exists() {
210 paths.push(sidecar);
211 }
212 }
213 paths
214 .into_iter()
215 .map(|path| {
216 let metadata =
217 std::fs::metadata(&path).map_err(|error| StoreError::Backend(error.to_string()))?;
218 Ok(SnapshotSource {
219 path,
220 len: metadata.len(),
221 modified: metadata.modified().ok(),
222 })
223 })
224 .collect()
225}
226
227pub(crate) fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
228 let mut sidecar = path.as_os_str().to_os_string();
229 sidecar.push(suffix);
230 PathBuf::from(sidecar)
231}
232
233fn initialize_session_schema(transaction: &Transaction<'_>) -> StoreResult<()> {
234 let has_legacy_schema_version = transaction
235 .query_row(
236 "SELECT EXISTS(
237 SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = 'schema_version'
238 )",
239 [],
240 |row| row.get::<_, bool>(0),
241 )
242 .map_err(map_sql)?;
243 let previous_schema_version = if has_legacy_schema_version {
244 transaction
245 .query_row("SELECT MAX(version) FROM schema_version", [], |row| {
246 row.get::<_, Option<i64>>(0)
247 })
248 .map_err(map_sql)?
249 .unwrap_or_default()
250 } else {
251 0
252 };
253 if previous_schema_version > SCHEMA_VERSION {
254 return Err(StoreError::SchemaIncompatible {
255 schema: SCHEMA_NAME.to_string(),
256 stored: previous_schema_version,
257 supported: SCHEMA_VERSION,
258 });
259 }
260
261 transaction
262 .execute_batch(
263 "CREATE TABLE IF NOT EXISTS sessions (
264 id TEXT PRIMARY KEY,
265 tenant_id TEXT,
266 persona TEXT,
267 parent_session_id TEXT,
268 title TEXT,
269 title_pinned INTEGER NOT NULL DEFAULT 0,
270 cwd TEXT,
271 model TEXT,
272 session_type TEXT,
273 project_scope TEXT,
274 usage_input INTEGER NOT NULL DEFAULT 0,
275 usage_output INTEGER NOT NULL DEFAULT 0,
276 usage_cost_usd_micros INTEGER NOT NULL DEFAULT 0,
277 created_at_ms INTEGER NOT NULL,
278 created_at TEXT NOT NULL,
279 updated_at_ms INTEGER NOT NULL,
280 updated_at TEXT NOT NULL,
281 status TEXT NOT NULL,
282 event_count INTEGER NOT NULL DEFAULT 0,
283 last_event_id INTEGER,
284 chain_root_hash TEXT,
285 closed_at_ms INTEGER,
286 closed_at TEXT,
287 soft_deleted_at_ms INTEGER,
288 ttl_seconds INTEGER,
289 tags_json TEXT NOT NULL DEFAULT '[]',
290 attributes_json TEXT NOT NULL DEFAULT '{}',
291 next_event_id INTEGER NOT NULL DEFAULT 1
292 );
293 CREATE INDEX IF NOT EXISTS sessions_tenant_created
294 ON sessions(tenant_id, created_at_ms);
295 CREATE INDEX IF NOT EXISTS sessions_status
296 ON sessions(status);
297 CREATE INDEX IF NOT EXISTS sessions_parent
298 ON sessions(parent_session_id);
299 CREATE TABLE IF NOT EXISTS session_events (
300 session_id TEXT NOT NULL,
301 event_id INTEGER NOT NULL,
302 tenant_id TEXT,
303 parent_event_id INTEGER,
304 actor TEXT,
305 kind TEXT NOT NULL,
306 custom_kind TEXT,
307 payload_json TEXT NOT NULL,
308 tags_json TEXT NOT NULL DEFAULT '[]',
309 headers_json TEXT NOT NULL DEFAULT '{}',
310 ts_ms INTEGER NOT NULL,
311 ts TEXT NOT NULL,
312 record_hash TEXT NOT NULL,
313 prev_hash TEXT,
314 signature_json TEXT,
315 PRIMARY KEY (session_id, event_id),
316 FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE
317 );
318 CREATE INDEX IF NOT EXISTS session_events_ts
319 ON session_events(session_id, ts_ms);
320 CREATE VIRTUAL TABLE IF NOT EXISTS session_events_fts USING fts5(
321 session_id UNINDEXED,
322 event_id UNINDEXED,
323 tenant_id UNINDEXED,
324 project_scope UNINDEXED,
325 text,
326 tokenize = 'unicode61 remove_diacritics 2'
327 );
328 CREATE TABLE IF NOT EXISTS session_event_vectors (
329 session_id TEXT NOT NULL,
330 event_id INTEGER NOT NULL,
331 backend TEXT NOT NULL,
332 dim INTEGER NOT NULL,
333 embedding BLOB NOT NULL,
334 PRIMARY KEY (session_id, event_id),
335 FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE
336 );
337 CREATE TABLE IF NOT EXISTS session_imports (
338 source_id TEXT PRIMARY KEY,
339 source_digest TEXT NOT NULL,
340 session_id TEXT NOT NULL,
341 event_count INTEGER NOT NULL
342 );
343 CREATE TABLE IF NOT EXISTS session_tags (
344 session_id TEXT NOT NULL,
345 tag TEXT NOT NULL,
346 PRIMARY KEY (session_id, tag),
347 FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE
348 );
349 CREATE INDEX IF NOT EXISTS session_tags_by_tag
350 ON session_tags(tag, session_id);
351 CREATE TABLE IF NOT EXISTS session_snapshots (
352 id TEXT PRIMARY KEY,
353 session_id TEXT NOT NULL,
354 captured_at_ms INTEGER NOT NULL,
355 captured_at TEXT NOT NULL,
356 body_json TEXT NOT NULL,
357 FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE
358 );",
359 )
360 .map_err(map_sql)?;
361 ensure_session_column(transaction, "title", "TEXT")?;
362 ensure_session_column(transaction, "title_pinned", "INTEGER NOT NULL DEFAULT 0")?;
365 ensure_session_column(transaction, "cwd", "TEXT")?;
366 ensure_session_column(transaction, "model", "TEXT")?;
367 ensure_session_column(transaction, "session_type", "TEXT")?;
368 ensure_session_column(transaction, "project_scope", "TEXT")?;
369 ensure_session_column(transaction, "usage_input", "INTEGER NOT NULL DEFAULT 0")?;
370 ensure_session_column(transaction, "usage_output", "INTEGER NOT NULL DEFAULT 0")?;
371 ensure_session_column(
372 transaction,
373 "usage_cost_usd_micros",
374 "INTEGER NOT NULL DEFAULT 0",
375 )?;
376 transaction
377 .execute_batch(
378 "CREATE INDEX IF NOT EXISTS sessions_project_updated
379 ON sessions(project_scope, updated_at_ms);",
380 )
381 .map_err(map_sql)?;
382 if previous_schema_version < SCHEMA_VERSION {
383 transaction
387 .execute_batch(
388 "DELETE FROM session_events
389 WHERE NOT EXISTS (SELECT 1 FROM sessions WHERE sessions.id = session_events.session_id);
390 DELETE FROM session_tags
391 WHERE NOT EXISTS (SELECT 1 FROM sessions WHERE sessions.id = session_tags.session_id);
392 DELETE FROM session_snapshots
393 WHERE NOT EXISTS (SELECT 1 FROM sessions WHERE sessions.id = session_snapshots.session_id);
394 DELETE FROM session_event_vectors
395 WHERE NOT EXISTS (SELECT 1 FROM sessions WHERE sessions.id = session_event_vectors.session_id);",
396 )
397 .map_err(map_sql)?;
398 }
399 if has_legacy_schema_version {
400 transaction
401 .execute_batch("DROP TABLE schema_version;")
402 .map_err(map_sql)?;
403 }
404 Ok(())
405}
406
407fn ensure_session_column(
408 transaction: &Transaction<'_>,
409 column: &str,
410 sql_type: &str,
411) -> StoreResult<()> {
412 let exists = transaction
413 .prepare("PRAGMA table_info(sessions)")
414 .map_err(map_sql)?
415 .query_map([], |row| row.get::<_, String>(1))
416 .map_err(map_sql)?
417 .collect::<Result<Vec<_>, _>>()
418 .map_err(map_sql)?
419 .iter()
420 .any(|name| name == column);
421 if !exists {
422 transaction
423 .execute_batch(&format!(
424 "ALTER TABLE sessions ADD COLUMN {column} {sql_type};"
425 ))
426 .map_err(map_sql)?;
427 }
428 Ok(())
429}
430
431fn map_sql(error: rusqlite::Error) -> StoreError {
432 let message = error.to_string();
433 match sqlite_contention(&error) {
434 Some(SqliteContention::Busy) => StoreError::Contention {
435 kind: StoreContention::DatabaseBusy,
436 message,
437 },
438 Some(SqliteContention::Locked) => StoreError::Contention {
439 kind: StoreContention::DatabaseLocked,
440 message,
441 },
442 None => StoreError::Backend(message),
443 }
444}
445
446fn map_initialization(error: harn_sqlite::InitializationError<StoreError>) -> StoreError {
447 match error {
448 harn_sqlite::InitializationError::NewerSchemaVersion {
449 name,
450 stored,
451 supported,
452 } => StoreError::SchemaIncompatible {
453 schema: name.to_string(),
454 stored,
455 supported,
456 },
457 harn_sqlite::InitializationError::Initialize(error) => error,
458 other => StoreError::Backend(other.to_string()),
459 }
460}
461
462fn write_transaction(conn: &mut Connection) -> StoreResult<Transaction<'_>> {
463 conn.transaction_with_behavior(TransactionBehavior::Immediate)
469 .map_err(map_sql)
470}
471
472fn map_create_sql(error: rusqlite::Error, session_id: &str) -> StoreError {
473 if matches!(
474 error,
475 rusqlite::Error::SqliteFailure(ref inner, _)
476 if inner.code == rusqlite::ErrorCode::ConstraintViolation
477 ) {
478 StoreError::AlreadyExists(session_id.to_string())
479 } else {
480 map_sql(error)
481 }
482}
483
484fn kind_to_sql(kind: &SessionEventKind) -> (String, Option<String>) {
485 match kind {
486 SessionEventKind::Custom { custom_type } => {
487 ("custom".to_string(), Some(custom_type.clone()))
488 }
489 other => (other.discriminator().to_string(), None),
490 }
491}
492
493fn kind_from_sql(kind: &str, custom_kind: Option<String>) -> StoreResult<SessionEventKind> {
494 Ok(match kind {
495 "message" => SessionEventKind::Message,
496 "tool_call" => SessionEventKind::ToolCall,
497 "tool_result" => SessionEventKind::ToolResult,
498 "plan" => SessionEventKind::Plan,
499 "compaction" => SessionEventKind::Compaction,
500 "system_reminder" => SessionEventKind::SystemReminder,
501 "hypothesis" => SessionEventKind::Hypothesis,
502 "receipt" => SessionEventKind::Receipt,
503 "reminder" => SessionEventKind::Reminder,
504 "permission_decision" => SessionEventKind::PermissionDecision,
505 "custom" => SessionEventKind::Custom {
506 custom_type: custom_kind.unwrap_or_default(),
507 },
508 other => {
509 return Err(StoreError::Backend(format!(
510 "unknown event kind '{other}' in storage"
511 )))
512 }
513 })
514}
515
516fn status_to_sql(status: SessionStatus) -> &'static str {
517 match status {
518 SessionStatus::Open => "open",
519 SessionStatus::Closed => "closed",
520 SessionStatus::SoftDeleted => "soft_deleted",
521 SessionStatus::HardDeleted => "hard_deleted",
522 }
523}
524
525fn status_from_sql(value: &str) -> StoreResult<SessionStatus> {
526 Ok(match value {
527 "open" => SessionStatus::Open,
528 "closed" => SessionStatus::Closed,
529 "soft_deleted" => SessionStatus::SoftDeleted,
530 "hard_deleted" => SessionStatus::HardDeleted,
531 other => {
532 return Err(StoreError::Backend(format!(
533 "unknown session status '{other}' in storage"
534 )))
535 }
536 })
537}
538
539fn session_type_to_sql(session_type: SessionType) -> &'static str {
540 match session_type {
541 SessionType::User => "user",
542 SessionType::Subagent => "subagent",
543 SessionType::Scheduled => "scheduled",
544 }
545}
546
547fn session_type_from_sql(value: &str) -> StoreResult<SessionType> {
548 match value {
549 "user" => Ok(SessionType::User),
550 "subagent" => Ok(SessionType::Subagent),
551 "scheduled" => Ok(SessionType::Scheduled),
552 other => Err(StoreError::Backend(format!(
553 "unknown session type '{other}' in storage"
554 ))),
555 }
556}
557
558fn insert_session_tags(conn: &Connection, session_id: &str, tags: &[String]) -> StoreResult<()> {
559 if tags.is_empty() {
560 return Ok(());
561 }
562 let mut stmt = conn
563 .prepare("INSERT OR IGNORE INTO session_tags (session_id, tag) VALUES (?1, ?2)")
564 .map_err(map_sql)?;
565 for tag in tags {
566 stmt.execute(params![session_id, tag]).map_err(map_sql)?;
567 }
568 Ok(())
569}
570
571fn insert_session(
572 conn: &Connection,
573 meta: &SessionMeta,
574 next_event_id: EventId,
575) -> StoreResult<()> {
576 let tags_json = serde_json::to_string(&meta.tags).unwrap_or_else(|_| "[]".into());
577 let attrs_json = serde_json::to_string(&meta.attributes).unwrap_or_else(|_| "{}".into());
578 conn.execute(
579 "INSERT INTO sessions (
580 id, tenant_id, persona, parent_session_id, title, cwd, model,
581 session_type, project_scope, usage_input, usage_output,
582 usage_cost_usd_micros, created_at_ms, created_at, updated_at_ms,
583 updated_at, status, event_count, last_event_id, chain_root_hash,
584 closed_at_ms, closed_at, soft_deleted_at_ms, ttl_seconds, tags_json,
585 attributes_json, next_event_id, title_pinned
586 ) VALUES (
587 ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14,
588 ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24, ?25, ?26, ?27,
589 ?28
590 )",
591 params![
592 meta.id,
593 meta.tenant_id,
594 meta.persona,
595 meta.parent_session_id,
596 meta.title,
597 meta.cwd,
598 meta.model,
599 meta.session_type.map(session_type_to_sql),
600 meta.project_scope,
601 meta.usage_input as i64,
602 meta.usage_output as i64,
603 meta.usage_cost_usd_micros as i64,
604 meta.created_at_ms,
605 meta.created_at,
606 meta.updated_at_ms,
607 meta.updated_at,
608 status_to_sql(meta.status),
609 meta.event_count as i64,
610 meta.last_event_id.map(|value| value as i64),
611 meta.chain_root_hash,
612 meta.closed_at_ms,
613 meta.closed_at,
614 meta.soft_deleted_at_ms,
615 meta.ttl_seconds.map(|value| value as i64),
616 tags_json,
617 attrs_json,
618 next_event_id as i64,
619 meta.title_pinned,
620 ],
621 )
622 .map_err(|error| map_create_sql(error, &meta.id))?;
623 insert_session_tags(conn, &meta.id, &meta.tags)?;
624 Ok(())
625}
626
627pub(crate) fn read_session_meta(
628 conn: &Connection,
629 session_id: &str,
630) -> StoreResult<(SessionMeta, EventId)> {
631 let row = conn
632 .query_row(
633 "SELECT tenant_id, persona, parent_session_id, title, cwd, model,
634 session_type, project_scope, usage_input, usage_output,
635 usage_cost_usd_micros, created_at_ms, created_at,
636 updated_at_ms, updated_at, status, event_count, last_event_id,
637 chain_root_hash, closed_at_ms, closed_at, soft_deleted_at_ms,
638 ttl_seconds, tags_json, attributes_json, next_event_id,
639 title_pinned
640 FROM sessions WHERE id = ?1",
641 params![session_id],
642 |row| {
643 Ok((
644 row.get::<_, Option<String>>(0)?,
645 row.get::<_, Option<String>>(1)?,
646 row.get::<_, Option<String>>(2)?,
647 row.get::<_, Option<String>>(3)?,
648 row.get::<_, Option<String>>(4)?,
649 row.get::<_, Option<String>>(5)?,
650 row.get::<_, Option<String>>(6)?,
651 row.get::<_, Option<String>>(7)?,
652 row.get::<_, i64>(8)?,
653 row.get::<_, i64>(9)?,
654 row.get::<_, i64>(10)?,
655 row.get::<_, i64>(11)?,
656 row.get::<_, String>(12)?,
657 row.get::<_, i64>(13)?,
658 row.get::<_, String>(14)?,
659 row.get::<_, String>(15)?,
660 row.get::<_, i64>(16)?,
661 row.get::<_, Option<i64>>(17)?,
662 row.get::<_, Option<String>>(18)?,
663 row.get::<_, Option<i64>>(19)?,
664 row.get::<_, Option<String>>(20)?,
665 row.get::<_, Option<i64>>(21)?,
666 row.get::<_, Option<i64>>(22)?,
667 row.get::<_, String>(23)?,
668 row.get::<_, String>(24)?,
669 row.get::<_, i64>(25)?,
670 row.get::<_, bool>(26)?,
671 ))
672 },
673 )
674 .optional()
675 .map_err(map_sql)?
676 .ok_or_else(|| StoreError::NotFound(session_id.to_string()))?;
677 let (
678 tenant_id,
679 persona,
680 parent_session_id,
681 title,
682 cwd,
683 model,
684 session_type,
685 project_scope,
686 usage_input,
687 usage_output,
688 usage_cost_usd_micros,
689 created_at_ms,
690 created_at,
691 updated_at_ms,
692 updated_at,
693 status,
694 event_count,
695 last_event_id,
696 chain_root_hash,
697 closed_at_ms,
698 closed_at,
699 soft_deleted_at_ms,
700 ttl_seconds,
701 tags_json,
702 attrs_json,
703 next_event_id,
704 title_pinned,
705 ) = row;
706 let tags = serde_json::from_str(&tags_json).unwrap_or_default();
707 let attributes = serde_json::from_str(&attrs_json).unwrap_or_default();
708 let meta = SessionMeta {
709 id: session_id.to_string(),
710 tenant_id,
711 persona,
712 parent_session_id,
713 title,
714 title_pinned,
715 cwd,
716 model,
717 session_type: session_type
718 .as_deref()
719 .map(session_type_from_sql)
720 .transpose()?,
721 project_scope,
722 usage_input: usage_input as u64,
723 usage_output: usage_output as u64,
724 usage_cost_usd_micros: usage_cost_usd_micros as u64,
725 created_at_ms,
726 created_at,
727 updated_at_ms,
728 updated_at,
729 status: status_from_sql(&status)?,
730 event_count: event_count as usize,
731 last_event_id: last_event_id.map(|value| value as EventId),
732 chain_root_hash,
733 closed_at_ms,
734 closed_at,
735 soft_deleted_at_ms,
736 ttl_seconds: ttl_seconds.map(|value| value as u64),
737 tags,
738 attributes,
739 };
740 Ok((meta, next_event_id as EventId))
741}
742
743fn insert_event(conn: &Connection, event: &StoredEvent) -> StoreResult<()> {
744 let (kind, custom_kind) = kind_to_sql(&event.kind);
745 let payload_json = serde_json::to_string(&event.payload)
746 .map_err(|error| StoreError::Backend(error.to_string()))?;
747 let tags_json = serde_json::to_string(&event.tags).unwrap_or_else(|_| "[]".into());
748 let headers_json = serde_json::to_string(&event.headers).unwrap_or_else(|_| "{}".into());
749 let signature_json = event
750 .signed_by
751 .as_ref()
752 .map(|sig| serde_json::to_string(sig).unwrap_or_else(|_| "null".into()));
753 conn.execute(
754 "INSERT INTO session_events (
755 session_id, event_id, tenant_id, parent_event_id, actor, kind, custom_kind,
756 payload_json, tags_json, headers_json, ts_ms, ts, record_hash, prev_hash,
757 signature_json
758 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)",
759 params![
760 event.session_id,
761 event.event_id as i64,
762 event.tenant_id,
763 event.parent_event_id.map(|value| value as i64),
764 event.actor,
765 kind,
766 custom_kind,
767 payload_json,
768 tags_json,
769 headers_json,
770 event.ts_ms,
771 event.ts,
772 event.record_hash,
773 event.prev_hash,
774 signature_json,
775 ],
776 )
777 .map_err(map_sql)?;
778 Ok(())
779}
780
781fn insert_search_rows(
782 conn: &Connection,
783 hooks: &StoreHooks,
784 meta: &SessionMeta,
785 event: &StoredEvent,
786) -> StoreResult<()> {
787 let document = redacted_search_document(hooks.redaction.as_ref(), meta, event);
788 conn.execute(
789 "INSERT INTO session_events_fts (
790 session_id, event_id, tenant_id, project_scope, text
791 ) VALUES (?1, ?2, ?3, ?4, ?5)",
792 params![
793 event.session_id,
794 event.event_id as i64,
795 event.tenant_id,
796 meta.project_scope,
797 document,
798 ],
799 )
800 .map_err(map_sql)?;
801 let vector = hooks.embedder.embed(&document);
802 if vector.len() != hooks.embedder.dim() {
803 return Err(StoreError::Backend(format!(
804 "embedding backend '{}' returned dimension {}, expected {}",
805 hooks.embedder.name(),
806 vector.len(),
807 hooks.embedder.dim()
808 )));
809 }
810 conn.execute(
811 "INSERT OR REPLACE INTO session_event_vectors (
812 session_id, event_id, backend, dim, embedding
813 ) VALUES (?1, ?2, ?3, ?4, ?5)",
814 params![
815 event.session_id,
816 event.event_id as i64,
817 hooks.embedder.name(),
818 hooks.embedder.dim() as i64,
819 vector_blob(&vector),
820 ],
821 )
822 .map_err(map_sql)?;
823 Ok(())
824}
825
826fn rebuild_search_index_if_needed(conn: &mut Connection, hooks: &StoreHooks) -> StoreResult<()> {
827 let event_count = conn
828 .query_row("SELECT COUNT(*) FROM session_events", [], |row| {
829 row.get::<_, i64>(0)
830 })
831 .map_err(map_sql)?;
832 let fts_count = conn
833 .query_row("SELECT COUNT(*) FROM session_events_fts", [], |row| {
834 row.get::<_, i64>(0)
835 })
836 .map_err(map_sql)?;
837 let compatible_vector_count = conn
838 .query_row(
839 "SELECT COUNT(*) FROM session_event_vectors
840 WHERE backend = ?1 AND dim = ?2",
841 params![hooks.embedder.name(), hooks.embedder.dim() as i64],
842 |row| row.get::<_, i64>(0),
843 )
844 .map_err(map_sql)?;
845 if event_count == fts_count && event_count == compatible_vector_count {
846 return Ok(());
847 }
848
849 let session_ids = {
850 let mut stmt = conn
851 .prepare("SELECT id FROM sessions ORDER BY id ASC")
852 .map_err(map_sql)?;
853 let ids = stmt
854 .query_map([], |row| row.get::<_, String>(0))
855 .map_err(map_sql)?
856 .collect::<Result<Vec<_>, _>>()
857 .map_err(map_sql)?;
858 ids
859 };
860 let mut rows = Vec::new();
861 for session_id in session_ids {
862 let (meta, _) = read_session_meta(conn, &session_id)?;
863 let mut events = load_all_events(conn, &session_id)?;
864 redact_stored_events(hooks, &mut events)?;
865 rows.extend(events.into_iter().map(|event| (meta.clone(), event)));
866 }
867
868 let tx = write_transaction(conn)?;
869 tx.execute("DELETE FROM session_events_fts", [])
870 .map_err(map_sql)?;
871 tx.execute("DELETE FROM session_event_vectors", [])
872 .map_err(map_sql)?;
873 for (meta, event) in &rows {
874 insert_search_rows(&tx, hooks, meta, event)?;
875 }
876 tx.commit().map_err(map_sql)
877}
878
879fn read_event(row: &rusqlite::Row) -> Result<StoredEvent, rusqlite::Error> {
880 let session_id: String = row.get(0)?;
881 let event_id: i64 = row.get(1)?;
882 let tenant_id: Option<String> = row.get(2)?;
883 let parent_event_id: Option<i64> = row.get(3)?;
884 let actor: Option<String> = row.get(4)?;
885 let kind: String = row.get(5)?;
886 let custom_kind: Option<String> = row.get(6)?;
887 let payload_json: String = row.get(7)?;
888 let tags_json: String = row.get(8)?;
889 let headers_json: String = row.get(9)?;
890 let ts_ms: i64 = row.get(10)?;
891 let ts: String = row.get(11)?;
892 let record_hash: String = row.get(12)?;
893 let prev_hash: Option<String> = row.get(13)?;
894 let signature_json: Option<String> = row.get(14)?;
895 let resolved_kind = kind_from_sql(&kind, custom_kind).map_err(|error| {
896 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, error.into())
897 })?;
898 let payload = serde_json::from_str(&payload_json).map_err(|error| {
899 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, error.into())
900 })?;
901 let tags = serde_json::from_str(&tags_json).unwrap_or_default();
902 let headers = serde_json::from_str(&headers_json).unwrap_or_default();
903 let signed_by = signature_json
904 .as_ref()
905 .map(|value| serde_json::from_str::<EventSignature>(value))
906 .transpose()
907 .map_err(|error| {
908 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, error.into())
909 })?;
910 Ok(StoredEvent {
911 event_id: event_id as EventId,
912 session_id,
913 tenant_id,
914 parent_event_id: parent_event_id.map(|value| value as EventId),
915 actor,
916 kind: resolved_kind,
917 payload,
918 tags,
919 headers,
920 ts_ms,
921 ts,
922 record_hash,
923 prev_hash,
924 signed_by,
925 })
926}
927
928fn load_all_events(conn: &Connection, session_id: &str) -> StoreResult<Vec<StoredEvent>> {
929 let mut stmt = conn
930 .prepare(
931 "SELECT session_id, event_id, tenant_id, parent_event_id, actor, kind,
932 custom_kind, payload_json, tags_json, headers_json, ts_ms, ts,
933 record_hash, prev_hash, signature_json
934 FROM session_events WHERE session_id = ?1 ORDER BY event_id ASC",
935 )
936 .map_err(map_sql)?;
937 let rows = stmt
938 .query_map(params![session_id], read_event)
939 .map_err(map_sql)?;
940 let mut out = Vec::new();
941 for row in rows {
942 out.push(row.map_err(map_sql)?);
943 }
944 Ok(out)
945}
946
947fn read_import(conn: &Connection, source_id: &str) -> StoreResult<Option<ImportResult>> {
948 conn.query_row(
949 "SELECT source_digest, session_id, event_count
950 FROM session_imports WHERE source_id = ?1",
951 params![source_id],
952 |row| {
953 Ok(ImportResult {
954 source_id: source_id.to_string(),
955 source_digest: row.get(0)?,
956 session_id: row.get(1)?,
957 event_count: row.get::<_, i64>(2)? as usize,
958 imported: false,
959 })
960 },
961 )
962 .optional()
963 .map_err(map_sql)
964}
965
966fn append_in_tx(
973 conn: &Connection,
974 hooks: &StoreHooks,
975 session_id: &str,
976 mut event: AppendEvent,
977) -> StoreResult<StoredEvent> {
978 prepare_append_event(hooks, &mut event)?;
979 let (mut meta, next_event_id) = read_session_meta(conn, session_id)?;
980 super::memory_helpers::validate_open(&meta)?;
981 if let Some(parent_event_id) = event.parent_event_id {
982 let exists: bool = conn
983 .query_row(
984 "SELECT 1 FROM session_events WHERE session_id = ?1 AND event_id = ?2",
985 params![session_id, parent_event_id as i64],
986 |_| Ok(true),
987 )
988 .optional()
989 .map_err(map_sql)?
990 .unwrap_or(false);
991 if !exists {
992 return Err(StoreError::InvalidInput(format!(
993 "parent_event_id {parent_event_id} not present in session"
994 )));
995 }
996 }
997 let prev_hash: Option<String> = conn
998 .query_row(
999 "SELECT record_hash FROM session_events
1000 WHERE session_id = ?1 ORDER BY event_id DESC LIMIT 1",
1001 params![session_id],
1002 |row| row.get(0),
1003 )
1004 .optional()
1005 .map_err(map_sql)?;
1006 let (ts_ms, ts) = now_ms_and_rfc3339();
1007 let mut stored = StoredEvent {
1008 event_id: next_event_id,
1009 session_id: session_id.to_string(),
1010 tenant_id: meta.tenant_id.clone(),
1011 parent_event_id: event.parent_event_id,
1012 actor: event.actor,
1013 kind: event.kind,
1014 payload: event.payload,
1015 tags: event.tags,
1016 headers: event.headers,
1017 ts_ms,
1018 ts: ts.clone(),
1019 record_hash: String::new(),
1020 prev_hash,
1021 signed_by: None,
1022 };
1023 stored.record_hash = compute_record_hash(&stored);
1024 if let Some(signer) = hooks.event_signer.as_ref() {
1025 stored.signed_by = Some(signer.sign_event(&stored));
1026 }
1027 insert_event(conn, &stored)?;
1028 insert_search_rows(conn, hooks, &meta, &stored)?;
1029 let prev_root = meta.chain_root_hash.clone().unwrap_or_else(chain_root_init);
1030 let chain_root = chain_root_fold(&prev_root, &stored.record_hash);
1031 meta.event_count = meta.event_count.saturating_add(1);
1032 meta.last_event_id = Some(next_event_id);
1033 meta.chain_root_hash = Some(chain_root);
1034 meta.updated_at_ms = ts_ms;
1035 meta.updated_at = ts;
1036 conn.execute(
1037 "UPDATE sessions SET event_count = ?1, last_event_id = ?2,
1038 chain_root_hash = ?3, updated_at_ms = ?4,
1039 updated_at = ?5, next_event_id = ?6 WHERE id = ?7",
1040 params![
1041 meta.event_count as i64,
1042 meta.last_event_id.map(|value| value as i64),
1043 meta.chain_root_hash,
1044 meta.updated_at_ms,
1045 meta.updated_at,
1046 (next_event_id + 1) as i64,
1047 session_id,
1048 ],
1049 )
1050 .map_err(map_sql)?;
1051 Ok(stored)
1052}
1053
1054#[path = "sqlite/operations.rs"]
1055mod operations;
1056
1057#[cfg(test)]
1058#[path = "sqlite_tests.rs"]
1059mod tests;