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