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