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 khive_storage::blob::ContentRef;
10use rusqlite::{Connection, OptionalExtension};
11use std::path::PathBuf;
12
13use crate::error::SqliteError;
14use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
15
16// =============================================================================
17// Legacy per-service migration API (preserved for backward compatibility)
18// =============================================================================
19
20/// A single legacy migration step within a `ServiceSchemaPlan`.
21pub struct Migration {
22    /// Unique identifier for this migration.
23    pub id: &'static str,
24    /// SQL to apply (forward direction).
25    pub up_sql: &'static str,
26    /// SQL to revert (optional).
27    pub down_sql: Option<&'static str>,
28    /// Optional predicate: returns true if migration was already applied
29    /// through a mechanism other than the migration tracker.
30    pub is_already_applied: Option<fn(&Connection) -> bool>,
31}
32
33/// A pack-scoped schema plan containing migrations for SQLite and Postgres.
34pub struct ServiceSchemaPlan {
35    /// Service name used as a key in the `_schema_versions` tracking table.
36    pub service: &'static str,
37    /// SQLite-specific migration steps, applied in order.
38    pub sqlite: &'static [Migration],
39    /// Postgres-specific migration steps (reserved for future use).
40    pub postgres: &'static [Migration],
41}
42
43const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
44
45/// Apply a pack-scoped schema plan, tracking each migration in `_schema_versions`.
46pub fn apply_schema_plan(conn: &Connection, plan: &ServiceSchemaPlan) -> Result<(), SqliteError> {
47    conn.execute_batch(SCHEMA_VERSION_TABLE)?;
48
49    for migration in plan.sqlite {
50        // Check if custom predicate says it's already applied
51        if let Some(check) = migration.is_already_applied {
52            if check(conn) {
53                continue;
54            }
55        }
56
57        // Check if tracked as applied
58        let already: bool = conn.query_row(
59            "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
60            rusqlite::params![plan.service, migration.id],
61            |row| row.get(0),
62        )?;
63
64        if already {
65            continue;
66        }
67
68        let tx =
69            rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
70        tx.execute_batch(migration.up_sql)?;
71
72        tx.execute(
73            "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
74            rusqlite::params![
75                plan.service,
76                migration.id,
77                chrono::Utc::now().timestamp_micros(),
78            ],
79        )?;
80        tx.commit()?;
81    }
82
83    Ok(())
84}
85
86// =============================================================================
87// Versioned migration system
88// =============================================================================
89
90/// A single forward-only schema migration.
91///
92/// Migrations are applied in order from the current DB version to the target
93/// version. Each migration runs in its own transaction; a failure rolls back
94/// that migration and leaves the DB at the prior version.
95pub struct VersionedMigration {
96    /// Monotonically increasing version number, starting at 1.
97    pub version: u32,
98    /// Short human-readable name for the migration (used in the audit table).
99    pub name: &'static str,
100    /// SQL to apply this migration. May contain multiple statements separated
101    /// by semicolons; `execute_batch` runs them all.
102    pub up: &'static str,
103}
104
105// V1: complete schema, loaded from sql/schema.sql.
106// Fresh-start repo (v0.2.8) — all schema in one migration, no incremental versions.
107const V1_UP: &str = include_str!("../sql/schema.sql");
108
109const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
110
111const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
112
113const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
114
115const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
116
117const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
118
119const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
120
121const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
122
123const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
124
125const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
126
127const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
128
129const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
130
131const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
132
133const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
134
135const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
136
137const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
138
139const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
140
141const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
142
143const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
144
145const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
146
147const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
148
149const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
150
151/// Core schema version reserved for ADR-121's attachments-first cutover.
152pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
153
154/// The latest schema version this build's migration chain produces.
155///
156/// Terminal-version assertions belong on this, not on a hardcoded number:
157/// a literal decays into a wrong claim the next time a migration is added.
158pub fn latest_schema_version() -> u32 {
159    MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
160}
161
162/// DDL for the `ann_write_log` delta table.
163///
164/// Shared between migration V11 and the belt-and-suspenders creation in
165/// `StorageBackend::vectors_for_namespace` (same pattern as
166/// [`EMBEDDING_MODELS_DDL`]): every database that hosts `vec_*` tables must
167/// also have the write log, or vector writes would fail on databases opened
168/// without `run_migrations()`. The `.sql` file is `IF NOT EXISTS`-idempotent.
169pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
170
171/// DDL for the `ann_write_log` model/kind/field-leading index (ADR-118 §"Cost
172/// bound"), shared between migration V12 and the belt-and-suspenders creation
173/// in `StorageBackend::vectors_for_namespace` for the same reason as
174/// [`ANN_WRITE_LOG_DDL`].
175pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
176
177/// Idempotent DDL for pending ANN-consumer lifecycle metadata (#1479).
178///
179/// The V18 migration additionally translates legacy zero-watermark rows once.
180/// This constant deliberately contains only idempotent DDL: vector-store open
181/// paths may execute it repeatedly and must never demote a valid active
182/// checkpoint at sequence zero back to pending.
183pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
184
185/// DDL for the `_embedding_models` registry table.
186///
187/// Shared between the V1 schema and the belt-and-suspenders creation in
188/// `StorageBackend::vectors_for_namespace`. Both sites reference this constant so
189/// the schema cannot silently diverge if the registry evolves.
190pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
191
192/// Canonical versioned migration ledger in ascending order.
193///
194/// [`run_migrations`] applies the ordinary prefix and may complete V21 through
195/// its zero-legacy-reference fast path. A legacy V20 database records V21 only
196/// when [`finalize_attachment_cutover`] commits the application-assisted
197/// cutover.
198pub const MIGRATIONS: &[VersionedMigration] = &[
199    VersionedMigration {
200        version: 1,
201        name: "initial_schema",
202        up: V1_UP,
203    },
204    VersionedMigration {
205        version: 2,
206        name: "narrow_fts_sections_update_trigger",
207        up: V2_UP,
208    },
209    VersionedMigration {
210        version: 3,
211        name: "backfill_domain_mirror_atoms",
212        up: V3_UP,
213    },
214    VersionedMigration {
215        version: 4,
216        name: "fts_consolidation",
217        up: V4_UP,
218    },
219    VersionedMigration {
220        version: 5,
221        name: "unique_comm_message_external_id",
222        up: V5_UP,
223    },
224    VersionedMigration {
225        version: 6,
226        name: "brain_retune_driver",
227        up: V6_UP,
228    },
229    VersionedMigration {
230        version: 7,
231        name: "notes_seq",
232        up: V7_UP,
233    },
234    VersionedMigration {
235        version: 8,
236        name: "notes_seq_repair",
237        up: V8_UP,
238    },
239    VersionedMigration {
240        version: 9,
241        name: "entities_name_ci_index",
242        up: V9_UP,
243    },
244    VersionedMigration {
245        version: 10,
246        name: "entities_content_ref",
247        up: V10_UP,
248    },
249    VersionedMigration {
250        version: 11,
251        name: "ann_write_log",
252        up: V11_UP,
253    },
254    VersionedMigration {
255        version: 12,
256        name: "ann_write_log_model_seq_index",
257        up: V12_UP,
258    },
259    VersionedMigration {
260        version: 13,
261        name: "list_cursor_sequences",
262        up: V13_UP,
263    },
264    VersionedMigration {
265        version: 14,
266        name: "graph_edges_id_unique",
267        up: V14_UP,
268    },
269    VersionedMigration {
270        version: 15,
271        name: "serve_ledger_attribution",
272        up: V15_UP,
273    },
274    VersionedMigration {
275        version: 16,
276        name: "gtd_dependency_cycle_guards",
277        up: V16_UP,
278    },
279    VersionedMigration {
280        version: 17,
281        name: "agents_ddl",
282        up: V17_UP,
283    },
284    VersionedMigration {
285        version: 18,
286        name: "ann_consumer_pending",
287        up: V18_UP,
288    },
289    VersionedMigration {
290        version: 19,
291        name: "list_cursor_backfill_repair",
292        up: V19_UP,
293    },
294    VersionedMigration {
295        version: 20,
296        name: "blob_gc_claims",
297        up: V20_UP,
298    },
299    VersionedMigration {
300        version: ATTACHMENT_CUTOVER_VERSION,
301        name: "attachments_first_class",
302        // V21 is coordinated rather than an unconditional SQL migration.
303        // The runner special-cases it below; exposing the stage DDL here keeps
304        // the ledger entry self-describing for migration inspection tooling.
305        up: V21_STAGE_UP,
306    },
307];
308
309/// Durable state of ADR-121's boot-gated, two-stage attachment cutover.
310#[derive(Clone, Copy, Debug, Eq, PartialEq)]
311pub enum AttachmentCutoverStatus {
312    /// V20 is current and no stage marker has been committed.
313    Pending,
314    /// Stage 1 committed; boot must finish verified pack-owned attachments.
315    Incomplete,
316    /// V21, the attachment fences, and the attachment-only schema committed.
317    Complete,
318}
319
320fn schema_object_exists(
321    conn: &Connection,
322    object_type: &str,
323    name: &str,
324) -> Result<bool, SqliteError> {
325    conn.query_row(
326        "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
327        rusqlite::params![object_type, name],
328        |row| row.get(0),
329    )
330    .map_err(Into::into)
331}
332
333fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
334    conn.query_row(
335        "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
336        rusqlite::params![table, column],
337        |row| row.get(0),
338    )
339    .map_err(Into::into)
340}
341
342fn require_attachment_schema_objects(
343    conn: &Connection,
344    objects: &[(&str, &str)],
345    phase: &str,
346) -> Result<(), SqliteError> {
347    for (object_type, name) in objects {
348        if !schema_object_exists(conn, object_type, name)? {
349            return Err(SqliteError::InvalidData(format!(
350                "attachment cutover {phase} state is missing {object_type} {name:?}"
351            )));
352        }
353    }
354    Ok(())
355}
356
357fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
358    require_attachment_schema_objects(
359        conn,
360        &[
361            ("table", "attachments"),
362            ("index", "idx_attachments_content_ref"),
363        ],
364        "incomplete",
365    )?;
366    require_legacy_attachment_fences(conn)
367}
368
369fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
370    require_attachment_schema_objects(
371        conn,
372        &[
373            ("table", "attachments"),
374            ("table", "blob_gc_claims"),
375            ("index", "idx_attachments_content_ref"),
376            ("index", "idx_blob_gc_claims_content_ref"),
377            ("trigger", "attachments_reject_claimed_blob_insert"),
378            ("trigger", "attachments_reject_claimed_blob_update"),
379        ],
380        "complete",
381    )?;
382    if schema_column_exists(conn, "entities", "content_ref")? {
383        return Err(SqliteError::InvalidData(
384            "attachment cutover is complete but entities.content_ref still exists".into(),
385        ));
386    }
387    for (object_type, name) in [
388        ("index", "idx_entities_content_ref"),
389        ("trigger", "entities_reject_claimed_blob_insert"),
390        ("trigger", "entities_reject_claimed_blob_update"),
391    ] {
392        if schema_object_exists(conn, object_type, name)? {
393            return Err(SqliteError::InvalidData(format!(
394                "attachment cutover is complete but legacy {object_type} {name:?} still exists"
395            )));
396        }
397    }
398    Ok(())
399}
400
401/// Inspect the coordinated V21 state without mutating the connection.
402///
403/// The marker and migration ledger form one state machine. Impossible pairs
404/// fail closed instead of being guessed into a resumable state.
405pub fn attachment_cutover_status(
406    conn: &Connection,
407) -> Result<AttachmentCutoverStatus, SqliteError> {
408    let version = read_schema_version(conn)?;
409    let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
410    if !marker_table {
411        if version >= ATTACHMENT_CUTOVER_VERSION {
412            return Err(SqliteError::InvalidData(format!(
413                "migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
414            )));
415        }
416        if schema_object_exists(conn, "table", "attachments")? {
417            return Err(SqliteError::InvalidData(
418                "attachments table exists without the durable attachment cutover marker".into(),
419            ));
420        }
421        return Ok(AttachmentCutoverStatus::Pending);
422    }
423
424    let marker: Option<(String, Option<i64>)> = conn
425        .query_row(
426            "SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
427            [],
428            |row| Ok((row.get(0)?, row.get(1)?)),
429        )
430        .optional()?;
431    match marker {
432        Some((state, None)) if state == "incomplete" => {
433            if version >= ATTACHMENT_CUTOVER_VERSION {
434                Err(SqliteError::InvalidData(format!(
435                    "attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
436                )))
437            } else {
438                validate_incomplete_attachment_schema(conn)?;
439                Ok(AttachmentCutoverStatus::Incomplete)
440            }
441        }
442        Some((state, Some(_))) if state == "complete" => {
443            // Later migrations (V22+) are recorded on top of a completed
444            // cutover in the normal course; only a ledger BELOW V21 beside a
445            // complete marker is an impossible pair.
446            if version >= ATTACHMENT_CUTOVER_VERSION {
447                validate_complete_attachment_schema(conn)?;
448                Ok(AttachmentCutoverStatus::Complete)
449            } else {
450                Err(SqliteError::InvalidData(format!(
451                    "attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
452                )))
453            }
454        }
455        Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
456            "invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
457        ))),
458        None => Err(SqliteError::InvalidData(
459            "attachment cutover marker table exists without its singleton row".into(),
460        )),
461    }
462}
463
464fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
465    if !schema_column_exists(conn, "entities", "content_ref")? {
466        return Err(SqliteError::InvalidData(
467            "attachment cutover requires legacy entities.content_ref until finalization".into(),
468        ));
469    }
470    for (object_type, name) in [
471        ("table", "blob_gc_claims"),
472        ("index", "idx_blob_gc_claims_content_ref"),
473        ("index", "idx_entities_content_ref"),
474        ("trigger", "entities_reject_claimed_blob_insert"),
475        ("trigger", "entities_reject_claimed_blob_update"),
476    ] {
477        if !schema_object_exists(conn, object_type, name)? {
478            return Err(SqliteError::InvalidData(format!(
479                "attachment cutover requires legacy {object_type} {name:?} until finalization"
480            )));
481        }
482    }
483    Ok(())
484}
485
486// length() and GLOB both stop scanning at an embedded NUL, so a value of 64
487// hex characters followed by a NUL and arbitrary trailing bytes would pass
488// both. Deriving the canonical byte width from the connection's own text
489// encoding (rather than assuming UTF-8) keeps this arm correct on a database
490// pinned to UTF-16 and fails closed if the probe returns something else,
491// matching the pattern already used by `validate_blob_gc_evidence`.
492fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
493    let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
494    if !(1..=4).contains(&width) {
495        return Err(SqliteError::InvalidData(format!(
496            "the text-encoding width probe returned {width}; refusing canonicality validation"
497        )));
498    }
499    Ok(width * 64)
500}
501
502fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
503    let canonical_bytes = canonical_content_ref_byte_width(conn)?;
504    let invalid: Option<String> = conn
505        .query_row(
506            "SELECT id FROM entities \
507             WHERE content_ref IS NOT NULL \
508               AND (typeof(content_ref) <> 'text' \
509                 OR length(content_ref) <> 64 \
510                 OR length(CAST(content_ref AS BLOB)) <> ?1 \
511                 OR content_ref GLOB '*[^0-9a-f]*') \
512             LIMIT 1",
513            [canonical_bytes],
514            |row| row.get(0),
515        )
516        .optional()?;
517    if let Some(id) = invalid {
518        return Err(SqliteError::InvalidData(format!(
519            "entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
520        )));
521    }
522    Ok(())
523}
524
525fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
526    let canonical_bytes = canonical_content_ref_byte_width(conn)?;
527    for (table, identity) in [
528        ("attachments", "record_uuid"),
529        ("blob_gc_claims", "root_key"),
530    ] {
531        let sql = format!(
532            "SELECT {identity} FROM {table} \
533             WHERE typeof(content_ref) <> 'text' \
534                OR length(content_ref) <> 64 \
535                OR length(CAST(content_ref AS BLOB)) <> ?1 \
536                OR content_ref GLOB '*[^0-9a-f]*' \
537             LIMIT 1"
538        );
539        let invalid: Option<String> = conn
540            .query_row(&sql, [canonical_bytes], |row| row.get(0))
541            .optional()?;
542        if let Some(owner) = invalid {
543            return Err(SqliteError::InvalidData(format!(
544                "{table}.content_ref for {identity} {owner:?} is not canonical"
545            )));
546        }
547    }
548    Ok(())
549}
550
551fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
552    let dangling: Option<(String, String)> = conn
553        .query_row(
554            "SELECT record_uuid, substrate FROM attachments AS attachment \
555             WHERE (substrate = 'entity' AND NOT EXISTS ( \
556                       SELECT 1 FROM entities WHERE id = attachment.record_uuid \
557                   )) \
558                OR (substrate = 'note' AND NOT EXISTS ( \
559                       SELECT 1 FROM notes WHERE id = attachment.record_uuid \
560                   )) \
561             LIMIT 1",
562            [],
563            |row| Ok((row.get(0)?, row.get(1)?)),
564        )
565        .optional()?;
566    if let Some((record_uuid, substrate)) = dangling {
567        return Err(SqliteError::InvalidData(format!(
568            "attachment role references absent {substrate} record {record_uuid:?}"
569        )));
570    }
571    Ok(())
572}
573
574fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
575    let conflict: Option<String> = conn
576        .query_row(
577            "SELECT entity.id FROM entities AS entity \
578             LEFT JOIN attachments AS attachment \
579               ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
580             WHERE entity.content_ref IS NOT NULL \
581               AND (attachment.record_uuid IS NULL \
582                 OR attachment.substrate <> 'entity' \
583                 OR attachment.content_ref <> entity.content_ref) \
584             LIMIT 1",
585            [],
586            |row| row.get(0),
587        )
588        .optional()?;
589    if let Some(record_uuid) = conflict {
590        return Err(SqliteError::InvalidData(format!(
591            "legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
592        )));
593    }
594    Ok(())
595}
596
597fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
598    require_legacy_attachment_fences(conn)?;
599    conn.execute_batch(V21_STAGE_UP)?;
600    validate_canonical_legacy_refs(conn)?;
601    validate_canonical_attachment_and_claim_refs(conn)?;
602
603    let conflict: Option<String> = conn
604        .query_row(
605            "SELECT entity.id FROM entities AS entity \
606             JOIN attachments AS attachment \
607               ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
608             WHERE entity.content_ref IS NOT NULL \
609               AND (attachment.substrate <> 'entity' \
610                 OR attachment.content_ref <> entity.content_ref) \
611             LIMIT 1",
612            [],
613            |row| row.get(0),
614        )
615        .optional()?;
616    if let Some(record_uuid) = conflict {
617        return Err(SqliteError::InvalidData(format!(
618            "existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
619        )));
620    }
621
622    conn.execute(
623        "INSERT INTO attachments \
624         (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
625         SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
626         FROM entities WHERE content_ref IS NOT NULL \
627         ON CONFLICT(record_uuid, role) DO NOTHING",
628        [],
629    )?;
630    validate_legacy_content_backfill(conn)?;
631
632    // The caller holds the canonical database GC owner, so every preexisting
633    // claim is abandoned. Clearing happens before the durable incomplete
634    // marker is exposed and remains inside this one transaction.
635    conn.execute("DELETE FROM blob_gc_claims", [])?;
636    conn.execute(
637        "INSERT INTO attachment_cutover_state \
638         (singleton, state, started_at, completed_at) \
639         VALUES (1, 'incomplete', ?1, NULL) \
640         ON CONFLICT(singleton) DO NOTHING",
641        [now],
642    )?;
643    Ok(())
644}
645
646/// Commit stage 1 of the coordinated V21 migration.
647///
648/// The caller must hold [`crate::stores::blob::DatabaseGcOwnerGuard`] for the
649/// canonical database before entering this function and retain it through
650/// application backfill and finalization. This function owns one IMMEDIATE
651/// SQLite transaction; a failure leaves neither its DDL nor marker visible.
652pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
653    match attachment_cutover_status(conn)? {
654        AttachmentCutoverStatus::Complete => return Ok(()),
655        AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
656    }
657    if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
658        return Err(SqliteError::InvalidData(format!(
659            "attachment cutover stage requires canonical V{} schema",
660            ATTACHMENT_CUTOVER_VERSION - 1
661        )));
662    }
663
664    let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
665    let status = attachment_cutover_status(&tx)?;
666    if status == AttachmentCutoverStatus::Complete {
667        return Ok(());
668    }
669    stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
670    tx.commit()?;
671    Ok(())
672}
673
674/// Add one host-verified pack-owned attachment during V21 stage 2.
675///
676/// This helper is deliberately transaction-neutral: it neither begins nor
677/// commits a transaction. The boot coordinator can therefore apply the full
678/// verified vector in one caller-owned IMMEDIATE transaction while retaining
679/// the canonical database GC owner. Reapplying the same role and content is
680/// idempotent; a different substrate or digest for that role fails closed.
681#[allow(clippy::too_many_arguments)]
682pub fn apply_generic_verified_attachment(
683    conn: &Connection,
684    record_uuid: &str,
685    substrate: &str,
686    role: &str,
687    content_ref: &ContentRef,
688    media_type: Option<&str>,
689    size_bytes: Option<u64>,
690    created_at: i64,
691) -> Result<(), SqliteError> {
692    if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
693        return Err(SqliteError::InvalidData(
694            "verified application attachments may only be applied while V21 cutover is incomplete"
695                .into(),
696        ));
697    }
698    if role.is_empty() || role.chars().any(char::is_control) {
699        return Err(SqliteError::InvalidData(
700            "attachment role must be non-empty and contain no control characters".into(),
701        ));
702    }
703    let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
704        SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
705    })?;
706    let owner_table = match substrate {
707        "entity" => "entities",
708        "note" => "notes",
709        other => {
710            return Err(SqliteError::InvalidData(format!(
711                "attachment substrate must be 'entity' or 'note', got {other:?}"
712            )))
713        }
714    };
715    let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
716    let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
717    if !owner_exists {
718        return Err(SqliteError::InvalidData(format!(
719            "cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
720        )));
721    }
722    let claimed: bool = conn.query_row(
723        "SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
724        [content_ref.as_str()],
725        |row| row.get(0),
726    )?;
727    if claimed {
728        return Err(SqliteError::InvalidData(format!(
729            "cannot attach claimed content_ref {} during V21 cutover",
730            content_ref.as_str()
731        )));
732    }
733
734    let changed = conn.execute(
735        "INSERT INTO attachments \
736         (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
737         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
738         ON CONFLICT(record_uuid, role) DO UPDATE SET \
739             media_type = excluded.media_type, \
740             size_bytes = excluded.size_bytes, \
741             created_at = excluded.created_at \
742         WHERE attachments.substrate = excluded.substrate \
743           AND attachments.content_ref = excluded.content_ref",
744        rusqlite::params![
745            record_uuid,
746            substrate,
747            role,
748            content_ref.as_str(),
749            media_type,
750            size_bytes,
751            created_at,
752        ],
753    )?;
754    if changed == 0 {
755        return Err(SqliteError::InvalidData(format!(
756            "attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
757        )));
758    }
759    Ok(())
760}
761
762fn finalize_attachment_cutover_on_connection(
763    conn: &Connection,
764    now: i64,
765) -> Result<(), SqliteError> {
766    require_legacy_attachment_fences(conn)?;
767    validate_canonical_legacy_refs(conn)?;
768    validate_canonical_attachment_and_claim_refs(conn)?;
769    validate_attachment_record_owners(conn)?;
770    validate_legacy_content_backfill(conn)?;
771
772    let remaining_claims: i64 =
773        conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
774    if remaining_claims != 0 {
775        return Err(SqliteError::InvalidData(format!(
776            "attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
777        )));
778    }
779
780    let uncovered_model: Option<String> = conn
781        .query_row(
782            "SELECT model.id FROM entities AS model \
783             WHERE model.entity_type = 'moodboard_model' \
784               AND model.content_ref IS NOT NULL \
785               AND NOT EXISTS ( \
786                   SELECT 1 FROM attachments AS attachment \
787                   WHERE attachment.record_uuid = model.id \
788                     AND attachment.substrate = 'entity' \
789                     AND attachment.role = 'fann-network' \
790               ) \
791             LIMIT 1",
792            [],
793            |row| row.get(0),
794        )
795        .optional()?;
796    if let Some(record_uuid) = uncovered_model {
797        return Err(SqliteError::InvalidData(format!(
798            "moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
799        )));
800    }
801
802    conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
803    conn.execute_batch(
804        "DROP TRIGGER entities_reject_claimed_blob_insert; \
805         DROP TRIGGER entities_reject_claimed_blob_update; \
806         DROP INDEX idx_entities_content_ref; \
807         ALTER TABLE entities DROP COLUMN content_ref;",
808    )?;
809    conn.execute(
810        "UPDATE attachment_cutover_state \
811         SET state = 'complete', completed_at = ?1 \
812         WHERE singleton = 1 AND state = 'incomplete'",
813        [now],
814    )?;
815    Ok(())
816}
817
818fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
819    let migration = MIGRATIONS
820        .iter()
821        .find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
822        .expect("V21 migration must be registered");
823    conn.execute(
824        "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
825        rusqlite::params![migration.version, migration.name, now],
826    )?;
827    Ok(())
828}
829
830/// Atomically switch an explicitly staged database to attachment-only V21.
831///
832/// The caller must still hold the canonical database GC owner. The exclusive
833/// transition revalidates every legacy and attachment reference, verifies
834/// moodboard model role coverage, replaces the claim fences, removes the old
835/// column, marks the cutover complete, and records V21 in one transaction.
836pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
837    if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
838        return Ok(());
839    }
840    let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
841    match attachment_cutover_status(&tx)? {
842        AttachmentCutoverStatus::Complete => return Ok(()),
843        AttachmentCutoverStatus::Pending => {
844            return Err(SqliteError::InvalidData(
845                "attachment cutover must complete stage 1 before finalization".into(),
846            ))
847        }
848        AttachmentCutoverStatus::Incomplete => {}
849    }
850    let now = chrono::Utc::now().timestamp_micros();
851    finalize_attachment_cutover_on_connection(&tx, now)?;
852    record_attachment_cutover_migration(&tx, now)?;
853    tx.commit()?;
854    Ok(())
855}
856
857/// Read the ordered migration ledger prefix without interpreting its rows.
858fn read_applied_migration_ledger(
859    conn: &Connection,
860    through_version: u32,
861) -> Result<Vec<(u32, String)>, SqliteError> {
862    let mut stmt = conn.prepare(
863        "SELECT version, name FROM _schema_migrations \
864         WHERE version <= ?1 ORDER BY version ASC",
865    )?;
866    let rows = stmt
867        .query_map([through_version], |row| {
868            Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
869        })?
870        .collect::<Result<Vec<_>, _>>()?;
871    Ok(rows)
872}
873
874/// Require the applied versions through `through_version` to be the exact
875/// contiguous canonical prefix of [`MIGRATIONS`]. A matching `MAX(version)` is
876/// insufficient: a missing middle row or foreign version can expose a
877/// materially different schema while retaining the same maximum.
878fn validate_applied_migration_versions(
879    applied: &[(u32, String)],
880    through_version: u32,
881) -> Result<(), SqliteError> {
882    let expected: Vec<&VersionedMigration> = MIGRATIONS
883        .iter()
884        .filter(|migration| migration.version <= through_version)
885        .collect();
886    let mut applied_index = 0;
887
888    for migration in expected {
889        let Some((version, applied_name)) = applied.get(applied_index) else {
890            return Err(SqliteError::InvalidData(format!(
891                "migration history is missing version {} ('{}'); the applied ledger must be \
892                 the exact contiguous canonical sequence through version {through_version}",
893                migration.version, migration.name,
894            )));
895        };
896        if *version < migration.version {
897            return Err(SqliteError::InvalidData(format!(
898                "migration history contains unknown version {version} recorded as \
899                 '{applied_name}'; the applied ledger must contain only canonical versions"
900            )));
901        }
902        if *version > migration.version {
903            return Err(SqliteError::InvalidData(format!(
904                "migration history is missing version {} ('{}'); found version {version} \
905                 next instead",
906                migration.version, migration.name,
907            )));
908        }
909        applied_index += 1;
910    }
911
912    if let Some((version, name)) = applied.get(applied_index) {
913        return Err(SqliteError::InvalidData(format!(
914            "migration history contains unknown version {version} recorded as '{name}'; \
915             the applied ledger must contain only canonical versions"
916        )));
917    }
918
919    Ok(())
920}
921
922fn validate_applied_migration_names(
923    applied: &[(u32, String)],
924    through_version: u32,
925    allow_known_v19_repairs: bool,
926) -> Result<(), SqliteError> {
927    for ((version, applied_name), migration) in applied.iter().zip(
928        MIGRATIONS
929            .iter()
930            .filter(|migration| migration.version <= through_version),
931    ) {
932        debug_assert_eq!(*version, migration.version);
933        if migration.name != applied_name.as_str() {
934            if allow_known_v19_repairs && matches!(*version, 13 | 14) {
935                continue;
936            }
937            return Err(SqliteError::InvalidData(format!(
938                "migration version {version} is recorded under name '{applied_name}', \
939                 expected '{expected}'. This database's migration history does not match \
940                 the current binary; recreate it from the current schema or repair the \
941                 specific known divergence via a dedicated migration.",
942                expected = migration.name,
943            )));
944        }
945    }
946
947    Ok(())
948}
949
950/// Confirm the complete applied ledger is the canonical prefix, including
951/// names. The only historical V13/V14 name divergence is repaired by V19
952/// before this validator runs for a pre-V19 database.
953fn validate_applied_migration_ledger(
954    conn: &Connection,
955    through_version: u32,
956) -> Result<(), SqliteError> {
957    let applied = read_applied_migration_ledger(conn, through_version)?;
958    validate_applied_migration_versions(&applied, through_version)?;
959    validate_applied_migration_names(&applied, through_version, false)
960}
961
962const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
963
964/// Read the applied schema version from an open connection **without** running
965/// migrations. Returns 0 when the `_schema_migrations` ledger is absent (an
966/// un-migrated or empty database); any other failure (BUSY, IO) propagates —
967/// collapsing it to 0 would misreport a live database as un-migrated. Never
968/// writes.
969pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
970    match conn.query_row(
971        "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
972        [],
973        |row| row.get(0),
974    ) {
975        Ok(version) => Ok(version),
976        Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
977            if msg.contains("no such table: _schema_migrations") =>
978        {
979            Ok(0)
980        }
981        Err(e) => Err(e.into()),
982    }
983}
984
985/// Open `path` read-only and report its applied schema version without creating
986/// or migrating the file. The caller must ensure `path` exists — opening a
987/// missing file read-only errors rather than creating it. This is the path used
988/// by schema-inspection commands that must not mutate the database.
989pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
990    let conn = crate::pool::open_read_only_snapshot_connection(path)?;
991    read_schema_version(&conn)
992}
993
994/// Open `path` read-only and require the exact canonical current ledger and
995/// physical cutover state, without creating files or applying migrations.
996pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
997    let conn = crate::pool::open_read_only_snapshot_connection(path)?;
998    validate_schema_is_current(&conn)
999}
1000
1001/// Require an already-open database to match this build's latest core schema
1002/// without applying migrations.
1003///
1004/// A read-only snapshot behind the current migration set cannot be repaired in
1005/// place, while a snapshot ahead of the binary may contain schema this build
1006/// does not understand. Both directions fail with an actionable diagnostic; an
1007/// exact match performs no writes.
1008pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
1009    let current_version = read_schema_version(conn)?;
1010    let latest_version = latest_schema_version();
1011
1012    if current_version < latest_version {
1013        return Err(SqliteError::InvalidData(format!(
1014            "read-only database schema version {current_version} is behind the latest known \
1015             migration {latest_version}; migrate a writable copy with this build before opening \
1016             the snapshot read-only"
1017        )));
1018    }
1019    if current_version > latest_version {
1020        return Err(SqliteError::InvalidData(format!(
1021            "read-only database schema version {current_version} is ahead of the latest known \
1022             migration {latest_version}; use a compatible newer build or recreate the snapshot"
1023        )));
1024    }
1025
1026    // Numeric equality alone is not enough: a database can carry the current
1027    // maximum version under renamed or foreign migration ledger entries while
1028    // exposing a materially different schema. Writable boot runs this same
1029    // closed-name validation in `run_migrations_locked`; snapshot inspection
1030    // must not accept a history that ordinary boot would reject merely because
1031    // it cannot repair it in place.
1032    validate_applied_migration_ledger(conn, current_version)?;
1033    // `>=`, not `==`: later migrations (V22+) record on top of a completed
1034    // cutover, and a ledger at the latest version must not exempt the
1035    // physical cutover state from validation.
1036    if current_version >= ATTACHMENT_CUTOVER_VERSION
1037        && attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
1038    {
1039        return Err(SqliteError::InvalidData(
1040            "read-only database has not completed the V21 attachment cutover".into(),
1041        ));
1042    }
1043
1044    Ok(current_version)
1045}
1046
1047#[cfg(test)]
1048pub(crate) mod test_sync {
1049    use std::sync::atomic::AtomicU32;
1050    use std::sync::{Arc, Barrier, Mutex};
1051
1052    /// When set, `run_migrations_locked` parks after its initial (stale)
1053    /// ledger read until every racing thread has arrived — forcing the
1054    /// contended interleaving the concurrent-boot test asserts on.
1055    pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
1056    /// Counts entries into the under-lock sibling fast-forward branch.
1057    pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
1058    /// Set by the SQLite busy handler installed on participating connections:
1059    /// `true` means SQLite itself reported a blocked lock acquisition to the
1060    /// loser — actual contention, not merely an intended attempt.
1061    pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
1062        std::sync::atomic::AtomicBool::new(false);
1063
1064    /// Busy handler for participating test connections: records that SQLite
1065    /// observed a busy acquisition, then keeps retrying.
1066    pub(crate) fn record_busy(_count: i32) -> bool {
1067        BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
1068        std::thread::sleep(std::time::Duration::from_millis(1));
1069        true
1070    }
1071
1072    /// Set by the winner immediately before committing its first migration
1073    /// transaction — i.e. before the write lock is first released.
1074    pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
1075        std::sync::atomic::AtomicBool::new(false);
1076    /// Recorded by the loser when its first `BEGIN IMMEDIATE` returns: whether
1077    /// the winner had already committed at that moment. `true` is direct
1078    /// evidence the loser's lock acquisition blocked across the winner's held
1079    /// write lock rather than the two calls serializing by scheduler accident.
1080    pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
1081        std::sync::atomic::AtomicBool::new(false);
1082
1083    std::thread_local! {
1084        /// Opt-in flag: only threads that set this participate in the barrier,
1085        /// so unrelated tests migrating in parallel are never parked.
1086        pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
1087            const { std::cell::Cell::new(false) };
1088        /// Whether this thread has already instrumented its first BEGIN.
1089        pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
1090            const { std::cell::Cell::new(false) };
1091    }
1092}
1093
1094/// Apply the ordinary unapplied migration prefix in order.
1095///
1096/// The operation is idempotent and each ordinary migration runs in its own
1097/// transaction. V21 is the application-assisted exception: a database with no
1098/// legacy content references may complete V21 atomically here, while a legacy
1099/// database stops successfully at V20 so the async host can stage, verify, and
1100/// finalize the attachment cutover. Errors on a non-contiguous migration array,
1101/// a non-canonical applied ledger, or a failed migration.
1102fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
1103    let configured = conn.path().unwrap_or_default();
1104    let raw_path = if configured.is_empty() {
1105        conn.query_row(
1106            "SELECT file FROM pragma_database_list WHERE name = 'main'",
1107            [],
1108            |row| row.get::<_, String>(0),
1109        )?
1110    } else {
1111        configured.to_string()
1112    };
1113
1114    if raw_path.is_empty() {
1115        return Ok(None);
1116    }
1117    std::fs::canonicalize(&raw_path)
1118        .map(Some)
1119        .map_err(SqliteError::Io)
1120}
1121
1122fn validate_database_gc_owner(
1123    conn: &Connection,
1124    owner: &DatabaseGcOwnerGuard,
1125) -> Result<(), SqliteError> {
1126    let connection_path = canonical_connection_database_path(conn)?;
1127    if owner.database_path() != connection_path.as_deref() {
1128        return Err(SqliteError::InvalidData(format!(
1129            "database GC owner targets {:?}, but migration connection targets {:?}",
1130            owner.database_path(),
1131            connection_path.as_deref(),
1132        )));
1133    }
1134    Ok(())
1135}
1136
1137pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
1138    let database_path = canonical_connection_database_path(conn)?;
1139    if let Some(database_path) = database_path {
1140        // This raw API may have been handed a connection behind an opaque pool
1141        // writer guard. Never wait here and invert the canonical
1142        // owner-before-writer order; fail closed and direct production callers
1143        // to `StorageBackend::prepare_core_schema` instead.
1144        let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
1145            SqliteError::InvalidData(format!(
1146                "failed to acquire database GC owner before schema migration: {error}"
1147            ))
1148        })?;
1149        return run_migrations_with_database_gc_owner(conn, &owner);
1150    }
1151
1152    // A raw in-memory connection has no durable/cross-process GC domain. The
1153    // production in-memory backend still uses the owner-aware path below.
1154    run_migrations_with_busy_timeout(conn)
1155}
1156
1157pub(crate) fn run_migrations_with_database_gc_owner(
1158    conn: &mut Connection,
1159    owner: &DatabaseGcOwnerGuard,
1160) -> Result<u32, SqliteError> {
1161    validate_database_gc_owner(conn, owner)?;
1162    run_migrations_with_busy_timeout(conn)
1163}
1164
1165fn run_migrations_with_busy_timeout(conn: &mut Connection) -> Result<u32, SqliteError> {
1166    // Concurrent boots (multiple processes migrating the same file) contend on
1167    // the write lock below; a short hot-path busy_timeout cannot wait out a
1168    // sibling's migration. Raise-only to a 5s floor — never reduce a caller
1169    // whose configured timeout is already longer — and restore after.
1170    let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
1171    let raised = prior_busy_ms < 5_000;
1172    if raised {
1173        conn.busy_timeout(std::time::Duration::from_secs(5))?;
1174    }
1175    let result = run_migrations_locked(conn);
1176    if raised {
1177        let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
1178    }
1179    result
1180}
1181
1182fn run_migrations_locked(conn: &mut Connection) -> Result<u32, SqliteError> {
1183    conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
1184
1185    let current_version: u32 = read_schema_version(conn)?;
1186
1187    // Deterministic-contention hook: parks every caller after the stale ledger
1188    // read (no lock held) until all racing test threads have observed it, so
1189    // they are then released to compete for the IMMEDIATE write lock below.
1190    #[cfg(test)]
1191    if test_sync::PARTICIPATE.with(|p| p.get()) {
1192        // Replaces the busy_timeout raised by `run_migrations` on this test
1193        // connection: records SQLite-observed contention, then keeps retrying.
1194        conn.busy_handler(Some(test_sync::record_busy))?;
1195        let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
1196        if let Some(barrier) = barrier {
1197            barrier.wait();
1198        }
1199    }
1200
1201    // A database whose recorded version is ahead of the latest known migration
1202    // predates the consolidated V1 baseline (ADR-015) — e.g. it still carries the
1203    // pre-consolidation V2..V22 ledger — or was written by a newer build. Either
1204    // way the baseline schema would be silently skipped, leaving the process on a
1205    // stale schema. Fail loudly instead of corrupting silently.
1206    let latest_version = latest_schema_version();
1207    if current_version > latest_version {
1208        return Err(SqliteError::InvalidData(format!(
1209            "database schema version {current_version} is ahead of the latest known migration \
1210             {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
1211             written by a newer build. Recreate it from the current schema; in-place downgrade is \
1212             not supported."
1213        )));
1214    }
1215
1216    // Every writable upgrade starts from an exact canonical version sequence;
1217    // fail before applying new migrations if MAX(version) hides a missing or
1218    // foreign row. A pre-V19 database may still carry the V13/V14 name
1219    // divergence that V19 exists to repair. Name validation therefore permits
1220    // exactly those two rows before V19, while every unrelated mismatch still
1221    // fails before any new migration is applied.
1222    let applied = read_applied_migration_ledger(conn, current_version)?;
1223    validate_applied_migration_versions(&applied, current_version)?;
1224    validate_applied_migration_names(&applied, current_version, current_version < 19)?;
1225
1226    let mut applied_version = current_version;
1227    // Floor advanced when a sibling's work is observed under the write lock,
1228    // so a losing process skips the remaining already-applied migrations
1229    // without opening a transaction for each.
1230    let mut skip_through = current_version;
1231
1232    for migration in MIGRATIONS {
1233        if migration.version <= skip_through {
1234            applied_version = applied_version.max(migration.version);
1235            continue;
1236        }
1237
1238        // IMMEDIATE: take the write lock up front so concurrent boots serialize
1239        // here instead of failing mid-migration when a DEFERRED transaction
1240        // upgrades to a write.
1241        #[cfg(test)]
1242        let instrumented_first_begin = test_sync::PARTICIPATE.with(|p| p.get())
1243            && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
1244        #[cfg(test)]
1245        if instrumented_first_begin {
1246            test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
1247        }
1248        let tx = conn
1249            .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1250            .map_err(|e| SqliteError::Migration {
1251                version: migration.version,
1252                error: e.to_string(),
1253            })?;
1254
1255        // Re-check under the write lock: a sibling process may have applied
1256        // this migration (and possibly later ones) while we waited. Running
1257        // its DDL again would fail; fast-forward past everything it applied.
1258        let sibling_version: u32 = tx
1259            .query_row(
1260                "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1261                [],
1262                |row| row.get(0),
1263            )
1264            .map_err(|e| SqliteError::Migration {
1265                version: migration.version,
1266                error: e.to_string(),
1267            })?;
1268        #[cfg(test)]
1269        if instrumented_first_begin {
1270            use std::sync::atomic::Ordering::SeqCst;
1271            if sibling_version == 0 {
1272                // Winner: hold the write lock until SQLite has reported a
1273                // busy acquisition to the loser (its busy handler fired) —
1274                // proof the loser's BEGIN is actually blocked on this held
1275                // lock, not merely intended. Bounded so a regression fails
1276                // the assertion instead of hanging the test.
1277                let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1278                while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline
1279                {
1280                    std::thread::yield_now();
1281                }
1282            } else {
1283                // Loser: our first BEGIN just returned. Record whether the
1284                // winner had already committed — true means we blocked across
1285                // its held lock.
1286                test_sync::LOSER_SAW_WINNER_COMMIT
1287                    .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
1288            }
1289        }
1290
1291        // The ahead-of-latest guard above ran on a pre-lock read; a newer
1292        // build may have committed a version past ours while we waited for
1293        // the write lock. Accepting it (clamped) would return Ok on a schema
1294        // this binary does not understand — reject it the same way.
1295        if sibling_version > latest_version {
1296            return Err(SqliteError::InvalidData(format!(
1297                "database schema version {sibling_version} is ahead of the latest known \
1298                 migration {latest_version} (committed by a concurrent process while this \
1299                 one waited for the migration write lock). This build cannot run against \
1300                 the newer schema; upgrade the binary or recreate the database."
1301            )));
1302        }
1303
1304        if sibling_version >= migration.version {
1305            #[cfg(test)]
1306            test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1307            skip_through = sibling_version.min(latest_version);
1308            applied_version = applied_version.max(migration.version);
1309            continue;
1310        }
1311
1312        if migration.version == ATTACHMENT_CUTOVER_VERSION {
1313            let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
1314                version: migration.version,
1315                error: e.to_string(),
1316            })?;
1317            let legacy_refs: i64 = tx
1318                .query_row(
1319                    "SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
1320                    [],
1321                    |row| row.get(0),
1322                )
1323                .map_err(|e| SqliteError::Migration {
1324                    version: migration.version,
1325                    error: e.to_string(),
1326                })?;
1327
1328            // V21 belongs to the boot coordinator. Ordinary backend open may
1329            // finish the degenerate zero-ref case atomically, but it must not
1330            // expose a dual-source interval or eagerly stage a legacy DB.
1331            if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
1332                drop(tx);
1333                break;
1334            }
1335            if status != AttachmentCutoverStatus::Pending {
1336                return Err(SqliteError::Migration {
1337                    version: migration.version,
1338                    error: format!("unexpected attachment cutover state {status:?}"),
1339                });
1340            }
1341
1342            let now = chrono::Utc::now().timestamp_micros();
1343            stage_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1344                SqliteError::Migration {
1345                    version: migration.version,
1346                    error: e.to_string(),
1347                }
1348            })?;
1349            finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1350                SqliteError::Migration {
1351                    version: migration.version,
1352                    error: e.to_string(),
1353                }
1354            })?;
1355        } else {
1356            tx.execute_batch(migration.up)
1357                .map_err(|e| SqliteError::Migration {
1358                    version: migration.version,
1359                    error: e.to_string(),
1360                })?;
1361        }
1362
1363        // V19's repair contract includes normalizing the two known-divergent
1364        // recorded names. `_schema_migrations` is created and owned by this
1365        // runner (not by any migration file), so the normalization lives
1366        // here, in the same transaction that applies V19's SQL. Exact,
1367        // closed set — versions 13 and 14 only; any other (version, name)
1368        // mismatch still fails startup via validate_applied_migration_ledger.
1369        if migration.version == 19 {
1370            tx.execute_batch(
1371                "UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
1372                 UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
1373            )
1374            .map_err(|e| SqliteError::Migration {
1375                version: migration.version,
1376                error: e.to_string(),
1377            })?;
1378        }
1379
1380        let now = chrono::Utc::now().timestamp_micros();
1381        tx.execute(
1382            "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
1383             ON CONFLICT(version) DO NOTHING",
1384            rusqlite::params![migration.version, migration.name, now],
1385        )
1386        .map_err(|e| SqliteError::Migration {
1387            version: migration.version,
1388            error: e.to_string(),
1389        })?;
1390
1391        #[cfg(test)]
1392        if instrumented_first_begin {
1393            test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
1394        }
1395
1396        tx.commit().map_err(|e| SqliteError::Migration {
1397            version: migration.version,
1398            error: e.to_string(),
1399        })?;
1400
1401        applied_version = migration.version;
1402    }
1403
1404    // Validate again after the loop: our own commits and any under-lock
1405    // sibling fast-forward must both leave the exact canonical ledger, not
1406    // merely advance its maximum version.
1407    validate_applied_migration_ledger(conn, applied_version)?;
1408
1409    Ok(applied_version)
1410}
1411
1412#[derive(Debug)]
1413pub struct EmbeddingModelRegistryRecord {
1414    /// Vector engine name (e.g. `"paraphrase"`).
1415    pub engine_name: String,
1416    /// Model identifier (e.g. `"all-minilm-l6-v2"`).
1417    pub model_id: String,
1418    /// Canonical deduplication key combining engine and model.
1419    pub key_version: String,
1420    /// Embedding dimensionality.
1421    pub dimensions: u32,
1422    /// Lifecycle status (`"active"` or `"superseded"`).
1423    pub status: String,
1424    /// Epoch timestamp when the model was activated.
1425    pub activated_at: Option<i64>,
1426    /// Epoch timestamp when the model was superseded.
1427    pub superseded_at: Option<i64>,
1428}
1429
1430/// Query the `_embedding_models` registry.
1431///
1432/// Opens the database at `db` (defaults to `~/.khive/khive.db`) and
1433/// returns all registry rows, optionally filtered by `engine_name`.
1434/// Returns an empty vec if the database or table does not exist.
1435pub fn query_embedding_models(
1436    db: Option<&std::path::Path>,
1437    engine_filter: Option<&str>,
1438) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1439    let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
1440        std::env::var("HOME")
1441            .map(std::path::PathBuf::from)
1442            .unwrap_or_else(|_| std::path::PathBuf::from("."))
1443            .join(".khive/khive.db")
1444    });
1445    if !path.exists() {
1446        return Ok(Vec::new());
1447    }
1448    let conn = Connection::open_with_flags(
1449        path,
1450        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
1451            | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
1452            | rusqlite::OpenFlags::SQLITE_OPEN_URI,
1453    )?;
1454    query_embedding_models_conn(&conn, engine_filter)
1455}
1456
1457/// Query `_embedding_models` from an existing connection (testable without a file).
1458///
1459/// Returns an empty vec if the table does not exist.
1460pub(crate) fn query_embedding_models_conn(
1461    conn: &Connection,
1462    engine_filter: Option<&str>,
1463) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1464    let exists: bool = conn.query_row(
1465        "SELECT COUNT(*) > 0 FROM sqlite_master \
1466         WHERE type='table' AND name='_embedding_models'",
1467        [],
1468        |row| row.get(0),
1469    )?;
1470    if !exists {
1471        return Ok(Vec::new());
1472    }
1473
1474    let sql = if engine_filter.is_some() {
1475        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1476         FROM _embedding_models WHERE engine_name = ?1 \
1477         ORDER BY engine_name, activated_at IS NULL, activated_at"
1478    } else {
1479        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1480         FROM _embedding_models \
1481         ORDER BY engine_name, activated_at IS NULL, activated_at"
1482    };
1483    let mut stmt = conn.prepare(sql)?;
1484    let map_row = |row: &rusqlite::Row<'_>| {
1485        let dim_raw: i64 = row.get(3)?;
1486        let dimensions = u32::try_from(dim_raw).map_err(|_| {
1487            rusqlite::Error::FromSqlConversionFailure(
1488                3,
1489                rusqlite::types::Type::Integer,
1490                Box::new(std::io::Error::other(format!(
1491                    "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
1492                    u32::MAX,
1493                ))),
1494            )
1495        })?;
1496        Ok(EmbeddingModelRegistryRecord {
1497            engine_name: row.get(0)?,
1498            model_id: row.get(1)?,
1499            key_version: row.get(2)?,
1500            dimensions,
1501            status: row.get(4)?,
1502            activated_at: row.get(5)?,
1503            superseded_at: row.get(6)?,
1504        })
1505    };
1506
1507    if let Some(engine) = engine_filter {
1508        stmt.query_map([engine], map_row)?
1509            .collect::<Result<Vec<_>, _>>()
1510            .map_err(Into::into)
1511    } else {
1512        stmt.query_map([], map_row)?
1513            .collect::<Result<Vec<_>, _>>()
1514            .map_err(Into::into)
1515    }
1516}
1517
1518// =============================================================================
1519// Tests
1520// =============================================================================
1521
1522#[cfg(test)]
1523#[path = "migrations_tests.rs"]
1524mod tests;