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//!
9//! Raw migration entry points require a connection opened by rusqlite;
10//! externally owned handles wrapped with `Connection::from_handle` are not
11//! supported. They reject inherited transactions. If rollback and close cannot
12//! establish an outcome, the original connection is retired before its volume
13//! lease is released and the caller receives `WriterSettlementUnknown`.
14
15use khive_storage::blob::ContentRef;
16use rusqlite::{Connection, OptionalExtension};
17use std::path::{Path, PathBuf};
18
19use crate::error::SqliteError;
20use crate::pool::WriteAdmission;
21use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
22
23#[path = "raw_migration_settlement.rs"]
24mod raw_migration_settlement;
25use raw_migration_settlement::{RawMigrationTransactions, RawMigrationWriteUnit};
26
27/// The write side of a migration run.
28///
29/// ADR-154 section 4 admits each migration transaction on its own: the volume
30/// lease is taken before the transaction begins and released once it has
31/// settled, so a multi-version upgrade never holds the volume for the whole
32/// run. Implementations own that per-call lease and settlement.
33pub(crate) trait MigrationTransactions {
34    /// Run one write transaction, or one autocommit write, under a fresh
35    /// admission.
36    fn admitted<T>(
37        &mut self,
38        operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
39    ) -> Result<T, SqliteError>;
40}
41
42/// Captured disk-guard settings and shared volume-lock directory for a SQLite
43/// writer.
44///
45/// Construction validates configuration without reading or writing a database.
46/// Cooperating callers must choose the same absolute volume-lock directory.
47#[derive(Clone, Debug)]
48pub struct MigrationWritePolicy {
49    disk_guard: crate::EffectiveDiskGuardConfig,
50    volume_lock_dir: PathBuf,
51}
52
53impl MigrationWritePolicy {
54    pub fn new(
55        disk_guard: crate::EffectiveDiskGuardConfig,
56        volume_lock_dir: impl Into<PathBuf>,
57    ) -> Result<Self, SqliteError> {
58        disk_guard.validate()?;
59        let volume_lock_dir = volume_lock_dir.into();
60        if !volume_lock_dir.is_absolute() {
61            return Err(SqliteError::InvalidConfig(
62                "migration volume-lock directory must be absolute".to_string(),
63            ));
64        }
65        Ok(Self {
66            disk_guard,
67            volume_lock_dir,
68        })
69    }
70
71    /// Resolve both settings from the process environment. The lock directory
72    /// follows [`crate::default_volume_lock_dir`], the same per-user default the
73    /// public SQLite constructors use.
74    pub fn from_environment() -> Result<Self, SqliteError> {
75        let disk_guard = crate::DiskGuardEnvironment::capture().resolve(None, None)?;
76        Self::new(disk_guard, crate::default_volume_lock_dir()?)
77    }
78
79    pub fn disk_guard_config(&self) -> crate::EffectiveDiskGuardConfig {
80        self.disk_guard
81    }
82
83    pub fn volume_lock_dir(&self) -> &Path {
84        &self.volume_lock_dir
85    }
86}
87
88#[path = "session_identity_migration.rs"]
89mod session_identity_migration;
90
91mod memory_visibility;
92
93// =============================================================================
94// Legacy per-service migration API (preserved for backward compatibility)
95// =============================================================================
96
97/// A single legacy migration step within a `ServiceSchemaPlan`.
98pub struct Migration {
99    /// Unique identifier for this migration.
100    pub id: &'static str,
101    /// SQL to apply (forward direction).
102    pub up_sql: &'static str,
103    /// SQL to revert (optional).
104    pub down_sql: Option<&'static str>,
105    /// Optional read-only predicate: returns true if migration was already
106    /// applied through a mechanism other than the migration tracker. It runs
107    /// while this connection holds the SQLite write lock.
108    pub is_already_applied: Option<fn(&Connection) -> bool>,
109}
110
111/// A pack-scoped schema plan containing migrations for SQLite and Postgres.
112pub struct ServiceSchemaPlan {
113    /// Service name used as a key in the `_schema_versions` tracking table.
114    pub service: &'static str,
115    /// SQLite-specific migration steps, applied in order.
116    pub sqlite: &'static [Migration],
117    /// Postgres-specific migration steps (reserved for future use).
118    pub postgres: &'static [Migration],
119}
120
121const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
122
123/// Apply a pack-scoped schema plan, tracking each migration in `_schema_versions`.
124/// A mutable connection lets the wrapper retire an unsettled original handle.
125pub fn apply_schema_plan(
126    conn: &mut Connection,
127    plan: &ServiceSchemaPlan,
128) -> Result<(), SqliteError> {
129    let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
130    apply_schema_plan_with_admission(
131        &mut RawMigrationTransactions::new(conn, &admission),
132        plan,
133        &admission,
134    )
135}
136
137/// Apply a raw connection's schema plan with a captured policy and lock path.
138pub fn apply_schema_plan_with_policy(
139    conn: &mut Connection,
140    plan: &ServiceSchemaPlan,
141    policy: &MigrationWritePolicy,
142) -> Result<(), SqliteError> {
143    let admission =
144        WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
145    apply_schema_plan_with_admission(
146        &mut RawMigrationTransactions::new(conn, &admission),
147        plan,
148        &admission,
149    )
150}
151
152/// Apply a service schema plan. The tracking-table bootstrap and each
153/// migration are separate admitted write units (ADR-154 section 4).
154pub(crate) fn apply_schema_plan_with_admission(
155    writes: &mut impl MigrationTransactions,
156    plan: &ServiceSchemaPlan,
157    admission: &WriteAdmission,
158) -> Result<(), SqliteError> {
159    writes.admitted(|conn| {
160        admission.check()?;
161        conn.execute_batch(SCHEMA_VERSION_TABLE)?;
162        require_autocommit(conn, "schema-version bootstrap")
163    })?;
164
165    for migration in plan.sqlite {
166        writes.admitted(|conn| apply_service_migration(conn, plan, migration, admission))?;
167    }
168
169    Ok(())
170}
171
172/// Apply one service migration in its own IMMEDIATE transaction unless its
173/// predicate or the ledger says it has already been applied.
174fn apply_service_migration(
175    conn: &Connection,
176    plan: &ServiceSchemaPlan,
177    migration: &Migration,
178    admission: &WriteAdmission,
179) -> Result<(), SqliteError> {
180    // Serialize the admission decision with other writers. Checking the
181    // predicate or ledger before BEGIN IMMEDIATE lets a second opener see
182    // stale state and replay a migration after the first one commits.
183    let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
184    if let Err(error) = admission.check() {
185        let rollback = tx.rollback();
186        return Err(capacity_refusal_after_rollback(
187            conn,
188            rollback,
189            error,
190            "service schema migration",
191        ));
192    }
193
194    // Check if custom predicate says it's already applied
195    if let Some(check) = migration.is_already_applied {
196        if check(&tx) {
197            return Ok(());
198        }
199    }
200
201    // Check if tracked as applied
202    let already: bool = tx.query_row(
203        "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
204        rusqlite::params![plan.service, migration.id],
205        |row| row.get(0),
206    )?;
207
208    if already {
209        return Ok(());
210    }
211
212    tx.execute_batch(migration.up_sql)?;
213
214    tx.execute(
215        "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
216        rusqlite::params![
217            plan.service,
218            migration.id,
219            chrono::Utc::now().timestamp_micros(),
220        ],
221    )?;
222    tx.commit()?;
223    Ok(())
224}
225
226fn require_autocommit(conn: &Connection, operation: &str) -> Result<(), SqliteError> {
227    if conn.is_autocommit() {
228        Ok(())
229    } else {
230        Err(SqliteError::InvalidData(format!(
231            "{operation} did not return the SQLite connection to autocommit"
232        )))
233    }
234}
235
236pub(crate) fn capacity_refusal_after_rollback(
237    conn: &Connection,
238    rollback: rusqlite::Result<()>,
239    refusal: SqliteError,
240    operation: &str,
241) -> SqliteError {
242    if let Err(error) = rollback {
243        return SqliteError::InvalidData(format!(
244            "{operation} capacity refusal could not roll back: {error}; \
245             initial refusal: {refusal}"
246        ));
247    }
248    if !conn.is_autocommit() {
249        return SqliteError::InvalidData(format!(
250            "{operation} capacity refusal rolled back without restoring autocommit; \
251             initial refusal: {refusal}"
252        ));
253    }
254    refusal
255}
256
257// =============================================================================
258// Versioned migration system
259// =============================================================================
260
261/// A single forward-only schema migration.
262///
263/// Migrations are applied in order from the current DB version to the target
264/// version. Each migration runs in its own transaction; a failure rolls back
265/// that migration and leaves the DB at the prior version.
266pub struct VersionedMigration {
267    /// Monotonically increasing version number, starting at 1.
268    pub version: u32,
269    /// Short human-readable name for the migration (used in the audit table).
270    pub name: &'static str,
271    /// SQL to apply this migration. May contain multiple statements separated
272    /// by semicolons; `execute_batch` runs them all.
273    pub up: &'static str,
274}
275
276// V1: complete schema, loaded from sql/schema.sql.
277// Fresh-start repo (v0.2.8) — all schema in one migration, no incremental versions.
278const V1_UP: &str = include_str!("../sql/schema.sql");
279
280const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
281
282const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
283
284const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
285
286const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
287
288const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
289
290const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
291
292const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
293
294const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
295
296const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
297
298const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
299
300const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
301
302const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
303
304const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
305
306const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
307
308const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
309
310const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
311
312const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
313
314const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
315
316const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
317
318const V22_UP: &str = include_str!("../sql/022-notes-unread-probe-recipient.sql");
319
320const V23_UP: &str = include_str!("../sql/023-fts-record-kind.sql");
321
322const V24_UP: &str = include_str!("../sql/024-fts-rowid-map.sql");
323
324const V25_UP: &str = include_str!("../sql/025-notes-unread-probe-recipient-direction.sql");
325
326const V26_UP: &str = include_str!("../sql/026-knowledge-fts-repair.sql");
327
328const V27_UP: &str = include_str!("../sql/027-notes-hot-property-indexes.sql");
329const V28_UP: &str = include_str!("../sql/028-notes-key.sql");
330
331const V29_UP: &str = include_str!("../sql/029-note-streams.sql");
332const V30_UP: &str = include_str!("../sql/030-tool-source-mounts.sql");
333const V31_UP: &str = include_str!("../sql/031-note-versions.sql");
334const V32_UP: &str = include_str!("../sql/032-knowledge-count-indexes.sql");
335const V33_UP: &str = include_str!("../sql/033-notes-message-recipient-direction.sql");
336const V34_UP: &str = include_str!("../sql/034-notes-namespace-created.sql");
337const V35_UP: &str = include_str!("../sql/035-notes-unread-probe-recipient-type-direction.sql");
338const V36_UP: &str = include_str!("../sql/036-events-operation-attribution.sql");
339const V37_UP: &str = include_str!("../sql/037-entity-versions.sql");
340const V38_UP: &str = include_str!("../sql/038-entities-legacy-type-index.sql");
341const V39_UP: &str = include_str!("../sql/039-knowledge-cursor-indexes.sql");
342const SESSION_IDENTITY_UP: &str = include_str!("../sql/040-session-source-scope.sql");
343const SESSION_IDENTITY_MIGRATION_NAME: &str = "session_source_scoped_identity";
344const V41_UP: &str = include_str!("../sql/041-sender-transport.sql");
345const V42_UP: &str = include_str!("../sql/042-comm-external-id-channel-scope.sql");
346const V43_UP: &str = include_str!("../sql/043-vector-provenance.sql");
347const V44_COLUMNS: &str = include_str!("../sql/044-comm-outbound-due-a-columns.sql");
348const V44_UP: &str = include_str!("../sql/044-comm-outbound-due-b-index.sql");
349
350/// V44 may follow a direct-store bootstrap, whose idempotent notes DDL
351/// already supplies the two columns. Backfill before building the index so
352/// existing future retries are excluded by its deadline range immediately.
353pub(crate) fn migrate_outbound_due_key(tx: &rusqlite::Transaction<'_>) -> rusqlite::Result<()> {
354    let has_column = |name: &str| -> rusqlite::Result<bool> {
355        tx.query_row(
356            "SELECT EXISTS(SELECT 1 FROM pragma_table_info('notes') WHERE name = ?1)",
357            [name],
358            |row| row.get(0),
359        )
360    };
361    match (has_column("strict_due_key")?, has_column("due_source")?) {
362        (false, false) => tx.execute_batch(V44_COLUMNS)?,
363        (true, true) => {}
364        _ => return Err(rusqlite::Error::InvalidQuery),
365    }
366
367    let mut after_id = String::new();
368    let mut first_page = true;
369    loop {
370        let rows: Vec<(String, String)> = {
371            let comparator = if first_page { ">=" } else { ">" };
372            let mut stmt = tx.prepare(&format!(
373                "SELECT id, json_extract(properties, '$.next_attempt_at') FROM notes \
374                 WHERE id {comparator} ?1 AND json_type(properties, '$.next_attempt_at') = 'text' \
375                 ORDER BY id LIMIT 500"
376            ))?;
377            let collected = stmt
378                .query_map([&after_id], |row| Ok((row.get(0)?, row.get(1)?)))?
379                .collect::<rusqlite::Result<_>>()?;
380            collected
381        };
382        if rows.is_empty() {
383            break;
384        }
385        after_id = rows.last().expect("nonempty V44 page").0.clone();
386        first_page = false;
387        for (id, source) in rows {
388            if let Some(key) = crate::pool::strict_rfc3339_key(&source) {
389                tx.execute(
390                    "UPDATE notes SET strict_due_key = ?1, due_source = ?2 \
391                     WHERE id = ?3 AND (strict_due_key IS NOT ?1 OR due_source IS NOT ?2)",
392                    rusqlite::params![key, source, id],
393                )?;
394            }
395        }
396    }
397    tx.execute_batch(V44_UP)
398}
399
400const RECIPIENT_TRANSPORT_VERSION: u32 = 45;
401const V45_UP: &str = include_str!("../sql/045-recipient-transport.sql");
402const V46_UP: &str = include_str!("../sql/046-memory-visibility-receipts.sql");
403const V47_UP: &str = include_str!("../sql/047-attachment-role-quarantine.sql");
404const V49_UP: &str = include_str!("../sql/049-git-note-property-indexes.sql");
405
406const V50_UP: &str = include_str!("../sql/050-entity-list-plans.sql");
407const V51_UP: &str = include_str!("../sql/051-schedule-core-indexes.sql");
408const V52_UP: &str = include_str!("../sql/052-comm-core-indexes.sql");
409const V53_UP: &str = include_str!("../sql/053-entity-kind-list-order.sql");
410const MEMORY_VISIBILITY_CUTOVER_VERSION: u32 = 54;
411const V54_UP: &str = include_str!("../sql/054-memory-visibility-epochs.sql");
412const V48_UP: &str = include_str!("../sql/048-acknowledgement-journal-a-table.sql");
413const ACKNOWLEDGEMENT_JOURNAL_INDEX: &str =
414    include_str!("../sql/048-acknowledgement-journal-b-index.sql");
415
416/// A ledger-tail replay keeps already-current journal state and retry metadata.
417pub(crate) fn migrate_acknowledgement_journal(
418    tx: &rusqlite::Transaction<'_>,
419) -> rusqlite::Result<()> {
420    let columns: i64 = tx.query_row(
421        "SELECT count(*) FROM pragma_table_info('comm_ack_work') \
422         WHERE name IN ('attempt_count','not_before','retirement_reason')",
423        [],
424        |row| row.get(0),
425    )?;
426    match columns {
427        0 => tx.execute_batch(V48_UP)?,
428        3 => {}
429        _ => return Err(rusqlite::Error::InvalidQuery),
430    }
431    tx.execute_batch(ACKNOWLEDGEMENT_JOURNAL_INDEX)
432}
433
434const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
435
436const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
437
438/// Core schema version reserved for ADR-121's attachments-first cutover.
439pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
440
441/// The latest schema version this build's migration chain produces.
442///
443/// Terminal-version assertions belong on this, not on a hardcoded number:
444/// a literal decays into a wrong claim the next time a migration is added.
445pub fn latest_schema_version() -> u32 {
446    MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
447}
448
449/// DDL for the `ann_write_log` delta table.
450///
451/// Shared between migration V11 and the belt-and-suspenders creation in
452/// `StorageBackend::vectors_for_namespace` (same pattern as
453/// [`EMBEDDING_MODELS_DDL`]): every database that hosts `vec_*` tables must
454/// also have the write log, or vector writes would fail on databases opened
455/// without `run_migrations()`. The `.sql` file is `IF NOT EXISTS`-idempotent.
456pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
457
458/// DDL for the `ann_write_log` model/kind/field-leading index (ADR-118 §"Cost
459/// bound"), shared between migration V12 and the belt-and-suspenders creation
460/// in `StorageBackend::vectors_for_namespace` for the same reason as
461/// [`ANN_WRITE_LOG_DDL`].
462pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
463
464/// Idempotent DDL for pending ANN-consumer lifecycle metadata (#1479).
465///
466/// The V18 migration additionally translates legacy zero-watermark rows once.
467/// This constant deliberately contains only idempotent DDL: vector-store open
468/// paths may execute it repeatedly and must never demote a valid active
469/// checkpoint at sequence zero back to pending.
470pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
471
472/// Sidecar DDL registered in the migration ledger. Direct vector-store
473/// construction leaves this schema to the versioned migration.
474pub const VECTOR_PROVENANCE_DDL: &str = V43_UP;
475
476/// DDL for the `_embedding_models` registry table.
477///
478/// Shared between the V1 schema and the belt-and-suspenders creation in
479/// `StorageBackend::vectors_for_namespace`. Both sites reference this constant so
480/// the schema cannot silently diverge if the registry evolves.
481pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
482
483/// Canonical versioned migration ledger in ascending order.
484///
485/// [`run_migrations`] applies the ordinary prefix and may complete V21 through
486/// its zero-legacy-reference fast path. A legacy V20 database records V21 only
487/// when [`finalize_attachment_cutover`] commits the application-assisted
488/// cutover.
489pub const MIGRATIONS: &[VersionedMigration] = &[
490    VersionedMigration {
491        version: 1,
492        name: "initial_schema",
493        up: V1_UP,
494    },
495    VersionedMigration {
496        version: 2,
497        name: "narrow_fts_sections_update_trigger",
498        up: V2_UP,
499    },
500    VersionedMigration {
501        version: 3,
502        name: "backfill_domain_mirror_atoms",
503        up: V3_UP,
504    },
505    VersionedMigration {
506        version: 4,
507        name: "fts_consolidation",
508        up: V4_UP,
509    },
510    VersionedMigration {
511        version: 5,
512        name: "unique_comm_message_external_id",
513        up: V5_UP,
514    },
515    VersionedMigration {
516        version: 6,
517        name: "brain_retune_driver",
518        up: V6_UP,
519    },
520    VersionedMigration {
521        version: 7,
522        name: "notes_seq",
523        up: V7_UP,
524    },
525    VersionedMigration {
526        version: 8,
527        name: "notes_seq_repair",
528        up: V8_UP,
529    },
530    VersionedMigration {
531        version: 9,
532        name: "entities_name_ci_index",
533        up: V9_UP,
534    },
535    VersionedMigration {
536        version: 10,
537        name: "entities_content_ref",
538        up: V10_UP,
539    },
540    VersionedMigration {
541        version: 11,
542        name: "ann_write_log",
543        up: V11_UP,
544    },
545    VersionedMigration {
546        version: 12,
547        name: "ann_write_log_model_seq_index",
548        up: V12_UP,
549    },
550    VersionedMigration {
551        version: 13,
552        name: "list_cursor_sequences",
553        up: V13_UP,
554    },
555    VersionedMigration {
556        version: 14,
557        name: "graph_edges_id_unique",
558        up: V14_UP,
559    },
560    VersionedMigration {
561        version: 15,
562        name: "serve_ledger_attribution",
563        up: V15_UP,
564    },
565    VersionedMigration {
566        version: 16,
567        name: "gtd_dependency_cycle_guards",
568        up: V16_UP,
569    },
570    VersionedMigration {
571        version: 17,
572        name: "agents_ddl",
573        up: V17_UP,
574    },
575    VersionedMigration {
576        version: 18,
577        name: "ann_consumer_pending",
578        up: V18_UP,
579    },
580    VersionedMigration {
581        version: 19,
582        name: "list_cursor_backfill_repair",
583        up: V19_UP,
584    },
585    VersionedMigration {
586        version: 20,
587        name: "blob_gc_claims",
588        up: V20_UP,
589    },
590    VersionedMigration {
591        version: ATTACHMENT_CUTOVER_VERSION,
592        name: "attachments_first_class",
593        // V21 is coordinated rather than an unconditional SQL migration.
594        // The runner special-cases it below; exposing the stage DDL here keeps
595        // the ledger entry self-describing for migration inspection tooling.
596        up: V21_STAGE_UP,
597    },
598    VersionedMigration {
599        version: 22,
600        name: "notes_unread_probe_recipient",
601        up: V22_UP,
602    },
603    VersionedMigration {
604        version: 23,
605        name: "fts_record_kind",
606        up: V23_UP,
607    },
608    VersionedMigration {
609        version: 24,
610        name: "fts_rowid_map",
611        up: V24_UP,
612    },
613    VersionedMigration {
614        version: 25,
615        name: "notes_unread_probe_recipient_direction",
616        up: V25_UP,
617    },
618    VersionedMigration {
619        version: 26,
620        name: "knowledge_fts_repair",
621        up: V26_UP,
622    },
623    VersionedMigration {
624        version: 27,
625        name: "notes_hot_property_indexes",
626        up: V27_UP,
627    },
628    VersionedMigration {
629        version: 28,
630        name: "notes_key",
631        up: V28_UP,
632    },
633    VersionedMigration {
634        version: 29,
635        name: "note_streams",
636        up: V29_UP,
637    },
638    VersionedMigration {
639        version: 30,
640        name: "tool_source_mounts",
641        up: V30_UP,
642    },
643    VersionedMigration {
644        version: 31,
645        name: "note_versions",
646        up: V31_UP,
647    },
648    VersionedMigration {
649        version: 32,
650        name: "knowledge_count_indexes",
651        up: V32_UP,
652    },
653    VersionedMigration {
654        version: 33,
655        name: "notes_message_recipient_direction",
656        up: V33_UP,
657    },
658    VersionedMigration {
659        version: 34,
660        name: "notes_namespace_created",
661        up: V34_UP,
662    },
663    VersionedMigration {
664        version: 35,
665        name: "notes_unread_probe_recipient_type_direction",
666        up: V35_UP,
667    },
668    VersionedMigration {
669        version: 36,
670        name: "events_operation_attribution",
671        up: V36_UP,
672    },
673    VersionedMigration {
674        version: 37,
675        name: "entity_versions",
676        up: V37_UP,
677    },
678    VersionedMigration {
679        version: 38,
680        name: "entities_legacy_type_index",
681        up: V38_UP,
682    },
683    VersionedMigration {
684        version: 39,
685        name: "knowledge_cursor_indexes",
686        up: V39_UP,
687    },
688    VersionedMigration {
689        version: 40,
690        name: SESSION_IDENTITY_MIGRATION_NAME,
691        up: SESSION_IDENTITY_UP,
692    },
693    VersionedMigration {
694        version: 41,
695        name: "sender_transport",
696        up: V41_UP,
697    },
698    VersionedMigration {
699        version: 42,
700        name: "comm_external_id_channel_scope",
701        up: V42_UP,
702    },
703    VersionedMigration {
704        version: 43,
705        name: "vector_provenance",
706        up: V43_UP,
707    },
708    VersionedMigration {
709        version: 44,
710        name: "comm_outbound_due",
711        up: V44_UP,
712    },
713    VersionedMigration {
714        version: RECIPIENT_TRANSPORT_VERSION,
715        name: "recipient_transport",
716        up: V45_UP,
717    },
718    VersionedMigration {
719        version: 46,
720        name: "memory_visibility_receipts",
721        up: V46_UP,
722    },
723    VersionedMigration {
724        version: 47,
725        name: "attachment_role_quarantine",
726        up: V47_UP,
727    },
728    VersionedMigration {
729        version: 48,
730        name: "acknowledgement_journal",
731        up: V48_UP,
732    },
733    VersionedMigration {
734        version: 49,
735        name: "git_note_property_indexes",
736        up: V49_UP,
737    },
738    VersionedMigration {
739        version: 50,
740        name: "entity_list_plans",
741        up: V50_UP,
742    },
743    VersionedMigration {
744        version: 51,
745        name: "schedule_core_indexes",
746        up: V51_UP,
747    },
748    VersionedMigration {
749        version: 52,
750        name: "comm_core_indexes",
751        up: V52_UP,
752    },
753    VersionedMigration {
754        version: 53,
755        name: "entity_kind_list_order",
756        up: V53_UP,
757    },
758    VersionedMigration {
759        version: MEMORY_VISIBILITY_CUTOVER_VERSION,
760        name: "memory_visibility_epochs",
761        up: V54_UP,
762    },
763];
764
765/// Durable state of ADR-121's boot-gated, two-stage attachment cutover.
766#[derive(Clone, Copy, Debug, Eq, PartialEq)]
767pub enum AttachmentCutoverStatus {
768    /// V20 is current and no stage marker has been committed.
769    Pending,
770    /// Stage 1 committed; boot must finish verified pack-owned attachments.
771    Incomplete,
772    /// V21, the attachment fences, and the attachment-only schema committed.
773    Complete,
774}
775
776fn schema_object_exists(
777    conn: &Connection,
778    object_type: &str,
779    name: &str,
780) -> Result<bool, SqliteError> {
781    conn.query_row(
782        "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
783        rusqlite::params![object_type, name],
784        |row| row.get(0),
785    )
786    .map_err(Into::into)
787}
788
789fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
790    conn.query_row(
791        "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
792        rusqlite::params![table, column],
793        |row| row.get(0),
794    )
795    .map_err(Into::into)
796}
797
798fn require_attachment_schema_objects(
799    conn: &Connection,
800    objects: &[(&str, &str)],
801    phase: &str,
802) -> Result<(), SqliteError> {
803    for (object_type, name) in objects {
804        if !schema_object_exists(conn, object_type, name)? {
805            return Err(SqliteError::InvalidData(format!(
806                "attachment cutover {phase} state is missing {object_type} {name:?}"
807            )));
808        }
809    }
810    Ok(())
811}
812
813fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
814    require_attachment_schema_objects(
815        conn,
816        &[
817            ("table", "attachments"),
818            ("index", "idx_attachments_content_ref"),
819        ],
820        "incomplete",
821    )?;
822    require_legacy_attachment_fences(conn)
823}
824
825fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
826    require_attachment_schema_objects(
827        conn,
828        &[
829            ("table", "attachments"),
830            ("table", "blob_gc_claims"),
831            ("index", "idx_attachments_content_ref"),
832            ("index", "idx_blob_gc_claims_content_ref"),
833            ("trigger", "attachments_reject_claimed_blob_insert"),
834            ("trigger", "attachments_reject_claimed_blob_update"),
835        ],
836        "complete",
837    )?;
838    if schema_column_exists(conn, "entities", "content_ref")? {
839        return Err(SqliteError::InvalidData(
840            "attachment cutover is complete but entities.content_ref still exists".into(),
841        ));
842    }
843    for (object_type, name) in [
844        ("index", "idx_entities_content_ref"),
845        ("trigger", "entities_reject_claimed_blob_insert"),
846        ("trigger", "entities_reject_claimed_blob_update"),
847    ] {
848        if schema_object_exists(conn, object_type, name)? {
849            return Err(SqliteError::InvalidData(format!(
850                "attachment cutover is complete but legacy {object_type} {name:?} still exists"
851            )));
852        }
853    }
854    Ok(())
855}
856
857/// Inspect the coordinated V21 state without mutating the connection.
858///
859/// The marker and migration ledger form one state machine. Impossible pairs
860/// fail closed instead of being guessed into a resumable state.
861pub fn attachment_cutover_status(
862    conn: &Connection,
863) -> Result<AttachmentCutoverStatus, SqliteError> {
864    let version = read_schema_version(conn)?;
865    let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
866    if !marker_table {
867        if version >= ATTACHMENT_CUTOVER_VERSION {
868            return Err(SqliteError::InvalidData(format!(
869                "migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
870            )));
871        }
872        if schema_object_exists(conn, "table", "attachments")? {
873            return Err(SqliteError::InvalidData(
874                "attachments table exists without the durable attachment cutover marker".into(),
875            ));
876        }
877        return Ok(AttachmentCutoverStatus::Pending);
878    }
879
880    let marker: Option<(String, Option<i64>)> = conn
881        .query_row(
882            "SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
883            [],
884            |row| Ok((row.get(0)?, row.get(1)?)),
885        )
886        .optional()?;
887    match marker {
888        Some((state, None)) if state == "incomplete" => {
889            if version >= ATTACHMENT_CUTOVER_VERSION {
890                Err(SqliteError::InvalidData(format!(
891                    "attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
892                )))
893            } else {
894                validate_incomplete_attachment_schema(conn)?;
895                Ok(AttachmentCutoverStatus::Incomplete)
896            }
897        }
898        Some((state, Some(_))) if state == "complete" => {
899            // Later migrations (V22+) are recorded on top of a completed
900            // cutover in the normal course; only a ledger BELOW V21 beside a
901            // complete marker is an impossible pair.
902            if version >= ATTACHMENT_CUTOVER_VERSION {
903                validate_complete_attachment_schema(conn)?;
904                Ok(AttachmentCutoverStatus::Complete)
905            } else {
906                Err(SqliteError::InvalidData(format!(
907                    "attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
908                )))
909            }
910        }
911        Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
912            "invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
913        ))),
914        None => Err(SqliteError::InvalidData(
915            "attachment cutover marker table exists without its singleton row".into(),
916        )),
917    }
918}
919
920fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
921    if !schema_column_exists(conn, "entities", "content_ref")? {
922        return Err(SqliteError::InvalidData(
923            "attachment cutover requires legacy entities.content_ref until finalization".into(),
924        ));
925    }
926    for (object_type, name) in [
927        ("table", "blob_gc_claims"),
928        ("index", "idx_blob_gc_claims_content_ref"),
929        ("index", "idx_entities_content_ref"),
930        ("trigger", "entities_reject_claimed_blob_insert"),
931        ("trigger", "entities_reject_claimed_blob_update"),
932    ] {
933        if !schema_object_exists(conn, object_type, name)? {
934            return Err(SqliteError::InvalidData(format!(
935                "attachment cutover requires legacy {object_type} {name:?} until finalization"
936            )));
937        }
938    }
939    Ok(())
940}
941
942// length() and GLOB both stop scanning at an embedded NUL, so a value of 64
943// hex characters followed by a NUL and arbitrary trailing bytes would pass
944// both. Deriving the canonical byte width from the connection's own text
945// encoding (rather than assuming UTF-8) keeps this arm correct on a database
946// pinned to UTF-16 and fails closed if the probe returns something else,
947// matching the pattern already used by `validate_blob_gc_evidence`.
948fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
949    let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
950    if !(1..=4).contains(&width) {
951        return Err(SqliteError::InvalidData(format!(
952            "the text-encoding width probe returned {width}; refusing canonicality validation"
953        )));
954    }
955    Ok(width * 64)
956}
957
958fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
959    let canonical_bytes = canonical_content_ref_byte_width(conn)?;
960    let invalid: Option<String> = conn
961        .query_row(
962            "SELECT id FROM entities \
963             WHERE content_ref IS NOT NULL \
964               AND (typeof(content_ref) <> 'text' \
965                 OR length(content_ref) <> 64 \
966                 OR length(CAST(content_ref AS BLOB)) <> ?1 \
967                 OR content_ref GLOB '*[^0-9a-f]*') \
968             LIMIT 1",
969            [canonical_bytes],
970            |row| row.get(0),
971        )
972        .optional()?;
973    if let Some(id) = invalid {
974        return Err(SqliteError::InvalidData(format!(
975            "entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
976        )));
977    }
978    Ok(())
979}
980
981fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
982    let canonical_bytes = canonical_content_ref_byte_width(conn)?;
983    for (table, identity) in [
984        ("attachments", "record_uuid"),
985        ("blob_gc_claims", "root_key"),
986    ] {
987        let sql = format!(
988            "SELECT {identity} FROM {table} \
989             WHERE typeof(content_ref) <> 'text' \
990                OR length(content_ref) <> 64 \
991                OR length(CAST(content_ref AS BLOB)) <> ?1 \
992                OR content_ref GLOB '*[^0-9a-f]*' \
993             LIMIT 1"
994        );
995        let invalid: Option<String> = conn
996            .query_row(&sql, [canonical_bytes], |row| row.get(0))
997            .optional()?;
998        if let Some(owner) = invalid {
999            return Err(SqliteError::InvalidData(format!(
1000                "{table}.content_ref for {identity} {owner:?} is not canonical"
1001            )));
1002        }
1003    }
1004    Ok(())
1005}
1006
1007fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
1008    let dangling: Option<(String, String)> = conn
1009        .query_row(
1010            "SELECT record_uuid, substrate FROM attachments AS attachment \
1011             WHERE (substrate = 'entity' AND NOT EXISTS ( \
1012                       SELECT 1 FROM entities WHERE id = attachment.record_uuid \
1013                   )) \
1014                OR (substrate = 'note' AND NOT EXISTS ( \
1015                       SELECT 1 FROM notes WHERE id = attachment.record_uuid \
1016                   )) \
1017             LIMIT 1",
1018            [],
1019            |row| Ok((row.get(0)?, row.get(1)?)),
1020        )
1021        .optional()?;
1022    if let Some((record_uuid, substrate)) = dangling {
1023        return Err(SqliteError::InvalidData(format!(
1024            "attachment role references absent {substrate} record {record_uuid:?}"
1025        )));
1026    }
1027    Ok(())
1028}
1029
1030fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
1031    let conflict: Option<String> = conn
1032        .query_row(
1033            "SELECT entity.id FROM entities AS entity \
1034             LEFT JOIN attachments AS attachment \
1035               ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
1036             WHERE entity.content_ref IS NOT NULL \
1037               AND (attachment.record_uuid IS NULL \
1038                 OR attachment.substrate <> 'entity' \
1039                 OR attachment.content_ref <> entity.content_ref) \
1040             LIMIT 1",
1041            [],
1042            |row| row.get(0),
1043        )
1044        .optional()?;
1045    if let Some(record_uuid) = conflict {
1046        return Err(SqliteError::InvalidData(format!(
1047            "legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
1048        )));
1049    }
1050    Ok(())
1051}
1052
1053fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
1054    require_legacy_attachment_fences(conn)?;
1055    conn.execute_batch(V21_STAGE_UP)?;
1056    validate_canonical_legacy_refs(conn)?;
1057    validate_canonical_attachment_and_claim_refs(conn)?;
1058
1059    let conflict: Option<String> = conn
1060        .query_row(
1061            "SELECT entity.id FROM entities AS entity \
1062             JOIN attachments AS attachment \
1063               ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
1064             WHERE entity.content_ref IS NOT NULL \
1065               AND (attachment.substrate <> 'entity' \
1066                 OR attachment.content_ref <> entity.content_ref) \
1067             LIMIT 1",
1068            [],
1069            |row| row.get(0),
1070        )
1071        .optional()?;
1072    if let Some(record_uuid) = conflict {
1073        return Err(SqliteError::InvalidData(format!(
1074            "existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
1075        )));
1076    }
1077
1078    conn.execute(
1079        "INSERT INTO attachments \
1080         (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
1081         SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
1082         FROM entities WHERE content_ref IS NOT NULL \
1083         ON CONFLICT(record_uuid, role) DO NOTHING",
1084        [],
1085    )?;
1086    validate_legacy_content_backfill(conn)?;
1087
1088    // The caller holds the canonical database GC owner, so every preexisting
1089    // claim is abandoned. Clearing happens before the durable incomplete
1090    // marker is exposed and remains inside this one transaction.
1091    conn.execute("DELETE FROM blob_gc_claims", [])?;
1092    conn.execute(
1093        "INSERT INTO attachment_cutover_state \
1094         (singleton, state, started_at, completed_at) \
1095         VALUES (1, 'incomplete', ?1, NULL) \
1096         ON CONFLICT(singleton) DO NOTHING",
1097        [now],
1098    )?;
1099    Ok(())
1100}
1101
1102/// Commit stage 1 of the coordinated V21 migration.
1103///
1104/// The caller must hold [`crate::stores::blob::DatabaseGcOwnerGuard`] for the
1105/// canonical database before entering this function and retain it through
1106/// application backfill and finalization. This function owns one IMMEDIATE
1107/// SQLite transaction; a failure leaves neither its DDL nor marker visible.
1108pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
1109    let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
1110    RawMigrationWriteUnit::new(conn, &admission)?
1111        .run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
1112}
1113
1114/// Stage attachment cutover with an explicit raw-connection policy.
1115pub fn stage_attachment_cutover_with_policy(
1116    conn: &mut Connection,
1117    policy: &MigrationWritePolicy,
1118) -> Result<(), SqliteError> {
1119    let admission =
1120        WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
1121    RawMigrationWriteUnit::new(conn, &admission)?
1122        .run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
1123}
1124
1125/// Stage attachment cutover using admission already held by a pooled caller.
1126pub(crate) fn stage_attachment_cutover_with_admission(
1127    conn: &mut Connection,
1128    admission: &WriteAdmission,
1129) -> Result<(), SqliteError> {
1130    match attachment_cutover_status(conn)? {
1131        AttachmentCutoverStatus::Complete => return Ok(()),
1132        AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
1133    }
1134    if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
1135        return Err(SqliteError::InvalidData(format!(
1136            "attachment cutover stage requires canonical V{} schema",
1137            ATTACHMENT_CUTOVER_VERSION - 1
1138        )));
1139    }
1140
1141    let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1142    if let Err(error) = admission.check() {
1143        let rollback = tx.rollback();
1144        return Err(capacity_refusal_after_rollback(
1145            conn,
1146            rollback,
1147            error,
1148            "attachment cutover stage",
1149        ));
1150    }
1151    let status = attachment_cutover_status(&tx)?;
1152    if status == AttachmentCutoverStatus::Complete {
1153        return Ok(());
1154    }
1155    stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
1156    tx.commit()?;
1157    Ok(())
1158}
1159
1160/// Add one host-verified pack-owned attachment during V21 stage 2.
1161///
1162/// This helper is deliberately transaction-neutral: it neither begins nor
1163/// commits a transaction. The boot coordinator can therefore apply the full
1164/// verified vector in one caller-owned IMMEDIATE transaction while retaining
1165/// the canonical database GC owner. Reapplying the same role and content is
1166/// idempotent; a different substrate or digest for that role fails closed.
1167#[allow(clippy::too_many_arguments)]
1168pub fn apply_generic_verified_attachment(
1169    conn: &Connection,
1170    record_uuid: &str,
1171    substrate: &str,
1172    role: &str,
1173    content_ref: &ContentRef,
1174    media_type: Option<&str>,
1175    size_bytes: Option<u64>,
1176    created_at: i64,
1177) -> Result<(), SqliteError> {
1178    if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
1179        return Err(SqliteError::InvalidData(
1180            "verified application attachments may only be applied while V21 cutover is incomplete"
1181                .into(),
1182        ));
1183    }
1184    if role.is_empty() || role.chars().any(char::is_control) {
1185        return Err(SqliteError::InvalidData(
1186            "attachment role must be non-empty and contain no control characters".into(),
1187        ));
1188    }
1189    let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
1190        SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
1191    })?;
1192    let owner_table = match substrate {
1193        "entity" => "entities",
1194        "note" => "notes",
1195        other => {
1196            return Err(SqliteError::InvalidData(format!(
1197                "attachment substrate must be 'entity' or 'note', got {other:?}"
1198            )))
1199        }
1200    };
1201    let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
1202    let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
1203    if !owner_exists {
1204        return Err(SqliteError::InvalidData(format!(
1205            "cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
1206        )));
1207    }
1208    let claimed: bool = conn.query_row(
1209        "SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
1210        [content_ref.as_str()],
1211        |row| row.get(0),
1212    )?;
1213    if claimed {
1214        return Err(SqliteError::InvalidData(format!(
1215            "cannot attach claimed content_ref {} during V21 cutover",
1216            content_ref.as_str()
1217        )));
1218    }
1219
1220    let changed = conn.execute(
1221        "INSERT INTO attachments \
1222         (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
1223         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
1224         ON CONFLICT(record_uuid, role) DO UPDATE SET \
1225             media_type = excluded.media_type, \
1226             size_bytes = excluded.size_bytes, \
1227             created_at = excluded.created_at \
1228         WHERE attachments.substrate = excluded.substrate \
1229           AND attachments.content_ref = excluded.content_ref",
1230        rusqlite::params![
1231            record_uuid,
1232            substrate,
1233            role,
1234            content_ref.as_str(),
1235            media_type,
1236            size_bytes,
1237            created_at,
1238        ],
1239    )?;
1240    if changed == 0 {
1241        return Err(SqliteError::InvalidData(format!(
1242            "attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
1243        )));
1244    }
1245    Ok(())
1246}
1247
1248fn finalize_attachment_cutover_on_connection(
1249    conn: &Connection,
1250    now: i64,
1251) -> Result<(), SqliteError> {
1252    require_legacy_attachment_fences(conn)?;
1253    validate_canonical_legacy_refs(conn)?;
1254    validate_canonical_attachment_and_claim_refs(conn)?;
1255    validate_attachment_record_owners(conn)?;
1256    validate_legacy_content_backfill(conn)?;
1257
1258    let remaining_claims: i64 =
1259        conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
1260    if remaining_claims != 0 {
1261        return Err(SqliteError::InvalidData(format!(
1262            "attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
1263        )));
1264    }
1265
1266    let uncovered_model: Option<String> = conn
1267        .query_row(
1268            "SELECT model.id FROM entities AS model \
1269             WHERE model.entity_type = 'moodboard_model' \
1270               AND model.content_ref IS NOT NULL \
1271               AND NOT EXISTS ( \
1272                   SELECT 1 FROM attachments AS attachment \
1273                   WHERE attachment.record_uuid = model.id \
1274                     AND attachment.substrate = 'entity' \
1275                     AND attachment.role = 'fann-network' \
1276               ) \
1277             LIMIT 1",
1278            [],
1279            |row| row.get(0),
1280        )
1281        .optional()?;
1282    if let Some(record_uuid) = uncovered_model {
1283        return Err(SqliteError::InvalidData(format!(
1284            "moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
1285        )));
1286    }
1287
1288    conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
1289    conn.execute_batch(
1290        "DROP TRIGGER entities_reject_claimed_blob_insert; \
1291         DROP TRIGGER entities_reject_claimed_blob_update; \
1292         DROP INDEX idx_entities_content_ref; \
1293         ALTER TABLE entities DROP COLUMN content_ref;",
1294    )?;
1295    conn.execute(
1296        "UPDATE attachment_cutover_state \
1297         SET state = 'complete', completed_at = ?1 \
1298         WHERE singleton = 1 AND state = 'incomplete'",
1299        [now],
1300    )?;
1301    Ok(())
1302}
1303
1304fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
1305    let migration = MIGRATIONS
1306        .iter()
1307        .find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
1308        .expect("V21 migration must be registered");
1309    conn.execute(
1310        "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
1311        rusqlite::params![migration.version, migration.name, now],
1312    )?;
1313    Ok(())
1314}
1315
1316/// Atomically switch an explicitly staged database to attachment-only V21.
1317///
1318/// The caller must still hold the canonical database GC owner. The exclusive
1319/// transition revalidates every legacy and attachment reference, verifies
1320/// moodboard model role coverage, replaces the claim fences, removes the old
1321/// column, marks the cutover complete, and records V21 in one transaction.
1322pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
1323    let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
1324    RawMigrationWriteUnit::new(conn, &admission)?
1325        .run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
1326}
1327
1328/// Finalize attachment cutover with an explicit raw-connection policy.
1329pub fn finalize_attachment_cutover_with_policy(
1330    conn: &mut Connection,
1331    policy: &MigrationWritePolicy,
1332) -> Result<(), SqliteError> {
1333    let admission =
1334        WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
1335    RawMigrationWriteUnit::new(conn, &admission)?
1336        .run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
1337}
1338
1339/// Finalize attachment cutover using admission already held by a pooled caller.
1340pub(crate) fn finalize_attachment_cutover_with_admission(
1341    conn: &mut Connection,
1342    admission: &WriteAdmission,
1343) -> Result<(), SqliteError> {
1344    if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
1345        return Ok(());
1346    }
1347    let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
1348    if let Err(error) = admission.check() {
1349        let rollback = tx.rollback();
1350        return Err(capacity_refusal_after_rollback(
1351            conn,
1352            rollback,
1353            error,
1354            "attachment cutover finalization",
1355        ));
1356    }
1357    match attachment_cutover_status(&tx)? {
1358        AttachmentCutoverStatus::Complete => return Ok(()),
1359        AttachmentCutoverStatus::Pending => {
1360            return Err(SqliteError::InvalidData(
1361                "attachment cutover must complete stage 1 before finalization".into(),
1362            ))
1363        }
1364        AttachmentCutoverStatus::Incomplete => {}
1365    }
1366    let now = chrono::Utc::now().timestamp_micros();
1367    finalize_attachment_cutover_on_connection(&tx, now)?;
1368    record_attachment_cutover_migration(&tx, now)?;
1369    tx.commit()?;
1370    Ok(())
1371}
1372
1373/// Read the ordered migration ledger prefix without interpreting its rows.
1374fn read_applied_migration_ledger(
1375    conn: &Connection,
1376    through_version: u32,
1377) -> Result<Vec<(u32, String)>, SqliteError> {
1378    let mut stmt = conn.prepare(
1379        "SELECT version, name FROM _schema_migrations \
1380         WHERE version <= ?1 ORDER BY version ASC",
1381    )?;
1382    let rows = stmt
1383        .query_map([through_version], |row| {
1384            Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
1385        })?
1386        .collect::<Result<Vec<_>, _>>()?;
1387    Ok(rows)
1388}
1389
1390/// Require the applied versions through `through_version` to be the exact
1391/// contiguous canonical prefix of [`MIGRATIONS`]. A matching `MAX(version)` is
1392/// insufficient: a missing middle row or foreign version can expose a
1393/// materially different schema while retaining the same maximum.
1394fn validate_applied_migration_versions(
1395    applied: &[(u32, String)],
1396    through_version: u32,
1397) -> Result<(), SqliteError> {
1398    let expected: Vec<&VersionedMigration> = MIGRATIONS
1399        .iter()
1400        .filter(|migration| migration.version <= through_version)
1401        .collect();
1402    let mut applied_index = 0;
1403
1404    for migration in expected {
1405        let Some((version, applied_name)) = applied.get(applied_index) else {
1406            return Err(SqliteError::InvalidData(format!(
1407                "migration history is missing version {} ('{}'); the applied ledger must be \
1408                 the exact contiguous canonical sequence through version {through_version}",
1409                migration.version, migration.name,
1410            )));
1411        };
1412        if *version < migration.version {
1413            return Err(SqliteError::InvalidData(format!(
1414                "migration history contains unknown version {version} recorded as \
1415                 '{applied_name}'; the applied ledger must contain only canonical versions"
1416            )));
1417        }
1418        if *version > migration.version {
1419            return Err(SqliteError::InvalidData(format!(
1420                "migration history is missing version {} ('{}'); found version {version} \
1421                 next instead",
1422                migration.version, migration.name,
1423            )));
1424        }
1425        applied_index += 1;
1426    }
1427
1428    if let Some((version, name)) = applied.get(applied_index) {
1429        return Err(SqliteError::InvalidData(format!(
1430            "migration history contains unknown version {version} recorded as '{name}'; \
1431             the applied ledger must contain only canonical versions"
1432        )));
1433    }
1434
1435    Ok(())
1436}
1437
1438fn validate_applied_migration_names(
1439    applied: &[(u32, String)],
1440    through_version: u32,
1441    allow_known_v19_repairs: bool,
1442) -> Result<(), SqliteError> {
1443    for ((version, applied_name), migration) in applied.iter().zip(
1444        MIGRATIONS
1445            .iter()
1446            .filter(|migration| migration.version <= through_version),
1447    ) {
1448        debug_assert_eq!(*version, migration.version);
1449        if migration.name != applied_name.as_str() {
1450            if allow_known_v19_repairs && matches!(*version, 13 | 14) {
1451                continue;
1452            }
1453            return Err(SqliteError::InvalidData(format!(
1454                "migration version {version} is recorded under name '{applied_name}', \
1455                 expected '{expected}'. This database's migration history does not match \
1456                 the current binary; recreate it from the current schema or repair the \
1457                 specific known divergence via a dedicated migration.",
1458                expected = migration.name,
1459            )));
1460        }
1461    }
1462
1463    Ok(())
1464}
1465
1466/// Confirm the complete applied ledger is the canonical prefix, including
1467/// names. The only historical V13/V14 name divergence is repaired by V19
1468/// before this validator runs for a pre-V19 database.
1469fn validate_applied_migration_ledger(
1470    conn: &Connection,
1471    through_version: u32,
1472) -> Result<(), SqliteError> {
1473    let applied = read_applied_migration_ledger(conn, through_version)?;
1474    validate_applied_migration_versions(&applied, through_version)?;
1475    validate_applied_migration_names(&applied, through_version, false)
1476}
1477
1478const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
1479
1480/// Read the applied schema version from an open connection **without** running
1481/// migrations. Returns 0 when the `_schema_migrations` ledger is absent (an
1482/// un-migrated or empty database); any other failure (BUSY, IO) propagates —
1483/// collapsing it to 0 would misreport a live database as un-migrated. Never
1484/// writes.
1485pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
1486    match conn.query_row(
1487        "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1488        [],
1489        |row| row.get(0),
1490    ) {
1491        Ok(version) => Ok(version),
1492        Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
1493            if msg.contains("no such table: _schema_migrations") =>
1494        {
1495            Ok(0)
1496        }
1497        Err(e) => Err(e.into()),
1498    }
1499}
1500
1501/// Open `path` read-only and report its applied schema version without creating
1502/// or migrating the file. The caller must ensure `path` exists — opening a
1503/// missing file read-only errors rather than creating it. This is the path used
1504/// by schema-inspection commands that must not mutate the database.
1505pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
1506    let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1507    read_schema_version(&conn)
1508}
1509
1510/// Open `path` read-only and require the exact canonical current ledger and
1511/// physical cutover state, without creating files or applying migrations.
1512pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
1513    let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1514    validate_schema_is_current(&conn)
1515}
1516
1517/// Require an already-open database to match this build's latest core schema
1518/// without applying migrations.
1519///
1520/// A read-only snapshot behind the current migration set cannot be repaired in
1521/// place, while a snapshot ahead of the binary may contain schema this build
1522/// does not understand. Both directions fail with an actionable diagnostic; an
1523/// exact match performs no writes.
1524pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
1525    let current_version = read_schema_version(conn)?;
1526    let latest_version = latest_schema_version();
1527
1528    if current_version < latest_version {
1529        return Err(SqliteError::InvalidData(format!(
1530            "read-only database schema version {current_version} is behind the latest known \
1531             migration {latest_version}; migrate a writable copy with this build before opening \
1532             the snapshot read-only"
1533        )));
1534    }
1535    if current_version > latest_version {
1536        return Err(SqliteError::InvalidData(format!(
1537            "read-only database schema version {current_version} is ahead of the latest known \
1538             migration {latest_version}; use a compatible newer build or recreate the snapshot"
1539        )));
1540    }
1541
1542    // Numeric equality alone is not enough: a database can carry the current
1543    // maximum version under renamed or foreign migration ledger entries while
1544    // exposing a materially different schema. Writable boot runs this same
1545    // closed-name validation in `bootstrap_migration_ledger`; snapshot inspection
1546    // must not accept a history that ordinary boot would reject merely because
1547    // it cannot repair it in place.
1548    validate_applied_migration_ledger(conn, current_version)?;
1549    // `>=`, not `==`: later migrations (V22+) record on top of a completed
1550    // cutover, and a ledger at the latest version must not exempt the
1551    // physical cutover state from validation.
1552    if current_version >= ATTACHMENT_CUTOVER_VERSION
1553        && attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
1554    {
1555        return Err(SqliteError::InvalidData(
1556            "read-only database has not completed the V21 attachment cutover".into(),
1557        ));
1558    }
1559
1560    if current_version >= MEMORY_VISIBILITY_CUTOVER_VERSION {
1561        memory_visibility::validate_cutover(conn)?;
1562    }
1563
1564    Ok(current_version)
1565}
1566
1567/// Require the complete core migration ledger and readable memory provenance
1568/// schema without applying migrations or inferring an individual note's epoch.
1569pub fn validate_memory_visibility_cutover(conn: &Connection) -> Result<(), SqliteError> {
1570    validate_schema_is_current(conn).map(|_| ())
1571}
1572
1573#[cfg(test)]
1574pub(crate) mod test_sync {
1575    use std::sync::atomic::AtomicU32;
1576    use std::sync::{Arc, Barrier, Mutex};
1577
1578    /// When set, `run_versioned_migrations` parks after its initial (stale)
1579    /// ledger read until every racing thread has arrived — forcing the
1580    /// contended interleaving the concurrent-boot test asserts on.
1581    pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
1582    /// Counts entries into the under-lock sibling fast-forward branch.
1583    pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
1584    /// Set by the SQLite busy handler installed on participating connections:
1585    /// `true` means SQLite itself reported a blocked lock acquisition to the
1586    /// loser — actual contention, not merely an intended attempt.
1587    pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
1588        std::sync::atomic::AtomicBool::new(false);
1589
1590    /// Busy handler for participating test connections: records that SQLite
1591    /// observed a busy acquisition, then keeps retrying.
1592    pub(crate) fn record_busy(_count: i32) -> bool {
1593        BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
1594        std::thread::sleep(std::time::Duration::from_millis(1));
1595        true
1596    }
1597
1598    /// Set by the winner immediately before committing its first migration
1599    /// transaction — i.e. before the write lock is first released.
1600    pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
1601        std::sync::atomic::AtomicBool::new(false);
1602    /// Recorded by the loser when its first `BEGIN IMMEDIATE` returns: whether
1603    /// the winner had already committed at that moment. `true` is direct
1604    /// evidence the loser's lock acquisition blocked across the winner's held
1605    /// write lock rather than the two calls serializing by scheduler accident.
1606    pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
1607        std::sync::atomic::AtomicBool::new(false);
1608
1609    std::thread_local! {
1610        /// Opt-in flag: only threads that set this participate in the barrier,
1611        /// so unrelated tests migrating in parallel are never parked.
1612        pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
1613            const { std::cell::Cell::new(false) };
1614        /// Whether this thread has already instrumented its first BEGIN.
1615        pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
1616            const { std::cell::Cell::new(false) };
1617    }
1618}
1619
1620/// Apply the ordinary unapplied migration prefix in order.
1621///
1622/// The operation is idempotent and each ordinary migration runs in its own
1623/// transaction. V21 is the application-assisted exception: a database with no
1624/// legacy content references may complete V21 atomically here, while a legacy
1625/// database stops successfully at V20 so the async host can stage, verify, and
1626/// finalize the attachment cutover. Errors on a non-contiguous migration array,
1627/// a non-canonical applied ledger, or a failed migration.
1628fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
1629    let configured = conn.path().unwrap_or_default();
1630    let raw_path = if configured.is_empty() {
1631        conn.query_row(
1632            "SELECT file FROM pragma_database_list WHERE name = 'main'",
1633            [],
1634            |row| row.get::<_, String>(0),
1635        )?
1636    } else {
1637        configured.to_string()
1638    };
1639
1640    if raw_path.is_empty() {
1641        return Ok(None);
1642    }
1643    let canonical = std::fs::canonicalize(&raw_path).map_err(SqliteError::Io)?;
1644    #[cfg(any(unix, windows))]
1645    crate::pool::opened_sqlite_file_identity(conn, &canonical)?;
1646    Ok(Some(canonical))
1647}
1648
1649fn validate_database_gc_owner(
1650    conn: &Connection,
1651    owner: &DatabaseGcOwnerGuard,
1652) -> Result<(), SqliteError> {
1653    let connection_path = canonical_connection_database_path(conn)?;
1654    if owner.database_path() != connection_path.as_deref() {
1655        return Err(SqliteError::InvalidData(format!(
1656            "database GC owner targets {:?}, but migration connection targets {:?}",
1657            owner.database_path(),
1658            connection_path.as_deref(),
1659        )));
1660    }
1661    Ok(())
1662}
1663
1664pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
1665    let database_path = canonical_connection_database_path(conn)?;
1666    let admission = WriteAdmission::for_canonical_path(database_path.clone())?;
1667    run_raw_migrations_with_admission(conn, database_path, &admission)
1668}
1669
1670/// Run migrations on a raw connection with a captured admission policy.
1671pub fn run_migrations_with_policy(
1672    conn: &mut Connection,
1673    policy: &MigrationWritePolicy,
1674) -> Result<u32, SqliteError> {
1675    let database_path = canonical_connection_database_path(conn)?;
1676    let admission = WriteAdmission::for_migration_policy(database_path.clone(), policy)?;
1677    run_raw_migrations_with_admission(conn, database_path, &admission)
1678}
1679
1680fn run_raw_migrations_with_admission(
1681    conn: &mut Connection,
1682    database_path: Option<PathBuf>,
1683    admission: &WriteAdmission,
1684) -> Result<u32, SqliteError> {
1685    if let Some(database_path) = database_path {
1686        // This raw API may have been handed a connection behind an opaque pool
1687        // writer guard. Never wait here and invert the canonical
1688        // owner-before-writer order; fail closed and direct production callers
1689        // to `StorageBackend::prepare_core_schema` instead.
1690        let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
1691            SqliteError::InvalidData(format!(
1692                "failed to acquire database GC owner before schema migration: {error}"
1693            ))
1694        })?;
1695        return run_migrations_with_database_gc_owner(
1696            &mut RawMigrationTransactions::new(conn, admission),
1697            &owner,
1698            admission,
1699        );
1700    }
1701
1702    // A raw in-memory connection has no durable/cross-process GC domain. The
1703    // production in-memory backend still uses the owner-aware path below.
1704    run_versioned_migrations(
1705        &mut RawMigrationTransactions::new(conn, admission),
1706        None,
1707        admission,
1708    )
1709}
1710
1711pub(crate) fn run_migrations_with_database_gc_owner(
1712    writes: &mut impl MigrationTransactions,
1713    owner: &DatabaseGcOwnerGuard,
1714    admission: &WriteAdmission,
1715) -> Result<u32, SqliteError> {
1716    run_versioned_migrations(writes, Some(owner), admission)
1717}
1718
1719/// Raise-only to a 5s busy_timeout floor for one admitted migration unit and
1720/// restore the caller's value after it. Concurrent boots (multiple processes
1721/// migrating the same file) contend on SQLite's write lock; a short hot-path
1722/// busy_timeout cannot wait out a sibling's migration, and a caller whose
1723/// configured timeout is already longer is never reduced.
1724fn with_migration_busy_timeout<T>(
1725    conn: &mut Connection,
1726    operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
1727) -> Result<T, SqliteError> {
1728    let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
1729    let raised = prior_busy_ms < 5_000;
1730    if raised {
1731        conn.busy_timeout(std::time::Duration::from_secs(5))?;
1732    }
1733    let result = operation(conn);
1734    if raised {
1735        let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
1736    }
1737    result
1738}
1739
1740/// What one admitted versioned-migration transaction did.
1741enum MigrationStep {
1742    /// This transaction applied the migration and committed.
1743    Applied,
1744    /// A sibling had already recorded this version or later; the value is the
1745    /// ledger maximum read under the write lock.
1746    AppliedBySibling(u32),
1747    /// V21 on a legacy database: left for the application-assisted cutover.
1748    Deferred,
1749}
1750
1751/// Bring the core schema to the latest version.
1752///
1753/// ADR-154 section 4 admits each migration transaction on its own, so the
1754/// ledger bootstrap with the pre-lock ledger read is one admitted unit and
1755/// every migration that still has to run is another, re-reading the ledger
1756/// under the write lock. No volume lease is held from one version to the next.
1757fn run_versioned_migrations(
1758    writes: &mut impl MigrationTransactions,
1759    owner: Option<&DatabaseGcOwnerGuard>,
1760    admission: &WriteAdmission,
1761) -> Result<u32, SqliteError> {
1762    let current_version = writes.admitted(|conn| {
1763        if let Some(owner) = owner {
1764            validate_database_gc_owner(conn, owner)?;
1765        }
1766        with_migration_busy_timeout(conn, |conn| bootstrap_migration_ledger(conn, admission))
1767    })?;
1768
1769    // Deterministic-contention hook: parks every caller after the stale ledger
1770    // read (no SQLite lock and no volume lease held) until all racing test
1771    // threads have observed it, so they are then released to compete for the
1772    // IMMEDIATE write lock of their first migration.
1773    #[cfg(test)]
1774    if test_sync::PARTICIPATE.with(|p| p.get()) {
1775        let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
1776        if let Some(barrier) = barrier {
1777            barrier.wait();
1778        }
1779    }
1780    let latest_version = latest_schema_version();
1781
1782    let mut applied_version = current_version;
1783    // Floor advanced when a sibling's work is observed under the write lock,
1784    // so a losing process skips the remaining already-applied migrations
1785    // without opening a transaction for each.
1786    let mut skip_through = current_version;
1787
1788    for migration in MIGRATIONS {
1789        if migration.version <= skip_through {
1790            applied_version = applied_version.max(migration.version);
1791            continue;
1792        }
1793        let step = writes.admitted(|conn| {
1794            with_migration_busy_timeout(conn, |conn| {
1795                apply_versioned_migration(conn, migration, admission, latest_version)
1796            })
1797        })?;
1798        match step {
1799            MigrationStep::Applied => applied_version = migration.version,
1800            MigrationStep::AppliedBySibling(sibling_version) => {
1801                skip_through = sibling_version.min(latest_version);
1802                applied_version = applied_version.max(migration.version);
1803            }
1804            MigrationStep::Deferred => break,
1805        }
1806    }
1807
1808    // Validate again after the loop: our own commits and any under-lock
1809    // sibling fast-forward must both leave the exact canonical ledger, not
1810    // merely advance its maximum version.
1811    writes.admitted(|conn| validate_applied_migration_ledger(conn, applied_version))?;
1812
1813    Ok(applied_version)
1814}
1815
1816/// Create the migration ledger if it is missing and validate what it records,
1817/// returning the recorded version. The read is taken before any migration
1818/// write lock; each migration re-reads the ledger under that lock.
1819fn bootstrap_migration_ledger(
1820    conn: &mut Connection,
1821    admission: &WriteAdmission,
1822) -> Result<u32, SqliteError> {
1823    admission.check()?;
1824    conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
1825    require_autocommit(conn, "migration-tracking bootstrap")?;
1826
1827    let current_version: u32 = read_schema_version(conn)?;
1828
1829    // A database whose recorded version is ahead of the latest known migration
1830    // predates the consolidated V1 baseline (ADR-015) — e.g. it still carries the
1831    // pre-consolidation V2..V22 ledger — or was written by a newer build. Either
1832    // way the baseline schema would be silently skipped, leaving the process on a
1833    // stale schema. Fail loudly instead of corrupting silently.
1834    let latest_version = latest_schema_version();
1835    if current_version > latest_version {
1836        return Err(SqliteError::InvalidData(format!(
1837            "database schema version {current_version} is ahead of the latest known migration \
1838             {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
1839             written by a newer build. Recreate it from the current schema; in-place downgrade is \
1840             not supported."
1841        )));
1842    }
1843
1844    // Every writable upgrade starts from an exact canonical version sequence;
1845    // fail before applying new migrations if MAX(version) hides a missing or
1846    // foreign row. A pre-V19 database may still carry the V13/V14 name
1847    // divergence that V19 exists to repair. Name validation therefore permits
1848    // exactly those two rows before V19, while every unrelated mismatch still
1849    // fails before any new migration is applied.
1850    let applied = read_applied_migration_ledger(conn, current_version)?;
1851    validate_applied_migration_versions(&applied, current_version)?;
1852    validate_applied_migration_names(&applied, current_version, current_version < 19)?;
1853
1854    Ok(current_version)
1855}
1856
1857/// Apply one versioned migration in its own IMMEDIATE transaction, probing
1858/// capacity after BEGIN and before its first statement (ADR-154 section 4).
1859fn apply_versioned_migration(
1860    conn: &mut Connection,
1861    migration: &VersionedMigration,
1862    admission: &WriteAdmission,
1863    latest_version: u32,
1864) -> Result<MigrationStep, SqliteError> {
1865    // IMMEDIATE: take the write lock up front so concurrent boots serialize
1866    // here instead of failing mid-migration when a DEFERRED transaction
1867    // upgrades to a write.
1868    #[cfg(test)]
1869    let instrumented_first_begin =
1870        test_sync::PARTICIPATE.with(|p| p.get()) && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
1871    #[cfg(test)]
1872    if instrumented_first_begin {
1873        test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
1874    }
1875    // Replaces the busy_timeout raised for this unit on a participating test
1876    // connection: records SQLite-observed contention, then keeps retrying.
1877    #[cfg(test)]
1878    if test_sync::PARTICIPATE.with(|p| p.get()) {
1879        conn.busy_handler(Some(test_sync::record_busy))?;
1880    }
1881    let tx = conn
1882        .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1883        .map_err(|e| SqliteError::Migration {
1884            version: migration.version,
1885            error: e.to_string(),
1886        })?;
1887    if let Err(error) = admission.check() {
1888        let rollback = tx.rollback();
1889        return Err(capacity_refusal_after_rollback(
1890            conn,
1891            rollback,
1892            error,
1893            "core schema migration",
1894        ));
1895    }
1896
1897    // Re-check under the write lock: a sibling process may have applied
1898    // this migration (and possibly later ones) while we waited. Running
1899    // its DDL again would fail; fast-forward past everything it applied.
1900    let sibling_version: u32 = tx
1901        .query_row(
1902            "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1903            [],
1904            |row| row.get(0),
1905        )
1906        .map_err(|e| SqliteError::Migration {
1907            version: migration.version,
1908            error: e.to_string(),
1909        })?;
1910    #[cfg(test)]
1911    if instrumented_first_begin {
1912        use std::sync::atomic::Ordering::SeqCst;
1913        if sibling_version == 0 {
1914            // Winner: hold the write lock until SQLite has reported a
1915            // busy acquisition to the loser (its busy handler fired) —
1916            // proof the loser's BEGIN is actually blocked on this held
1917            // lock, not merely intended. Bounded so a regression fails
1918            // the assertion instead of hanging the test.
1919            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1920            while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline {
1921                std::thread::yield_now();
1922            }
1923        } else {
1924            // Loser: our first BEGIN just returned. Record whether the
1925            // winner had already committed — true means we blocked across
1926            // its held lock.
1927            test_sync::LOSER_SAW_WINNER_COMMIT
1928                .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
1929        }
1930    }
1931
1932    // The ahead-of-latest guard above ran on a pre-lock read; a newer
1933    // build may have committed a version past ours while we waited for
1934    // the write lock. Accepting it (clamped) would return Ok on a schema
1935    // this binary does not understand — reject it the same way.
1936    if sibling_version > latest_version {
1937        return Err(SqliteError::InvalidData(format!(
1938            "database schema version {sibling_version} is ahead of the latest known \
1939             migration {latest_version} (committed by a concurrent process while this \
1940             one waited for the migration write lock). This build cannot run against \
1941             the newer schema; upgrade the binary or recreate the database."
1942        )));
1943    }
1944
1945    if sibling_version >= migration.version {
1946        #[cfg(test)]
1947        test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1948        return Ok(MigrationStep::AppliedBySibling(sibling_version));
1949    }
1950
1951    if migration.version == 46 {
1952        memory_visibility::capture_pre_v46(&tx).map_err(|error| SqliteError::Migration {
1953            version: migration.version,
1954            error: error.to_string(),
1955        })?;
1956        #[cfg(test)]
1957        memory_visibility::test_state::stop_at(memory_visibility::test_state::Stop::AfterCapture)?;
1958    }
1959
1960    if migration.version == ATTACHMENT_CUTOVER_VERSION {
1961        let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
1962            version: migration.version,
1963            error: e.to_string(),
1964        })?;
1965        let legacy_refs: i64 = tx
1966            .query_row(
1967                "SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
1968                [],
1969                |row| row.get(0),
1970            )
1971            .map_err(|e| SqliteError::Migration {
1972                version: migration.version,
1973                error: e.to_string(),
1974            })?;
1975
1976        // V21 belongs to the boot coordinator. Ordinary backend open may
1977        // finish the degenerate zero-ref case atomically, but it must not
1978        // expose a dual-source interval or eagerly stage a legacy DB.
1979        if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
1980            drop(tx);
1981            return Ok(MigrationStep::Deferred);
1982        }
1983        if status != AttachmentCutoverStatus::Pending {
1984            return Err(SqliteError::Migration {
1985                version: migration.version,
1986                error: format!("unexpected attachment cutover state {status:?}"),
1987            });
1988        }
1989
1990        let now = chrono::Utc::now().timestamp_micros();
1991        stage_attachment_cutover_on_connection(&tx, now).map_err(|e| SqliteError::Migration {
1992            version: migration.version,
1993            error: e.to_string(),
1994        })?;
1995        finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1996            SqliteError::Migration {
1997                version: migration.version,
1998                error: e.to_string(),
1999            }
2000        })?;
2001    } else if migration.version == 36 {
2002        // The events store DDL may already have added these columns.
2003        crate::stores::event::ensure_operation_attribution_columns(&tx).map_err(|error| {
2004            SqliteError::Migration {
2005                version: migration.version,
2006                error: error.to_string(),
2007            }
2008        })?;
2009    } else if migration.version == 44 {
2010        migrate_outbound_due_key(&tx).map_err(|error| SqliteError::Migration {
2011            version: migration.version,
2012            error: error.to_string(),
2013        })?;
2014    } else if migration.version == 48 {
2015        migrate_acknowledgement_journal(&tx).map_err(|error| SqliteError::Migration {
2016            version: migration.version,
2017            error: error.to_string(),
2018        })?;
2019    } else if migration.name == SESSION_IDENTITY_MIGRATION_NAME {
2020        tx.execute_batch(migration.up)
2021            .map_err(|error| SqliteError::Migration {
2022                version: migration.version,
2023                error: error.to_string(),
2024            })?;
2025        session_identity_migration::apply(&tx).map_err(|error| SqliteError::Migration {
2026            version: migration.version,
2027            error: error.to_string(),
2028        })?;
2029    } else {
2030        tx.execute_batch(migration.up)
2031            .map_err(|e| SqliteError::Migration {
2032                version: migration.version,
2033                error: e.to_string(),
2034            })?;
2035    }
2036
2037    let visibility_counts = if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
2038        Some(
2039            memory_visibility::cutover_counts(&tx).map_err(|error| SqliteError::Migration {
2040                version: migration.version,
2041                error: error.to_string(),
2042            })?,
2043        )
2044    } else {
2045        None
2046    };
2047
2048    // V19's repair contract includes normalizing the two known-divergent
2049    // recorded names. `_schema_migrations` is created and owned by this
2050    // runner (not by any migration file), so the normalization lives
2051    // here, in the same transaction that applies V19's SQL. Exact,
2052    // closed set — versions 13 and 14 only; any other (version, name)
2053    // mismatch still fails startup via validate_applied_migration_ledger.
2054    if migration.version == 19 {
2055        tx.execute_batch(
2056            "UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
2057             UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
2058        )
2059        .map_err(|e| SqliteError::Migration {
2060            version: migration.version,
2061            error: e.to_string(),
2062        })?;
2063    }
2064
2065    let now = chrono::Utc::now().timestamp_micros();
2066    tx.execute(
2067        "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
2068         ON CONFLICT(version) DO NOTHING",
2069        rusqlite::params![migration.version, migration.name, now],
2070    )
2071    .map_err(|e| SqliteError::Migration {
2072        version: migration.version,
2073        error: e.to_string(),
2074    })?;
2075
2076    #[cfg(test)]
2077    if instrumented_first_begin {
2078        test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
2079    }
2080
2081    #[cfg(test)]
2082    if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
2083        memory_visibility::test_state::stop_at(
2084            memory_visibility::test_state::Stop::BeforeCutoverCommit,
2085        )?;
2086    }
2087    tx.commit().map_err(|e| SqliteError::Migration {
2088        version: migration.version,
2089        error: e.to_string(),
2090    })?;
2091    if let Some(counts) = visibility_counts {
2092        memory_visibility::log_counts(&counts, conn.path().unwrap_or(":memory:"));
2093    }
2094    #[cfg(test)]
2095    if migration.version == 46 {
2096        memory_visibility::test_state::stop_at(
2097            memory_visibility::test_state::Stop::AfterV46Commit,
2098        )?;
2099    }
2100
2101    Ok(MigrationStep::Applied)
2102}
2103
2104#[derive(Debug)]
2105pub struct EmbeddingModelRegistryRecord {
2106    /// Vector engine name (e.g. `"paraphrase"`).
2107    pub engine_name: String,
2108    /// Model identifier (e.g. `"all-minilm-l6-v2"`).
2109    pub model_id: String,
2110    /// Canonical deduplication key combining engine and model.
2111    pub key_version: String,
2112    /// Embedding dimensionality.
2113    pub dimensions: u32,
2114    /// Lifecycle status (`"active"` or `"superseded"`).
2115    pub status: String,
2116    /// Epoch timestamp when the model was activated.
2117    pub activated_at: Option<i64>,
2118    /// Epoch timestamp when the model was superseded.
2119    pub superseded_at: Option<i64>,
2120}
2121
2122/// Query the `_embedding_models` registry.
2123///
2124/// Opens the database at `db` (defaults to `~/.khive/khive.db`) and
2125/// returns all registry rows, optionally filtered by `engine_name`.
2126/// Returns an empty vec if the database or table does not exist.
2127pub fn query_embedding_models(
2128    db: Option<&std::path::Path>,
2129    engine_filter: Option<&str>,
2130) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
2131    let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
2132        std::env::var("HOME")
2133            .map(std::path::PathBuf::from)
2134            .unwrap_or_else(|_| std::path::PathBuf::from("."))
2135            .join(".khive/khive.db")
2136    });
2137    if !path.exists() {
2138        return Ok(Vec::new());
2139    }
2140    let conn = Connection::open_with_flags(
2141        path,
2142        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
2143            | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
2144            | rusqlite::OpenFlags::SQLITE_OPEN_URI,
2145    )?;
2146    query_embedding_models_conn(&conn, engine_filter)
2147}
2148
2149/// Query `_embedding_models` from an existing connection (testable without a file).
2150///
2151/// Returns an empty vec if the table does not exist.
2152pub(crate) fn query_embedding_models_conn(
2153    conn: &Connection,
2154    engine_filter: Option<&str>,
2155) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
2156    let exists: bool = conn.query_row(
2157        "SELECT COUNT(*) > 0 FROM sqlite_master \
2158         WHERE type='table' AND name='_embedding_models'",
2159        [],
2160        |row| row.get(0),
2161    )?;
2162    if !exists {
2163        return Ok(Vec::new());
2164    }
2165
2166    let sql = if engine_filter.is_some() {
2167        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
2168         FROM _embedding_models WHERE engine_name = ?1 \
2169         ORDER BY engine_name, activated_at IS NULL, activated_at"
2170    } else {
2171        "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
2172         FROM _embedding_models \
2173         ORDER BY engine_name, activated_at IS NULL, activated_at"
2174    };
2175    let mut stmt = conn.prepare(sql)?;
2176    let map_row = |row: &rusqlite::Row<'_>| {
2177        let dim_raw: i64 = row.get(3)?;
2178        let dimensions = u32::try_from(dim_raw).map_err(|_| {
2179            rusqlite::Error::FromSqlConversionFailure(
2180                3,
2181                rusqlite::types::Type::Integer,
2182                Box::new(std::io::Error::other(format!(
2183                    "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
2184                    u32::MAX,
2185                ))),
2186            )
2187        })?;
2188        Ok(EmbeddingModelRegistryRecord {
2189            engine_name: row.get(0)?,
2190            model_id: row.get(1)?,
2191            key_version: row.get(2)?,
2192            dimensions,
2193            status: row.get(4)?,
2194            activated_at: row.get(5)?,
2195            superseded_at: row.get(6)?,
2196        })
2197    };
2198
2199    if let Some(engine) = engine_filter {
2200        stmt.query_map([engine], map_row)?
2201            .collect::<Result<Vec<_>, _>>()
2202            .map_err(Into::into)
2203    } else {
2204        stmt.query_map([], map_row)?
2205            .collect::<Result<Vec<_>, _>>()
2206            .map_err(Into::into)
2207    }
2208}
2209
2210// Test fixtures use the public policy constructor and a private lock namespace.
2211#[cfg(test)]
2212pub(crate) fn migration_test_policy() -> MigrationWritePolicy {
2213    MigrationWritePolicy::new(
2214        crate::DiskGuardEnvironment::capture()
2215            .resolve(None, None)
2216            .expect("test disk policy"),
2217        crate::PoolConfig::for_test()
2218            .volume_lock_dir
2219            .expect("test volume-lock directory"),
2220    )
2221    .expect("valid test migration policy")
2222}
2223
2224#[cfg(test)]
2225pub(crate) fn run_migrations_for_test(conn: &mut Connection) -> Result<u32, SqliteError> {
2226    run_migrations_with_policy(conn, &migration_test_policy())
2227}
2228
2229#[cfg(test)]
2230fn apply_schema_plan_for_test(
2231    conn: &mut Connection,
2232    plan: &ServiceSchemaPlan,
2233) -> Result<(), SqliteError> {
2234    apply_schema_plan_with_policy(conn, plan, &migration_test_policy())
2235}
2236
2237#[cfg(test)]
2238fn stage_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
2239    stage_attachment_cutover_with_policy(conn, &migration_test_policy())
2240}
2241
2242#[cfg(test)]
2243fn finalize_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
2244    finalize_attachment_cutover_with_policy(conn, &migration_test_policy())
2245}
2246
2247// =============================================================================
2248// Tests
2249// =============================================================================
2250
2251#[cfg(test)]
2252#[path = "entity_version_migration_measurement.rs"]
2253mod entity_version_measurement;
2254
2255#[cfg(test)]
2256#[path = "migrations_tests.rs"]
2257mod tests;
2258
2259#[cfg(test)]
2260#[path = "raw_migration_settlement_tests.rs"]
2261mod raw_settlement_tests;
2262
2263#[cfg(test)]
2264#[path = "git_note_index_migration_tests.rs"]
2265mod git_note_indexes;
2266
2267#[cfg(test)]
2268#[path = "entity_list_index_migration_tests.rs"]
2269mod entity_list_indexes;
2270
2271#[cfg(test)]
2272#[path = "schedule_core_index_migration_tests.rs"]
2273mod schedule_core_index_migration_tests;
2274
2275#[cfg(test)]
2276#[path = "comm_core_index_migration_tests.rs"]
2277mod comm_core_index_migration_tests;