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