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