1use rusqlite::Connection;
10
11use crate::error::SqliteError;
12
13pub struct Migration {
19 pub id: &'static str,
21 pub up_sql: &'static str,
23 pub down_sql: Option<&'static str>,
25 pub is_already_applied: Option<fn(&Connection) -> bool>,
28}
29
30pub struct ServiceSchemaPlan {
32 pub service: &'static str,
34 pub sqlite: &'static [Migration],
36 pub postgres: &'static [Migration],
38}
39
40const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
41
42pub 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 if let Some(check) = migration.is_already_applied {
49 if check(conn) {
50 continue;
51 }
52 }
53
54 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
83pub struct VersionedMigration {
93 pub version: u32,
95 pub name: &'static str,
97 pub up: &'static str,
100}
101
102const 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
136pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
144
145pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
150
151pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
157
158pub 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
244pub 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
267pub 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 pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
288 pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
290 pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
294 std::sync::atomic::AtomicBool::new(false);
295
296 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 pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
307 std::sync::atomic::AtomicBool::new(false);
308 pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
313 std::sync::atomic::AtomicBool::new(false);
314
315 std::thread_local! {
316 pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
319 const { std::cell::Cell::new(false) };
320 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 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 #[cfg(test)]
352 if test_sync::PARTICIPATE.with(|p| p.get()) {
353 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 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 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 #[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 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 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 test_sync::LOSER_SAW_WINNER_COMMIT
438 .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
439 }
440 }
441
442 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 pub engine_name: String,
500 pub model_id: String,
502 pub key_version: String,
504 pub dimensions: u32,
506 pub status: String,
508 pub activated_at: Option<i64>,
510 pub superseded_at: Option<i64>,
512}
513
514pub 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
536pub(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#[cfg(test)]
602#[path = "migrations_tests.rs"]
603mod tests;