Skip to main content

harn_session_store/
sqlite.rs

1//! SQLite-backed [`SessionStore`].
2//!
3//! Single-file durable backend suitable for self-hosted deployments and
4//! the TUI's persistent session DB. Schema versioning is intentionally
5//! minimal. File-backed databases use the shared Harn SQLite schema marker;
6//! pre-marker databases are upgraded from the original `schema_version` table
7//! in the same initialization transaction. The Postgres backend (issue #2500)
8//! follows the same shape so consumers can swap by config.
9
10use 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
44// v5 adds `sessions.title_pinned`. Adding a column has to move this number:
45// the shared initializer fast-paths out of schema setup entirely when the
46// recorded version already matches, so an unbumped column never reaches an
47// existing database.
48const 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    // Rows written before pinning existed carry no user choice, so 0 is the
247    // correct reading of their history rather than a placeholder.
248    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        // Foreign-key enforcement is connection-local and does not repair
268        // child rows orphaned by v1. This guarded cleanup runs once while
269        // upgrading from a pre-foreign-key schema.
270        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    // These operations read before they write. With SQLite's default DEFERRED
332    // transaction, two writers can both acquire read locks and then form an
333    // upgrade deadlock; SQLite intentionally skips the busy handler in that
334    // case and returns SQLITE_BUSY immediately. Acquire writer ownership at
335    // the boundary so the configured busy policy can serialize contenders.
336    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
831/// Core append logic, operating on a caller-owned connection (typically a
832/// transaction). Redacts, validates, links, signs (when an event signer
833/// is configured), inserts the event, and advances the session counters —
834/// but does **not** commit. `append` wraps this in its own transaction;
835/// `close` reuses it so the receipt insert, its signature, and the status
836/// flip all land in a single atomic transaction.
837fn 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;