use khive_storage::blob::ContentRef;
use rusqlite::{Connection, OptionalExtension};
use std::path::{Path, PathBuf};
use crate::error::SqliteError;
use crate::pool::WriteAdmission;
use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
#[path = "raw_migration_settlement.rs"]
mod raw_migration_settlement;
use raw_migration_settlement::{RawMigrationTransactions, RawMigrationWriteUnit};
pub(crate) trait MigrationTransactions {
fn admitted<T>(
&mut self,
operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
) -> Result<T, SqliteError>;
}
#[derive(Clone, Debug)]
pub struct MigrationWritePolicy {
disk_guard: crate::EffectiveDiskGuardConfig,
volume_lock_dir: PathBuf,
}
impl MigrationWritePolicy {
pub fn new(
disk_guard: crate::EffectiveDiskGuardConfig,
volume_lock_dir: impl Into<PathBuf>,
) -> Result<Self, SqliteError> {
disk_guard.validate()?;
let volume_lock_dir = volume_lock_dir.into();
if !volume_lock_dir.is_absolute() {
return Err(SqliteError::InvalidConfig(
"migration volume-lock directory must be absolute".to_string(),
));
}
Ok(Self {
disk_guard,
volume_lock_dir,
})
}
pub fn from_environment() -> Result<Self, SqliteError> {
let disk_guard = crate::DiskGuardEnvironment::capture().resolve(None, None)?;
Self::new(disk_guard, crate::default_volume_lock_dir()?)
}
pub fn disk_guard_config(&self) -> crate::EffectiveDiskGuardConfig {
self.disk_guard
}
pub fn volume_lock_dir(&self) -> &Path {
&self.volume_lock_dir
}
}
#[path = "session_identity_migration.rs"]
mod session_identity_migration;
mod memory_visibility;
pub struct Migration {
pub id: &'static str,
pub up_sql: &'static str,
pub down_sql: Option<&'static str>,
pub is_already_applied: Option<fn(&Connection) -> bool>,
}
pub struct ServiceSchemaPlan {
pub service: &'static str,
pub sqlite: &'static [Migration],
pub postgres: &'static [Migration],
}
const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
pub fn apply_schema_plan(
conn: &mut Connection,
plan: &ServiceSchemaPlan,
) -> Result<(), SqliteError> {
let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
apply_schema_plan_with_admission(
&mut RawMigrationTransactions::new(conn, &admission),
plan,
&admission,
)
}
pub fn apply_schema_plan_with_policy(
conn: &mut Connection,
plan: &ServiceSchemaPlan,
policy: &MigrationWritePolicy,
) -> Result<(), SqliteError> {
let admission =
WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
apply_schema_plan_with_admission(
&mut RawMigrationTransactions::new(conn, &admission),
plan,
&admission,
)
}
pub(crate) fn apply_schema_plan_with_admission(
writes: &mut impl MigrationTransactions,
plan: &ServiceSchemaPlan,
admission: &WriteAdmission,
) -> Result<(), SqliteError> {
writes.admitted(|conn| {
admission.check()?;
conn.execute_batch(SCHEMA_VERSION_TABLE)?;
require_autocommit(conn, "schema-version bootstrap")
})?;
for migration in plan.sqlite {
writes.admitted(|conn| apply_service_migration(conn, plan, migration, admission))?;
}
Ok(())
}
fn apply_service_migration(
conn: &Connection,
plan: &ServiceSchemaPlan,
migration: &Migration,
admission: &WriteAdmission,
) -> Result<(), SqliteError> {
let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
if let Err(error) = admission.check() {
let rollback = tx.rollback();
return Err(capacity_refusal_after_rollback(
conn,
rollback,
error,
"service schema migration",
));
}
if let Some(check) = migration.is_already_applied {
if check(&tx) {
return Ok(());
}
}
let already: bool = tx.query_row(
"SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
rusqlite::params![plan.service, migration.id],
|row| row.get(0),
)?;
if already {
return Ok(());
}
tx.execute_batch(migration.up_sql)?;
tx.execute(
"INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
rusqlite::params![
plan.service,
migration.id,
chrono::Utc::now().timestamp_micros(),
],
)?;
tx.commit()?;
Ok(())
}
fn require_autocommit(conn: &Connection, operation: &str) -> Result<(), SqliteError> {
if conn.is_autocommit() {
Ok(())
} else {
Err(SqliteError::InvalidData(format!(
"{operation} did not return the SQLite connection to autocommit"
)))
}
}
pub(crate) fn capacity_refusal_after_rollback(
conn: &Connection,
rollback: rusqlite::Result<()>,
refusal: SqliteError,
operation: &str,
) -> SqliteError {
if let Err(error) = rollback {
return SqliteError::InvalidData(format!(
"{operation} capacity refusal could not roll back: {error}; \
initial refusal: {refusal}"
));
}
if !conn.is_autocommit() {
return SqliteError::InvalidData(format!(
"{operation} capacity refusal rolled back without restoring autocommit; \
initial refusal: {refusal}"
));
}
refusal
}
pub struct VersionedMigration {
pub version: u32,
pub name: &'static str,
pub up: &'static str,
}
const V1_UP: &str = include_str!("../sql/schema.sql");
const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
const V22_UP: &str = include_str!("../sql/022-notes-unread-probe-recipient.sql");
const V23_UP: &str = include_str!("../sql/023-fts-record-kind.sql");
const V24_UP: &str = include_str!("../sql/024-fts-rowid-map.sql");
const V25_UP: &str = include_str!("../sql/025-notes-unread-probe-recipient-direction.sql");
const V26_UP: &str = include_str!("../sql/026-knowledge-fts-repair.sql");
const V27_UP: &str = include_str!("../sql/027-notes-hot-property-indexes.sql");
const V28_UP: &str = include_str!("../sql/028-notes-key.sql");
const V29_UP: &str = include_str!("../sql/029-note-streams.sql");
const V30_UP: &str = include_str!("../sql/030-tool-source-mounts.sql");
const V31_UP: &str = include_str!("../sql/031-note-versions.sql");
const V32_UP: &str = include_str!("../sql/032-knowledge-count-indexes.sql");
const V33_UP: &str = include_str!("../sql/033-notes-message-recipient-direction.sql");
const V34_UP: &str = include_str!("../sql/034-notes-namespace-created.sql");
const V35_UP: &str = include_str!("../sql/035-notes-unread-probe-recipient-type-direction.sql");
const V36_UP: &str = include_str!("../sql/036-events-operation-attribution.sql");
const V37_UP: &str = include_str!("../sql/037-entity-versions.sql");
const V38_UP: &str = include_str!("../sql/038-entities-legacy-type-index.sql");
const V39_UP: &str = include_str!("../sql/039-knowledge-cursor-indexes.sql");
const SESSION_IDENTITY_UP: &str = include_str!("../sql/040-session-source-scope.sql");
const SESSION_IDENTITY_MIGRATION_NAME: &str = "session_source_scoped_identity";
const V41_UP: &str = include_str!("../sql/041-sender-transport.sql");
const V42_UP: &str = include_str!("../sql/042-comm-external-id-channel-scope.sql");
const V43_UP: &str = include_str!("../sql/043-vector-provenance.sql");
const V44_COLUMNS: &str = include_str!("../sql/044-comm-outbound-due-a-columns.sql");
const V44_UP: &str = include_str!("../sql/044-comm-outbound-due-b-index.sql");
pub(crate) fn migrate_outbound_due_key(tx: &rusqlite::Transaction<'_>) -> rusqlite::Result<()> {
let has_column = |name: &str| -> rusqlite::Result<bool> {
tx.query_row(
"SELECT EXISTS(SELECT 1 FROM pragma_table_info('notes') WHERE name = ?1)",
[name],
|row| row.get(0),
)
};
match (has_column("strict_due_key")?, has_column("due_source")?) {
(false, false) => tx.execute_batch(V44_COLUMNS)?,
(true, true) => {}
_ => return Err(rusqlite::Error::InvalidQuery),
}
let mut after_id = String::new();
let mut first_page = true;
loop {
let rows: Vec<(String, String)> = {
let comparator = if first_page { ">=" } else { ">" };
let mut stmt = tx.prepare(&format!(
"SELECT id, json_extract(properties, '$.next_attempt_at') FROM notes \
WHERE id {comparator} ?1 AND json_type(properties, '$.next_attempt_at') = 'text' \
ORDER BY id LIMIT 500"
))?;
let collected = stmt
.query_map([&after_id], |row| Ok((row.get(0)?, row.get(1)?)))?
.collect::<rusqlite::Result<_>>()?;
collected
};
if rows.is_empty() {
break;
}
after_id = rows.last().expect("nonempty V44 page").0.clone();
first_page = false;
for (id, source) in rows {
if let Some(key) = crate::pool::strict_rfc3339_key(&source) {
tx.execute(
"UPDATE notes SET strict_due_key = ?1, due_source = ?2 \
WHERE id = ?3 AND (strict_due_key IS NOT ?1 OR due_source IS NOT ?2)",
rusqlite::params![key, source, id],
)?;
}
}
}
tx.execute_batch(V44_UP)
}
const RECIPIENT_TRANSPORT_VERSION: u32 = 45;
const V45_UP: &str = include_str!("../sql/045-recipient-transport.sql");
const V46_UP: &str = include_str!("../sql/046-memory-visibility-receipts.sql");
const V47_UP: &str = include_str!("../sql/047-attachment-role-quarantine.sql");
const V49_UP: &str = include_str!("../sql/049-git-note-property-indexes.sql");
const V50_UP: &str = include_str!("../sql/050-entity-list-plans.sql");
const V51_UP: &str = include_str!("../sql/051-schedule-core-indexes.sql");
const V52_UP: &str = include_str!("../sql/052-comm-core-indexes.sql");
const V53_UP: &str = include_str!("../sql/053-entity-kind-list-order.sql");
const MEMORY_VISIBILITY_CUTOVER_VERSION: u32 = 54;
const V54_UP: &str = include_str!("../sql/054-memory-visibility-epochs.sql");
const V48_UP: &str = include_str!("../sql/048-acknowledgement-journal-a-table.sql");
const ACKNOWLEDGEMENT_JOURNAL_INDEX: &str =
include_str!("../sql/048-acknowledgement-journal-b-index.sql");
pub(crate) fn migrate_acknowledgement_journal(
tx: &rusqlite::Transaction<'_>,
) -> rusqlite::Result<()> {
let columns: i64 = tx.query_row(
"SELECT count(*) FROM pragma_table_info('comm_ack_work') \
WHERE name IN ('attempt_count','not_before','retirement_reason')",
[],
|row| row.get(0),
)?;
match columns {
0 => tx.execute_batch(V48_UP)?,
3 => {}
_ => return Err(rusqlite::Error::InvalidQuery),
}
tx.execute_batch(ACKNOWLEDGEMENT_JOURNAL_INDEX)
}
const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
pub fn latest_schema_version() -> u32 {
MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
}
pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
pub const VECTOR_PROVENANCE_DDL: &str = V43_UP;
pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
pub const MIGRATIONS: &[VersionedMigration] = &[
VersionedMigration {
version: 1,
name: "initial_schema",
up: V1_UP,
},
VersionedMigration {
version: 2,
name: "narrow_fts_sections_update_trigger",
up: V2_UP,
},
VersionedMigration {
version: 3,
name: "backfill_domain_mirror_atoms",
up: V3_UP,
},
VersionedMigration {
version: 4,
name: "fts_consolidation",
up: V4_UP,
},
VersionedMigration {
version: 5,
name: "unique_comm_message_external_id",
up: V5_UP,
},
VersionedMigration {
version: 6,
name: "brain_retune_driver",
up: V6_UP,
},
VersionedMigration {
version: 7,
name: "notes_seq",
up: V7_UP,
},
VersionedMigration {
version: 8,
name: "notes_seq_repair",
up: V8_UP,
},
VersionedMigration {
version: 9,
name: "entities_name_ci_index",
up: V9_UP,
},
VersionedMigration {
version: 10,
name: "entities_content_ref",
up: V10_UP,
},
VersionedMigration {
version: 11,
name: "ann_write_log",
up: V11_UP,
},
VersionedMigration {
version: 12,
name: "ann_write_log_model_seq_index",
up: V12_UP,
},
VersionedMigration {
version: 13,
name: "list_cursor_sequences",
up: V13_UP,
},
VersionedMigration {
version: 14,
name: "graph_edges_id_unique",
up: V14_UP,
},
VersionedMigration {
version: 15,
name: "serve_ledger_attribution",
up: V15_UP,
},
VersionedMigration {
version: 16,
name: "gtd_dependency_cycle_guards",
up: V16_UP,
},
VersionedMigration {
version: 17,
name: "agents_ddl",
up: V17_UP,
},
VersionedMigration {
version: 18,
name: "ann_consumer_pending",
up: V18_UP,
},
VersionedMigration {
version: 19,
name: "list_cursor_backfill_repair",
up: V19_UP,
},
VersionedMigration {
version: 20,
name: "blob_gc_claims",
up: V20_UP,
},
VersionedMigration {
version: ATTACHMENT_CUTOVER_VERSION,
name: "attachments_first_class",
up: V21_STAGE_UP,
},
VersionedMigration {
version: 22,
name: "notes_unread_probe_recipient",
up: V22_UP,
},
VersionedMigration {
version: 23,
name: "fts_record_kind",
up: V23_UP,
},
VersionedMigration {
version: 24,
name: "fts_rowid_map",
up: V24_UP,
},
VersionedMigration {
version: 25,
name: "notes_unread_probe_recipient_direction",
up: V25_UP,
},
VersionedMigration {
version: 26,
name: "knowledge_fts_repair",
up: V26_UP,
},
VersionedMigration {
version: 27,
name: "notes_hot_property_indexes",
up: V27_UP,
},
VersionedMigration {
version: 28,
name: "notes_key",
up: V28_UP,
},
VersionedMigration {
version: 29,
name: "note_streams",
up: V29_UP,
},
VersionedMigration {
version: 30,
name: "tool_source_mounts",
up: V30_UP,
},
VersionedMigration {
version: 31,
name: "note_versions",
up: V31_UP,
},
VersionedMigration {
version: 32,
name: "knowledge_count_indexes",
up: V32_UP,
},
VersionedMigration {
version: 33,
name: "notes_message_recipient_direction",
up: V33_UP,
},
VersionedMigration {
version: 34,
name: "notes_namespace_created",
up: V34_UP,
},
VersionedMigration {
version: 35,
name: "notes_unread_probe_recipient_type_direction",
up: V35_UP,
},
VersionedMigration {
version: 36,
name: "events_operation_attribution",
up: V36_UP,
},
VersionedMigration {
version: 37,
name: "entity_versions",
up: V37_UP,
},
VersionedMigration {
version: 38,
name: "entities_legacy_type_index",
up: V38_UP,
},
VersionedMigration {
version: 39,
name: "knowledge_cursor_indexes",
up: V39_UP,
},
VersionedMigration {
version: 40,
name: SESSION_IDENTITY_MIGRATION_NAME,
up: SESSION_IDENTITY_UP,
},
VersionedMigration {
version: 41,
name: "sender_transport",
up: V41_UP,
},
VersionedMigration {
version: 42,
name: "comm_external_id_channel_scope",
up: V42_UP,
},
VersionedMigration {
version: 43,
name: "vector_provenance",
up: V43_UP,
},
VersionedMigration {
version: 44,
name: "comm_outbound_due",
up: V44_UP,
},
VersionedMigration {
version: RECIPIENT_TRANSPORT_VERSION,
name: "recipient_transport",
up: V45_UP,
},
VersionedMigration {
version: 46,
name: "memory_visibility_receipts",
up: V46_UP,
},
VersionedMigration {
version: 47,
name: "attachment_role_quarantine",
up: V47_UP,
},
VersionedMigration {
version: 48,
name: "acknowledgement_journal",
up: V48_UP,
},
VersionedMigration {
version: 49,
name: "git_note_property_indexes",
up: V49_UP,
},
VersionedMigration {
version: 50,
name: "entity_list_plans",
up: V50_UP,
},
VersionedMigration {
version: 51,
name: "schedule_core_indexes",
up: V51_UP,
},
VersionedMigration {
version: 52,
name: "comm_core_indexes",
up: V52_UP,
},
VersionedMigration {
version: 53,
name: "entity_kind_list_order",
up: V53_UP,
},
VersionedMigration {
version: MEMORY_VISIBILITY_CUTOVER_VERSION,
name: "memory_visibility_epochs",
up: V54_UP,
},
];
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum AttachmentCutoverStatus {
Pending,
Incomplete,
Complete,
}
fn schema_object_exists(
conn: &Connection,
object_type: &str,
name: &str,
) -> Result<bool, SqliteError> {
conn.query_row(
"SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
rusqlite::params![object_type, name],
|row| row.get(0),
)
.map_err(Into::into)
}
fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
conn.query_row(
"SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
rusqlite::params![table, column],
|row| row.get(0),
)
.map_err(Into::into)
}
fn require_attachment_schema_objects(
conn: &Connection,
objects: &[(&str, &str)],
phase: &str,
) -> Result<(), SqliteError> {
for (object_type, name) in objects {
if !schema_object_exists(conn, object_type, name)? {
return Err(SqliteError::InvalidData(format!(
"attachment cutover {phase} state is missing {object_type} {name:?}"
)));
}
}
Ok(())
}
fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
require_attachment_schema_objects(
conn,
&[
("table", "attachments"),
("index", "idx_attachments_content_ref"),
],
"incomplete",
)?;
require_legacy_attachment_fences(conn)
}
fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
require_attachment_schema_objects(
conn,
&[
("table", "attachments"),
("table", "blob_gc_claims"),
("index", "idx_attachments_content_ref"),
("index", "idx_blob_gc_claims_content_ref"),
("trigger", "attachments_reject_claimed_blob_insert"),
("trigger", "attachments_reject_claimed_blob_update"),
],
"complete",
)?;
if schema_column_exists(conn, "entities", "content_ref")? {
return Err(SqliteError::InvalidData(
"attachment cutover is complete but entities.content_ref still exists".into(),
));
}
for (object_type, name) in [
("index", "idx_entities_content_ref"),
("trigger", "entities_reject_claimed_blob_insert"),
("trigger", "entities_reject_claimed_blob_update"),
] {
if schema_object_exists(conn, object_type, name)? {
return Err(SqliteError::InvalidData(format!(
"attachment cutover is complete but legacy {object_type} {name:?} still exists"
)));
}
}
Ok(())
}
pub fn attachment_cutover_status(
conn: &Connection,
) -> Result<AttachmentCutoverStatus, SqliteError> {
let version = read_schema_version(conn)?;
let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
if !marker_table {
if version >= ATTACHMENT_CUTOVER_VERSION {
return Err(SqliteError::InvalidData(format!(
"migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
)));
}
if schema_object_exists(conn, "table", "attachments")? {
return Err(SqliteError::InvalidData(
"attachments table exists without the durable attachment cutover marker".into(),
));
}
return Ok(AttachmentCutoverStatus::Pending);
}
let marker: Option<(String, Option<i64>)> = conn
.query_row(
"SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
match marker {
Some((state, None)) if state == "incomplete" => {
if version >= ATTACHMENT_CUTOVER_VERSION {
Err(SqliteError::InvalidData(format!(
"attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
)))
} else {
validate_incomplete_attachment_schema(conn)?;
Ok(AttachmentCutoverStatus::Incomplete)
}
}
Some((state, Some(_))) if state == "complete" => {
if version >= ATTACHMENT_CUTOVER_VERSION {
validate_complete_attachment_schema(conn)?;
Ok(AttachmentCutoverStatus::Complete)
} else {
Err(SqliteError::InvalidData(format!(
"attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
)))
}
}
Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
"invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
))),
None => Err(SqliteError::InvalidData(
"attachment cutover marker table exists without its singleton row".into(),
)),
}
}
fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
if !schema_column_exists(conn, "entities", "content_ref")? {
return Err(SqliteError::InvalidData(
"attachment cutover requires legacy entities.content_ref until finalization".into(),
));
}
for (object_type, name) in [
("table", "blob_gc_claims"),
("index", "idx_blob_gc_claims_content_ref"),
("index", "idx_entities_content_ref"),
("trigger", "entities_reject_claimed_blob_insert"),
("trigger", "entities_reject_claimed_blob_update"),
] {
if !schema_object_exists(conn, object_type, name)? {
return Err(SqliteError::InvalidData(format!(
"attachment cutover requires legacy {object_type} {name:?} until finalization"
)));
}
}
Ok(())
}
fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
if !(1..=4).contains(&width) {
return Err(SqliteError::InvalidData(format!(
"the text-encoding width probe returned {width}; refusing canonicality validation"
)));
}
Ok(width * 64)
}
fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
let canonical_bytes = canonical_content_ref_byte_width(conn)?;
let invalid: Option<String> = conn
.query_row(
"SELECT id FROM entities \
WHERE content_ref IS NOT NULL \
AND (typeof(content_ref) <> 'text' \
OR length(content_ref) <> 64 \
OR length(CAST(content_ref AS BLOB)) <> ?1 \
OR content_ref GLOB '*[^0-9a-f]*') \
LIMIT 1",
[canonical_bytes],
|row| row.get(0),
)
.optional()?;
if let Some(id) = invalid {
return Err(SqliteError::InvalidData(format!(
"entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
)));
}
Ok(())
}
fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
let canonical_bytes = canonical_content_ref_byte_width(conn)?;
for (table, identity) in [
("attachments", "record_uuid"),
("blob_gc_claims", "root_key"),
] {
let sql = format!(
"SELECT {identity} FROM {table} \
WHERE typeof(content_ref) <> 'text' \
OR length(content_ref) <> 64 \
OR length(CAST(content_ref AS BLOB)) <> ?1 \
OR content_ref GLOB '*[^0-9a-f]*' \
LIMIT 1"
);
let invalid: Option<String> = conn
.query_row(&sql, [canonical_bytes], |row| row.get(0))
.optional()?;
if let Some(owner) = invalid {
return Err(SqliteError::InvalidData(format!(
"{table}.content_ref for {identity} {owner:?} is not canonical"
)));
}
}
Ok(())
}
fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
let dangling: Option<(String, String)> = conn
.query_row(
"SELECT record_uuid, substrate FROM attachments AS attachment \
WHERE (substrate = 'entity' AND NOT EXISTS ( \
SELECT 1 FROM entities WHERE id = attachment.record_uuid \
)) \
OR (substrate = 'note' AND NOT EXISTS ( \
SELECT 1 FROM notes WHERE id = attachment.record_uuid \
)) \
LIMIT 1",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
if let Some((record_uuid, substrate)) = dangling {
return Err(SqliteError::InvalidData(format!(
"attachment role references absent {substrate} record {record_uuid:?}"
)));
}
Ok(())
}
fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
let conflict: Option<String> = conn
.query_row(
"SELECT entity.id FROM entities AS entity \
LEFT JOIN attachments AS attachment \
ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
WHERE entity.content_ref IS NOT NULL \
AND (attachment.record_uuid IS NULL \
OR attachment.substrate <> 'entity' \
OR attachment.content_ref <> entity.content_ref) \
LIMIT 1",
[],
|row| row.get(0),
)
.optional()?;
if let Some(record_uuid) = conflict {
return Err(SqliteError::InvalidData(format!(
"legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
)));
}
Ok(())
}
fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
require_legacy_attachment_fences(conn)?;
conn.execute_batch(V21_STAGE_UP)?;
validate_canonical_legacy_refs(conn)?;
validate_canonical_attachment_and_claim_refs(conn)?;
let conflict: Option<String> = conn
.query_row(
"SELECT entity.id FROM entities AS entity \
JOIN attachments AS attachment \
ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
WHERE entity.content_ref IS NOT NULL \
AND (attachment.substrate <> 'entity' \
OR attachment.content_ref <> entity.content_ref) \
LIMIT 1",
[],
|row| row.get(0),
)
.optional()?;
if let Some(record_uuid) = conflict {
return Err(SqliteError::InvalidData(format!(
"existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
)));
}
conn.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
FROM entities WHERE content_ref IS NOT NULL \
ON CONFLICT(record_uuid, role) DO NOTHING",
[],
)?;
validate_legacy_content_backfill(conn)?;
conn.execute("DELETE FROM blob_gc_claims", [])?;
conn.execute(
"INSERT INTO attachment_cutover_state \
(singleton, state, started_at, completed_at) \
VALUES (1, 'incomplete', ?1, NULL) \
ON CONFLICT(singleton) DO NOTHING",
[now],
)?;
Ok(())
}
pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
RawMigrationWriteUnit::new(conn, &admission)?
.run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
}
pub fn stage_attachment_cutover_with_policy(
conn: &mut Connection,
policy: &MigrationWritePolicy,
) -> Result<(), SqliteError> {
let admission =
WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
RawMigrationWriteUnit::new(conn, &admission)?
.run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
}
pub(crate) fn stage_attachment_cutover_with_admission(
conn: &mut Connection,
admission: &WriteAdmission,
) -> Result<(), SqliteError> {
match attachment_cutover_status(conn)? {
AttachmentCutoverStatus::Complete => return Ok(()),
AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
}
if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
return Err(SqliteError::InvalidData(format!(
"attachment cutover stage requires canonical V{} schema",
ATTACHMENT_CUTOVER_VERSION - 1
)));
}
let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
if let Err(error) = admission.check() {
let rollback = tx.rollback();
return Err(capacity_refusal_after_rollback(
conn,
rollback,
error,
"attachment cutover stage",
));
}
let status = attachment_cutover_status(&tx)?;
if status == AttachmentCutoverStatus::Complete {
return Ok(());
}
stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
tx.commit()?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn apply_generic_verified_attachment(
conn: &Connection,
record_uuid: &str,
substrate: &str,
role: &str,
content_ref: &ContentRef,
media_type: Option<&str>,
size_bytes: Option<u64>,
created_at: i64,
) -> Result<(), SqliteError> {
if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
return Err(SqliteError::InvalidData(
"verified application attachments may only be applied while V21 cutover is incomplete"
.into(),
));
}
if role.is_empty() || role.chars().any(char::is_control) {
return Err(SqliteError::InvalidData(
"attachment role must be non-empty and contain no control characters".into(),
));
}
let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
})?;
let owner_table = match substrate {
"entity" => "entities",
"note" => "notes",
other => {
return Err(SqliteError::InvalidData(format!(
"attachment substrate must be 'entity' or 'note', got {other:?}"
)))
}
};
let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
if !owner_exists {
return Err(SqliteError::InvalidData(format!(
"cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
)));
}
let claimed: bool = conn.query_row(
"SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
[content_ref.as_str()],
|row| row.get(0),
)?;
if claimed {
return Err(SqliteError::InvalidData(format!(
"cannot attach claimed content_ref {} during V21 cutover",
content_ref.as_str()
)));
}
let changed = conn.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
ON CONFLICT(record_uuid, role) DO UPDATE SET \
media_type = excluded.media_type, \
size_bytes = excluded.size_bytes, \
created_at = excluded.created_at \
WHERE attachments.substrate = excluded.substrate \
AND attachments.content_ref = excluded.content_ref",
rusqlite::params![
record_uuid,
substrate,
role,
content_ref.as_str(),
media_type,
size_bytes,
created_at,
],
)?;
if changed == 0 {
return Err(SqliteError::InvalidData(format!(
"attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
)));
}
Ok(())
}
fn finalize_attachment_cutover_on_connection(
conn: &Connection,
now: i64,
) -> Result<(), SqliteError> {
require_legacy_attachment_fences(conn)?;
validate_canonical_legacy_refs(conn)?;
validate_canonical_attachment_and_claim_refs(conn)?;
validate_attachment_record_owners(conn)?;
validate_legacy_content_backfill(conn)?;
let remaining_claims: i64 =
conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
if remaining_claims != 0 {
return Err(SqliteError::InvalidData(format!(
"attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
)));
}
let uncovered_model: Option<String> = conn
.query_row(
"SELECT model.id FROM entities AS model \
WHERE model.entity_type = 'moodboard_model' \
AND model.content_ref IS NOT NULL \
AND NOT EXISTS ( \
SELECT 1 FROM attachments AS attachment \
WHERE attachment.record_uuid = model.id \
AND attachment.substrate = 'entity' \
AND attachment.role = 'fann-network' \
) \
LIMIT 1",
[],
|row| row.get(0),
)
.optional()?;
if let Some(record_uuid) = uncovered_model {
return Err(SqliteError::InvalidData(format!(
"moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
)));
}
conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
conn.execute_batch(
"DROP TRIGGER entities_reject_claimed_blob_insert; \
DROP TRIGGER entities_reject_claimed_blob_update; \
DROP INDEX idx_entities_content_ref; \
ALTER TABLE entities DROP COLUMN content_ref;",
)?;
conn.execute(
"UPDATE attachment_cutover_state \
SET state = 'complete', completed_at = ?1 \
WHERE singleton = 1 AND state = 'incomplete'",
[now],
)?;
Ok(())
}
fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
let migration = MIGRATIONS
.iter()
.find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
.expect("V21 migration must be registered");
conn.execute(
"INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
rusqlite::params![migration.version, migration.name, now],
)?;
Ok(())
}
pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
RawMigrationWriteUnit::new(conn, &admission)?
.run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
}
pub fn finalize_attachment_cutover_with_policy(
conn: &mut Connection,
policy: &MigrationWritePolicy,
) -> Result<(), SqliteError> {
let admission =
WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
RawMigrationWriteUnit::new(conn, &admission)?
.run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
}
pub(crate) fn finalize_attachment_cutover_with_admission(
conn: &mut Connection,
admission: &WriteAdmission,
) -> Result<(), SqliteError> {
if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
return Ok(());
}
let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
if let Err(error) = admission.check() {
let rollback = tx.rollback();
return Err(capacity_refusal_after_rollback(
conn,
rollback,
error,
"attachment cutover finalization",
));
}
match attachment_cutover_status(&tx)? {
AttachmentCutoverStatus::Complete => return Ok(()),
AttachmentCutoverStatus::Pending => {
return Err(SqliteError::InvalidData(
"attachment cutover must complete stage 1 before finalization".into(),
))
}
AttachmentCutoverStatus::Incomplete => {}
}
let now = chrono::Utc::now().timestamp_micros();
finalize_attachment_cutover_on_connection(&tx, now)?;
record_attachment_cutover_migration(&tx, now)?;
tx.commit()?;
Ok(())
}
fn read_applied_migration_ledger(
conn: &Connection,
through_version: u32,
) -> Result<Vec<(u32, String)>, SqliteError> {
let mut stmt = conn.prepare(
"SELECT version, name FROM _schema_migrations \
WHERE version <= ?1 ORDER BY version ASC",
)?;
let rows = stmt
.query_map([through_version], |row| {
Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
})?
.collect::<Result<Vec<_>, _>>()?;
Ok(rows)
}
fn validate_applied_migration_versions(
applied: &[(u32, String)],
through_version: u32,
) -> Result<(), SqliteError> {
let expected: Vec<&VersionedMigration> = MIGRATIONS
.iter()
.filter(|migration| migration.version <= through_version)
.collect();
let mut applied_index = 0;
for migration in expected {
let Some((version, applied_name)) = applied.get(applied_index) else {
return Err(SqliteError::InvalidData(format!(
"migration history is missing version {} ('{}'); the applied ledger must be \
the exact contiguous canonical sequence through version {through_version}",
migration.version, migration.name,
)));
};
if *version < migration.version {
return Err(SqliteError::InvalidData(format!(
"migration history contains unknown version {version} recorded as \
'{applied_name}'; the applied ledger must contain only canonical versions"
)));
}
if *version > migration.version {
return Err(SqliteError::InvalidData(format!(
"migration history is missing version {} ('{}'); found version {version} \
next instead",
migration.version, migration.name,
)));
}
applied_index += 1;
}
if let Some((version, name)) = applied.get(applied_index) {
return Err(SqliteError::InvalidData(format!(
"migration history contains unknown version {version} recorded as '{name}'; \
the applied ledger must contain only canonical versions"
)));
}
Ok(())
}
fn validate_applied_migration_names(
applied: &[(u32, String)],
through_version: u32,
allow_known_v19_repairs: bool,
) -> Result<(), SqliteError> {
for ((version, applied_name), migration) in applied.iter().zip(
MIGRATIONS
.iter()
.filter(|migration| migration.version <= through_version),
) {
debug_assert_eq!(*version, migration.version);
if migration.name != applied_name.as_str() {
if allow_known_v19_repairs && matches!(*version, 13 | 14) {
continue;
}
return Err(SqliteError::InvalidData(format!(
"migration version {version} is recorded under name '{applied_name}', \
expected '{expected}'. This database's migration history does not match \
the current binary; recreate it from the current schema or repair the \
specific known divergence via a dedicated migration.",
expected = migration.name,
)));
}
}
Ok(())
}
fn validate_applied_migration_ledger(
conn: &Connection,
through_version: u32,
) -> Result<(), SqliteError> {
let applied = read_applied_migration_ledger(conn, through_version)?;
validate_applied_migration_versions(&applied, through_version)?;
validate_applied_migration_names(&applied, through_version, false)
}
const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
match conn.query_row(
"SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
[],
|row| row.get(0),
) {
Ok(version) => Ok(version),
Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
if msg.contains("no such table: _schema_migrations") =>
{
Ok(0)
}
Err(e) => Err(e.into()),
}
}
pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
let conn = crate::pool::open_read_only_snapshot_connection(path)?;
read_schema_version(&conn)
}
pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
let conn = crate::pool::open_read_only_snapshot_connection(path)?;
validate_schema_is_current(&conn)
}
pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
let current_version = read_schema_version(conn)?;
let latest_version = latest_schema_version();
if current_version < latest_version {
return Err(SqliteError::InvalidData(format!(
"read-only database schema version {current_version} is behind the latest known \
migration {latest_version}; migrate a writable copy with this build before opening \
the snapshot read-only"
)));
}
if current_version > latest_version {
return Err(SqliteError::InvalidData(format!(
"read-only database schema version {current_version} is ahead of the latest known \
migration {latest_version}; use a compatible newer build or recreate the snapshot"
)));
}
validate_applied_migration_ledger(conn, current_version)?;
if current_version >= ATTACHMENT_CUTOVER_VERSION
&& attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
{
return Err(SqliteError::InvalidData(
"read-only database has not completed the V21 attachment cutover".into(),
));
}
if current_version >= MEMORY_VISIBILITY_CUTOVER_VERSION {
memory_visibility::validate_cutover(conn)?;
}
Ok(current_version)
}
pub fn validate_memory_visibility_cutover(conn: &Connection) -> Result<(), SqliteError> {
validate_schema_is_current(conn).map(|_| ())
}
#[cfg(test)]
pub(crate) mod test_sync {
use std::sync::atomic::AtomicU32;
use std::sync::{Arc, Barrier, Mutex};
pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
pub(crate) fn record_busy(_count: i32) -> bool {
BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(1));
true
}
pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
std::thread_local! {
pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
const { std::cell::Cell::new(false) };
pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
const { std::cell::Cell::new(false) };
}
}
fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
let configured = conn.path().unwrap_or_default();
let raw_path = if configured.is_empty() {
conn.query_row(
"SELECT file FROM pragma_database_list WHERE name = 'main'",
[],
|row| row.get::<_, String>(0),
)?
} else {
configured.to_string()
};
if raw_path.is_empty() {
return Ok(None);
}
let canonical = std::fs::canonicalize(&raw_path).map_err(SqliteError::Io)?;
#[cfg(any(unix, windows))]
crate::pool::opened_sqlite_file_identity(conn, &canonical)?;
Ok(Some(canonical))
}
fn validate_database_gc_owner(
conn: &Connection,
owner: &DatabaseGcOwnerGuard,
) -> Result<(), SqliteError> {
let connection_path = canonical_connection_database_path(conn)?;
if owner.database_path() != connection_path.as_deref() {
return Err(SqliteError::InvalidData(format!(
"database GC owner targets {:?}, but migration connection targets {:?}",
owner.database_path(),
connection_path.as_deref(),
)));
}
Ok(())
}
pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
let database_path = canonical_connection_database_path(conn)?;
let admission = WriteAdmission::for_canonical_path(database_path.clone())?;
run_raw_migrations_with_admission(conn, database_path, &admission)
}
pub fn run_migrations_with_policy(
conn: &mut Connection,
policy: &MigrationWritePolicy,
) -> Result<u32, SqliteError> {
let database_path = canonical_connection_database_path(conn)?;
let admission = WriteAdmission::for_migration_policy(database_path.clone(), policy)?;
run_raw_migrations_with_admission(conn, database_path, &admission)
}
fn run_raw_migrations_with_admission(
conn: &mut Connection,
database_path: Option<PathBuf>,
admission: &WriteAdmission,
) -> Result<u32, SqliteError> {
if let Some(database_path) = database_path {
let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
SqliteError::InvalidData(format!(
"failed to acquire database GC owner before schema migration: {error}"
))
})?;
return run_migrations_with_database_gc_owner(
&mut RawMigrationTransactions::new(conn, admission),
&owner,
admission,
);
}
run_versioned_migrations(
&mut RawMigrationTransactions::new(conn, admission),
None,
admission,
)
}
pub(crate) fn run_migrations_with_database_gc_owner(
writes: &mut impl MigrationTransactions,
owner: &DatabaseGcOwnerGuard,
admission: &WriteAdmission,
) -> Result<u32, SqliteError> {
run_versioned_migrations(writes, Some(owner), admission)
}
fn with_migration_busy_timeout<T>(
conn: &mut Connection,
operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
) -> Result<T, SqliteError> {
let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
let raised = prior_busy_ms < 5_000;
if raised {
conn.busy_timeout(std::time::Duration::from_secs(5))?;
}
let result = operation(conn);
if raised {
let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
}
result
}
enum MigrationStep {
Applied,
AppliedBySibling(u32),
Deferred,
}
fn run_versioned_migrations(
writes: &mut impl MigrationTransactions,
owner: Option<&DatabaseGcOwnerGuard>,
admission: &WriteAdmission,
) -> Result<u32, SqliteError> {
let current_version = writes.admitted(|conn| {
if let Some(owner) = owner {
validate_database_gc_owner(conn, owner)?;
}
with_migration_busy_timeout(conn, |conn| bootstrap_migration_ledger(conn, admission))
})?;
#[cfg(test)]
if test_sync::PARTICIPATE.with(|p| p.get()) {
let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
if let Some(barrier) = barrier {
barrier.wait();
}
}
let latest_version = latest_schema_version();
let mut applied_version = current_version;
let mut skip_through = current_version;
for migration in MIGRATIONS {
if migration.version <= skip_through {
applied_version = applied_version.max(migration.version);
continue;
}
let step = writes.admitted(|conn| {
with_migration_busy_timeout(conn, |conn| {
apply_versioned_migration(conn, migration, admission, latest_version)
})
})?;
match step {
MigrationStep::Applied => applied_version = migration.version,
MigrationStep::AppliedBySibling(sibling_version) => {
skip_through = sibling_version.min(latest_version);
applied_version = applied_version.max(migration.version);
}
MigrationStep::Deferred => break,
}
}
writes.admitted(|conn| validate_applied_migration_ledger(conn, applied_version))?;
Ok(applied_version)
}
fn bootstrap_migration_ledger(
conn: &mut Connection,
admission: &WriteAdmission,
) -> Result<u32, SqliteError> {
admission.check()?;
conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
require_autocommit(conn, "migration-tracking bootstrap")?;
let current_version: u32 = read_schema_version(conn)?;
let latest_version = latest_schema_version();
if current_version > latest_version {
return Err(SqliteError::InvalidData(format!(
"database schema version {current_version} is ahead of the latest known migration \
{latest_version}. This database predates the consolidated baseline (ADR-015) or was \
written by a newer build. Recreate it from the current schema; in-place downgrade is \
not supported."
)));
}
let applied = read_applied_migration_ledger(conn, current_version)?;
validate_applied_migration_versions(&applied, current_version)?;
validate_applied_migration_names(&applied, current_version, current_version < 19)?;
Ok(current_version)
}
fn apply_versioned_migration(
conn: &mut Connection,
migration: &VersionedMigration,
admission: &WriteAdmission,
latest_version: u32,
) -> Result<MigrationStep, SqliteError> {
#[cfg(test)]
let instrumented_first_begin =
test_sync::PARTICIPATE.with(|p| p.get()) && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
#[cfg(test)]
if instrumented_first_begin {
test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
}
#[cfg(test)]
if test_sync::PARTICIPATE.with(|p| p.get()) {
conn.busy_handler(Some(test_sync::record_busy))?;
}
let tx = conn
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
if let Err(error) = admission.check() {
let rollback = tx.rollback();
return Err(capacity_refusal_after_rollback(
conn,
rollback,
error,
"core schema migration",
));
}
let sibling_version: u32 = tx
.query_row(
"SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
[],
|row| row.get(0),
)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
#[cfg(test)]
if instrumented_first_begin {
use std::sync::atomic::Ordering::SeqCst;
if sibling_version == 0 {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline {
std::thread::yield_now();
}
} else {
test_sync::LOSER_SAW_WINNER_COMMIT
.store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
}
}
if sibling_version > latest_version {
return Err(SqliteError::InvalidData(format!(
"database schema version {sibling_version} is ahead of the latest known \
migration {latest_version} (committed by a concurrent process while this \
one waited for the migration write lock). This build cannot run against \
the newer schema; upgrade the binary or recreate the database."
)));
}
if sibling_version >= migration.version {
#[cfg(test)]
test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
return Ok(MigrationStep::AppliedBySibling(sibling_version));
}
if migration.version == 46 {
memory_visibility::capture_pre_v46(&tx).map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?;
#[cfg(test)]
memory_visibility::test_state::stop_at(memory_visibility::test_state::Stop::AfterCapture)?;
}
if migration.version == ATTACHMENT_CUTOVER_VERSION {
let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
let legacy_refs: i64 = tx
.query_row(
"SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
[],
|row| row.get(0),
)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
drop(tx);
return Ok(MigrationStep::Deferred);
}
if status != AttachmentCutoverStatus::Pending {
return Err(SqliteError::Migration {
version: migration.version,
error: format!("unexpected attachment cutover state {status:?}"),
});
}
let now = chrono::Utc::now().timestamp_micros();
stage_attachment_cutover_on_connection(&tx, now).map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
SqliteError::Migration {
version: migration.version,
error: e.to_string(),
}
})?;
} else if migration.version == 36 {
crate::stores::event::ensure_operation_attribution_columns(&tx).map_err(|error| {
SqliteError::Migration {
version: migration.version,
error: error.to_string(),
}
})?;
} else if migration.version == 44 {
migrate_outbound_due_key(&tx).map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?;
} else if migration.version == 48 {
migrate_acknowledgement_journal(&tx).map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?;
} else if migration.name == SESSION_IDENTITY_MIGRATION_NAME {
tx.execute_batch(migration.up)
.map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?;
session_identity_migration::apply(&tx).map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?;
} else {
tx.execute_batch(migration.up)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
}
let visibility_counts = if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
Some(
memory_visibility::cutover_counts(&tx).map_err(|error| SqliteError::Migration {
version: migration.version,
error: error.to_string(),
})?,
)
} else {
None
};
if migration.version == 19 {
tx.execute_batch(
"UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
}
let now = chrono::Utc::now().timestamp_micros();
tx.execute(
"INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
ON CONFLICT(version) DO NOTHING",
rusqlite::params![migration.version, migration.name, now],
)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
#[cfg(test)]
if instrumented_first_begin {
test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
}
#[cfg(test)]
if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
memory_visibility::test_state::stop_at(
memory_visibility::test_state::Stop::BeforeCutoverCommit,
)?;
}
tx.commit().map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
if let Some(counts) = visibility_counts {
memory_visibility::log_counts(&counts, conn.path().unwrap_or(":memory:"));
}
#[cfg(test)]
if migration.version == 46 {
memory_visibility::test_state::stop_at(
memory_visibility::test_state::Stop::AfterV46Commit,
)?;
}
Ok(MigrationStep::Applied)
}
#[derive(Debug)]
pub struct EmbeddingModelRegistryRecord {
pub engine_name: String,
pub model_id: String,
pub key_version: String,
pub dimensions: u32,
pub status: String,
pub activated_at: Option<i64>,
pub superseded_at: Option<i64>,
}
pub fn query_embedding_models(
db: Option<&std::path::Path>,
engine_filter: Option<&str>,
) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
std::env::var("HOME")
.map(std::path::PathBuf::from)
.unwrap_or_else(|_| std::path::PathBuf::from("."))
.join(".khive/khive.db")
});
if !path.exists() {
return Ok(Vec::new());
}
let conn = Connection::open_with_flags(
path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
| rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
| rusqlite::OpenFlags::SQLITE_OPEN_URI,
)?;
query_embedding_models_conn(&conn, engine_filter)
}
pub(crate) fn query_embedding_models_conn(
conn: &Connection,
engine_filter: Option<&str>,
) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
let exists: bool = conn.query_row(
"SELECT COUNT(*) > 0 FROM sqlite_master \
WHERE type='table' AND name='_embedding_models'",
[],
|row| row.get(0),
)?;
if !exists {
return Ok(Vec::new());
}
let sql = if engine_filter.is_some() {
"SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
FROM _embedding_models WHERE engine_name = ?1 \
ORDER BY engine_name, activated_at IS NULL, activated_at"
} else {
"SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
FROM _embedding_models \
ORDER BY engine_name, activated_at IS NULL, activated_at"
};
let mut stmt = conn.prepare(sql)?;
let map_row = |row: &rusqlite::Row<'_>| {
let dim_raw: i64 = row.get(3)?;
let dimensions = u32::try_from(dim_raw).map_err(|_| {
rusqlite::Error::FromSqlConversionFailure(
3,
rusqlite::types::Type::Integer,
Box::new(std::io::Error::other(format!(
"_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
u32::MAX,
))),
)
})?;
Ok(EmbeddingModelRegistryRecord {
engine_name: row.get(0)?,
model_id: row.get(1)?,
key_version: row.get(2)?,
dimensions,
status: row.get(4)?,
activated_at: row.get(5)?,
superseded_at: row.get(6)?,
})
};
if let Some(engine) = engine_filter {
stmt.query_map([engine], map_row)?
.collect::<Result<Vec<_>, _>>()
.map_err(Into::into)
} else {
stmt.query_map([], map_row)?
.collect::<Result<Vec<_>, _>>()
.map_err(Into::into)
}
}
#[cfg(test)]
pub(crate) fn migration_test_policy() -> MigrationWritePolicy {
MigrationWritePolicy::new(
crate::DiskGuardEnvironment::capture()
.resolve(None, None)
.expect("test disk policy"),
crate::PoolConfig::for_test()
.volume_lock_dir
.expect("test volume-lock directory"),
)
.expect("valid test migration policy")
}
#[cfg(test)]
pub(crate) fn run_migrations_for_test(conn: &mut Connection) -> Result<u32, SqliteError> {
run_migrations_with_policy(conn, &migration_test_policy())
}
#[cfg(test)]
fn apply_schema_plan_for_test(
conn: &mut Connection,
plan: &ServiceSchemaPlan,
) -> Result<(), SqliteError> {
apply_schema_plan_with_policy(conn, plan, &migration_test_policy())
}
#[cfg(test)]
fn stage_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
stage_attachment_cutover_with_policy(conn, &migration_test_policy())
}
#[cfg(test)]
fn finalize_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
finalize_attachment_cutover_with_policy(conn, &migration_test_policy())
}
#[cfg(test)]
#[path = "entity_version_migration_measurement.rs"]
mod entity_version_measurement;
#[cfg(test)]
#[path = "migrations_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "raw_migration_settlement_tests.rs"]
mod raw_settlement_tests;
#[cfg(test)]
#[path = "git_note_index_migration_tests.rs"]
mod git_note_indexes;
#[cfg(test)]
#[path = "entity_list_index_migration_tests.rs"]
mod entity_list_indexes;
#[cfg(test)]
#[path = "schedule_core_index_migration_tests.rs"]
mod schedule_core_index_migration_tests;
#[cfg(test)]
#[path = "comm_core_index_migration_tests.rs"]
mod comm_core_index_migration_tests;