Skip to main content

khive_db/
migrations.rs

1//! Schema migration system for the SQLite storage layer.
2//!
3//! Two APIs coexist:
4//! - **Legacy per-service migrations** (`ServiceSchemaPlan` / `apply_schema_plan`):
5//!   used by pack-scoped schemas.
6//! - **Versioned migrations** (`MIGRATIONS` / `run_migrations`): the forward-only
7//!   migration pipeline for the core tables.
8
9use rusqlite::Connection;
10
11use crate::error::SqliteError;
12
13// =============================================================================
14// Legacy per-service migration API (preserved for backward compatibility)
15// =============================================================================
16
17/// A single legacy migration step within a `ServiceSchemaPlan`.
18pub struct Migration {
19    /// Unique identifier for this migration.
20    pub id: &'static str,
21    /// SQL to apply (forward direction).
22    pub up_sql: &'static str,
23    /// SQL to revert (optional).
24    pub down_sql: Option<&'static str>,
25    /// Optional predicate: returns true if migration was already applied
26    /// through a mechanism other than the migration tracker.
27    pub is_already_applied: Option<fn(&Connection) -> bool>,
28}
29
30/// A pack-scoped schema plan containing migrations for SQLite and Postgres.
31pub struct ServiceSchemaPlan {
32    /// Service name used as a key in the `_schema_versions` tracking table.
33    pub service: &'static str,
34    /// SQLite-specific migration steps, applied in order.
35    pub sqlite: &'static [Migration],
36    /// Postgres-specific migration steps (reserved for future use).
37    pub postgres: &'static [Migration],
38}
39
40const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
41
42/// Apply a pack-scoped schema plan, tracking each migration in `_schema_versions`.
43pub fn apply_schema_plan(conn: &Connection, plan: &ServiceSchemaPlan) -> Result<(), SqliteError> {
44    conn.execute_batch(SCHEMA_VERSION_TABLE)?;
45
46    for migration in plan.sqlite {
47        // Check if custom predicate says it's already applied
48        if let Some(check) = migration.is_already_applied {
49            if check(conn) {
50                continue;
51            }
52        }
53
54        // Check if tracked as applied
55        let already: bool = conn.query_row(
56            "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
57            rusqlite::params![plan.service, migration.id],
58            |row| row.get(0),
59        )?;
60
61        if already {
62            continue;
63        }
64
65        let tx =
66            rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
67        tx.execute_batch(migration.up_sql)?;
68
69        tx.execute(
70            "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
71            rusqlite::params![
72                plan.service,
73                migration.id,
74                chrono::Utc::now().timestamp_micros(),
75            ],
76        )?;
77        tx.commit()?;
78    }
79
80    Ok(())
81}
82
83// =============================================================================
84// Versioned migration system
85// =============================================================================
86
87/// A single forward-only schema migration.
88///
89/// Migrations are applied in order from the current DB version to the target
90/// version. Each migration runs in its own transaction; a failure rolls back
91/// that migration and leaves the DB at the prior version.
92pub struct VersionedMigration {
93    /// Monotonically increasing version number, starting at 1.
94    pub version: u32,
95    /// Short human-readable name for the migration (used in the audit table).
96    pub name: &'static str,
97    /// SQL to apply this migration. May contain multiple statements separated
98    /// by semicolons; `execute_batch` runs them all.
99    pub up: &'static str,
100}
101
102// V1: complete schema, loaded from sql/schema.sql.
103// Fresh-start repo (v0.2.8) — all schema in one migration, no incremental versions.
104const V1_UP: &str = include_str!("../sql/schema.sql");
105
106const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
107
108const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
109
110const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
111
112const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
113
114const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
115
116const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
117
118const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
119
120const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
121
122const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
123
124const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
125
126const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
127
128const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
129
130const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
131
132const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
133
134const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
135
136/// DDL for the `ann_write_log` delta table.
137///
138/// Shared between migration V11 and the belt-and-suspenders creation in
139/// `StorageBackend::vectors_for_namespace` (same pattern as
140/// [`EMBEDDING_MODELS_DDL`]): every database that hosts `vec_*` tables must
141/// also have the write log, or vector writes would fail on databases opened
142/// without `run_migrations()`. The `.sql` file is `IF NOT EXISTS`-idempotent.
143pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
144
145/// DDL for the `ann_write_log` model/kind/field-leading index (ADR-118 §"Cost
146/// bound"), shared between migration V12 and the belt-and-suspenders creation
147/// in `StorageBackend::vectors_for_namespace` for the same reason as
148/// [`ANN_WRITE_LOG_DDL`].
149pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
150
151/// DDL for the `_embedding_models` registry table.
152///
153/// Shared between the V1 schema and the belt-and-suspenders creation in
154/// `StorageBackend::vectors_for_namespace`. Both sites reference this constant so
155/// the schema cannot silently diverge if the registry evolves.
156pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
157
158/// All versioned migrations in ascending order, applied by `run_migrations`.
159pub const MIGRATIONS: &[VersionedMigration] = &[
160    VersionedMigration {
161        version: 1,
162        name: "initial_schema",
163        up: V1_UP,
164    },
165    VersionedMigration {
166        version: 2,
167        name: "narrow_fts_sections_update_trigger",
168        up: V2_UP,
169    },
170    VersionedMigration {
171        version: 3,
172        name: "backfill_domain_mirror_atoms",
173        up: V3_UP,
174    },
175    VersionedMigration {
176        version: 4,
177        name: "fts_consolidation",
178        up: V4_UP,
179    },
180    VersionedMigration {
181        version: 5,
182        name: "unique_comm_message_external_id",
183        up: V5_UP,
184    },
185    VersionedMigration {
186        version: 6,
187        name: "brain_retune_driver",
188        up: V6_UP,
189    },
190    VersionedMigration {
191        version: 7,
192        name: "notes_seq",
193        up: V7_UP,
194    },
195    VersionedMigration {
196        version: 8,
197        name: "notes_seq_repair",
198        up: V8_UP,
199    },
200    VersionedMigration {
201        version: 9,
202        name: "entities_name_ci_index",
203        up: V9_UP,
204    },
205    VersionedMigration {
206        version: 10,
207        name: "entities_content_ref",
208        up: V10_UP,
209    },
210    VersionedMigration {
211        version: 11,
212        name: "ann_write_log",
213        up: V11_UP,
214    },
215    VersionedMigration {
216        version: 12,
217        name: "ann_write_log_model_seq_index",
218        up: V12_UP,
219    },
220    VersionedMigration {
221        version: 13,
222        name: "list_cursor_sequences",
223        up: V13_UP,
224    },
225    VersionedMigration {
226        version: 14,
227        name: "graph_edges_id_unique",
228        up: V14_UP,
229    },
230    VersionedMigration {
231        version: 15,
232        name: "serve_ledger_attribution",
233        up: V15_UP,
234    },
235    VersionedMigration {
236        version: 16,
237        name: "gtd_dependency_cycle_guards",
238        up: V16_UP,
239    },
240];
241
242const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
243
244/// Apply all unapplied migrations in order. Idempotent; each migration runs in its own transaction.
245/// Errors on non-contiguous version array or failed migration.
246/// Read the applied schema version from an open connection **without** running
247/// migrations. Returns 0 when the `_schema_migrations` ledger is absent (an
248/// un-migrated or empty database); any other failure (BUSY, IO) propagates —
249/// collapsing it to 0 would misreport a live database as un-migrated. Never
250/// writes.
251pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
252    match conn.query_row(
253        "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
254        [],
255        |row| row.get(0),
256    ) {
257        Ok(version) => Ok(version),
258        Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
259            if msg.contains("no such table: _schema_migrations") =>
260        {
261            Ok(0)
262        }
263        Err(e) => Err(e.into()),
264    }
265}
266
267/// Open `path` read-only and report its applied schema version without creating
268/// or migrating the file. The caller must ensure `path` exists — opening a
269/// missing file read-only errors rather than creating it. This is the path used
270/// by schema-inspection commands that must not mutate the database.
271pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
272    let conn = Connection::open_with_flags(
273        path,
274        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
275    )?;
276    read_schema_version(&conn)
277}
278
279#[cfg(test)]
280pub(crate) mod test_sync {
281    use std::sync::atomic::AtomicU32;
282    use std::sync::{Arc, Barrier, Mutex};
283
284    /// When set, `run_migrations_locked` parks after its initial (stale)
285    /// ledger read until every racing thread has arrived — forcing the
286    /// contended interleaving the concurrent-boot test asserts on.
287    pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
288    /// Counts entries into the under-lock sibling fast-forward branch.
289    pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
290    /// Set by the SQLite busy handler installed on participating connections:
291    /// `true` means SQLite itself reported a blocked lock acquisition to the
292    /// loser — actual contention, not merely an intended attempt.
293    pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
294        std::sync::atomic::AtomicBool::new(false);
295
296    /// Busy handler for participating test connections: records that SQLite
297    /// observed a busy acquisition, then keeps retrying.
298    pub(crate) fn record_busy(_count: i32) -> bool {
299        BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
300        std::thread::sleep(std::time::Duration::from_millis(1));
301        true
302    }
303
304    /// Set by the winner immediately before committing its first migration
305    /// transaction — i.e. before the write lock is first released.
306    pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
307        std::sync::atomic::AtomicBool::new(false);
308    /// Recorded by the loser when its first `BEGIN IMMEDIATE` returns: whether
309    /// the winner had already committed at that moment. `true` is direct
310    /// evidence the loser's lock acquisition blocked across the winner's held
311    /// write lock rather than the two calls serializing by scheduler accident.
312    pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
313        std::sync::atomic::AtomicBool::new(false);
314
315    std::thread_local! {
316        /// Opt-in flag: only threads that set this participate in the barrier,
317        /// so unrelated tests migrating in parallel are never parked.
318        pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
319            const { std::cell::Cell::new(false) };
320        /// Whether this thread has already instrumented its first BEGIN.
321        pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
322            const { std::cell::Cell::new(false) };
323    }
324}
325
326pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
327    // Concurrent boots (multiple processes migrating the same file) contend on
328    // the write lock below; a short hot-path busy_timeout cannot wait out a
329    // sibling's migration. Raise-only to a 5s floor — never reduce a caller
330    // whose configured timeout is already longer — and restore after.
331    let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
332    let raised = prior_busy_ms < 5_000;
333    if raised {
334        conn.busy_timeout(std::time::Duration::from_secs(5))?;
335    }
336    let result = run_migrations_locked(conn);
337    if raised {
338        let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
339    }
340    result
341}
342
343fn run_migrations_locked(conn: &mut Connection) -> Result<u32, SqliteError> {
344    conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
345
346    let current_version: u32 = read_schema_version(conn)?;
347
348    // Deterministic-contention hook: parks every caller after the stale ledger
349    // read (no lock held) until all racing test threads have observed it, so
350    // they are then released to compete for the IMMEDIATE write lock below.
351    #[cfg(test)]
352    if test_sync::PARTICIPATE.with(|p| p.get()) {
353        // Replaces the busy_timeout raised by `run_migrations` on this test
354        // connection: records SQLite-observed contention, then keeps retrying.
355        conn.busy_handler(Some(test_sync::record_busy))?;
356        let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
357        if let Some(barrier) = barrier {
358            barrier.wait();
359        }
360    }
361
362    // A database whose recorded version is ahead of the latest known migration
363    // predates the consolidated V1 baseline (ADR-015) — e.g. it still carries the
364    // pre-consolidation V2..V22 ledger — or was written by a newer build. Either
365    // way the baseline schema would be silently skipped, leaving the process on a
366    // stale schema. Fail loudly instead of corrupting silently.
367    let latest_version = MIGRATIONS.last().map(|m| m.version).unwrap_or(0);
368    if current_version > latest_version {
369        return Err(SqliteError::InvalidData(format!(
370            "database schema version {current_version} is ahead of the latest known migration \
371             {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
372             written by a newer build. Recreate it from the current schema; in-place downgrade is \
373             not supported."
374        )));
375    }
376
377    let mut applied_version = current_version;
378    // Floor advanced when a sibling's work is observed under the write lock,
379    // so a losing process skips the remaining already-applied migrations
380    // without opening a transaction for each.
381    let mut skip_through = current_version;
382
383    for migration in MIGRATIONS {
384        if migration.version <= skip_through {
385            applied_version = applied_version.max(migration.version);
386            continue;
387        }
388
389        // IMMEDIATE: take the write lock up front so concurrent boots serialize
390        // here instead of failing mid-migration when a DEFERRED transaction
391        // upgrades to a write.
392        #[cfg(test)]
393        let instrumented_first_begin = test_sync::PARTICIPATE.with(|p| p.get())
394            && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
395        #[cfg(test)]
396        if instrumented_first_begin {
397            test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
398        }
399        let tx = conn
400            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
401            .map_err(|e| SqliteError::Migration {
402                version: migration.version,
403                error: e.to_string(),
404            })?;
405
406        // Re-check under the write lock: a sibling process may have applied
407        // this migration (and possibly later ones) while we waited. Running
408        // its DDL again would fail; fast-forward past everything it applied.
409        let sibling_version: u32 = tx
410            .query_row(
411                "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
412                [],
413                |row| row.get(0),
414            )
415            .map_err(|e| SqliteError::Migration {
416                version: migration.version,
417                error: e.to_string(),
418            })?;
419        #[cfg(test)]
420        if instrumented_first_begin {
421            use std::sync::atomic::Ordering::SeqCst;
422            if sibling_version == 0 {
423                // Winner: hold the write lock until SQLite has reported a
424                // busy acquisition to the loser (its busy handler fired) —
425                // proof the loser's BEGIN is actually blocked on this held
426                // lock, not merely intended. Bounded so a regression fails
427                // the assertion instead of hanging the test.
428                let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
429                while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline
430                {
431                    std::thread::yield_now();
432                }
433            } else {
434                // Loser: our first BEGIN just returned. Record whether the
435                // winner had already committed — true means we blocked across
436                // its held lock.
437                test_sync::LOSER_SAW_WINNER_COMMIT
438                    .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
439            }
440        }
441
442        // The ahead-of-latest guard above ran on a pre-lock read; a newer
443        // build may have committed a version past ours while we waited for
444        // the write lock. Accepting it (clamped) would return Ok on a schema
445        // this binary does not understand — reject it the same way.
446        if sibling_version > latest_version {
447            return Err(SqliteError::InvalidData(format!(
448                "database schema version {sibling_version} is ahead of the latest known \
449                 migration {latest_version} (committed by a concurrent process while this \
450                 one waited for the migration write lock). This build cannot run against \
451                 the newer schema; upgrade the binary or recreate the database."
452            )));
453        }
454
455        if sibling_version >= migration.version {
456            #[cfg(test)]
457            test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
458            skip_through = sibling_version.min(latest_version);
459            applied_version = applied_version.max(migration.version);
460            continue;
461        }
462
463        tx.execute_batch(migration.up)
464            .map_err(|e| SqliteError::Migration {
465                version: migration.version,
466                error: e.to_string(),
467            })?;
468
469        let now = chrono::Utc::now().timestamp_micros();
470        tx.execute(
471            "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
472             ON CONFLICT(version) DO NOTHING",
473            rusqlite::params![migration.version, migration.name, now],
474        )
475        .map_err(|e| SqliteError::Migration {
476            version: migration.version,
477            error: e.to_string(),
478        })?;
479
480        #[cfg(test)]
481        if instrumented_first_begin {
482            test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
483        }
484
485        tx.commit().map_err(|e| SqliteError::Migration {
486            version: migration.version,
487            error: e.to_string(),
488        })?;
489
490        applied_version = migration.version;
491    }
492
493    Ok(applied_version)
494}
495
496#[derive(Debug)]
497pub struct EmbeddingModelRegistryRecord {
498    /// Vector engine name (e.g. `"paraphrase"`).
499    pub engine_name: String,
500    /// Model identifier (e.g. `"all-minilm-l6-v2"`).
501    pub model_id: String,
502    /// Canonical deduplication key combining engine and model.
503    pub key_version: String,
504    /// Embedding dimensionality.
505    pub dimensions: u32,
506    /// Lifecycle status (`"active"` or `"superseded"`).
507    pub status: String,
508    /// Epoch timestamp when the model was activated.
509    pub activated_at: Option<i64>,
510    /// Epoch timestamp when the model was superseded.
511    pub superseded_at: Option<i64>,
512}
513
514/// Query the `_embedding_models` registry.
515///
516/// Opens the database at `db` (defaults to `~/.khive/khive.db`) and
517/// returns all registry rows, optionally filtered by `engine_name`.
518/// Returns an empty vec if the database or table does not exist.
519pub fn query_embedding_models(
520    db: Option<&std::path::Path>,
521    engine_filter: Option<&str>,
522) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
523    let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
524        std::env::var("HOME")
525            .map(std::path::PathBuf::from)
526            .unwrap_or_else(|_| std::path::PathBuf::from("."))
527            .join(".khive/khive.db")
528    });
529    if !path.exists() {
530        return Ok(Vec::new());
531    }
532    let conn = Connection::open(path)?;
533    query_embedding_models_conn(&conn, engine_filter)
534}
535
536/// Query `_embedding_models` from an existing connection (testable without a file).
537///
538/// Returns an empty vec if the table does not exist.
539pub(crate) fn query_embedding_models_conn(
540    conn: &Connection,
541    engine_filter: Option<&str>,
542) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
543    let exists: bool = conn.query_row(
544        "SELECT COUNT(*) > 0 FROM sqlite_master \
545         WHERE type='table' AND name='_embedding_models'",
546        [],
547        |row| row.get(0),
548    )?;
549    if !exists {
550        return Ok(Vec::new());
551    }
552
553    let sql = if engine_filter.is_some() {
554        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
555         FROM _embedding_models WHERE engine_name = ?1 \
556         ORDER BY engine_name, activated_at IS NULL, activated_at"
557    } else {
558        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
559         FROM _embedding_models \
560         ORDER BY engine_name, activated_at IS NULL, activated_at"
561    };
562    let mut stmt = conn.prepare(sql)?;
563    let map_row = |row: &rusqlite::Row<'_>| {
564        let dim_raw: i64 = row.get(3)?;
565        let dimensions = u32::try_from(dim_raw).map_err(|_| {
566            rusqlite::Error::FromSqlConversionFailure(
567                3,
568                rusqlite::types::Type::Integer,
569                Box::new(std::io::Error::other(format!(
570                    "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
571                    u32::MAX,
572                ))),
573            )
574        })?;
575        Ok(EmbeddingModelRegistryRecord {
576            engine_name: row.get(0)?,
577            model_id: row.get(1)?,
578            key_version: row.get(2)?,
579            dimensions,
580            status: row.get(4)?,
581            activated_at: row.get(5)?,
582            superseded_at: row.get(6)?,
583        })
584    };
585
586    if let Some(engine) = engine_filter {
587        stmt.query_map([engine], map_row)?
588            .collect::<Result<Vec<_>, _>>()
589            .map_err(Into::into)
590    } else {
591        stmt.query_map([], map_row)?
592            .collect::<Result<Vec<_>, _>>()
593            .map_err(Into::into)
594    }
595}
596
597// =============================================================================
598// Tests
599// =============================================================================
600
601#[cfg(test)]
602#[path = "migrations_tests.rs"]
603mod tests;