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