use khive_storage::blob::ContentRef;
use rusqlite::{Connection, OptionalExtension};
use std::path::PathBuf;
use crate::error::SqliteError;
use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
#[path = "session_identity_migration.rs"]
mod session_identity_migration;
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: &Connection, plan: &ServiceSchemaPlan) -> Result<(), SqliteError> {
conn.execute_batch(SCHEMA_VERSION_TABLE)?;
for migration in plan.sqlite {
let tx =
rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
if let Some(check) = migration.is_already_applied {
if check(&tx) {
continue;
}
}
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 {
continue;
}
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(())
}
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 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 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,
},
];
#[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> {
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)?;
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> {
if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
return Ok(());
}
let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
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(),
));
}
Ok(current_version)
}
#[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);
}
std::fs::canonicalize(&raw_path)
.map(Some)
.map_err(SqliteError::Io)
}
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)?;
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(conn, &owner);
}
run_migrations_with_busy_timeout(conn)
}
pub(crate) fn run_migrations_with_database_gc_owner(
conn: &mut Connection,
owner: &DatabaseGcOwnerGuard,
) -> Result<u32, SqliteError> {
validate_database_gc_owner(conn, owner)?;
run_migrations_with_busy_timeout(conn)
}
fn run_migrations_with_busy_timeout(conn: &mut Connection) -> Result<u32, 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 = run_migrations_locked(conn);
if raised {
let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
}
result
}
fn run_migrations_locked(conn: &mut Connection) -> Result<u32, SqliteError> {
conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
let current_version: u32 = read_schema_version(conn)?;
#[cfg(test)]
if test_sync::PARTICIPATE.with(|p| p.get()) {
conn.busy_handler(Some(test_sync::record_busy))?;
let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
if let Some(barrier) = barrier {
barrier.wait();
}
}
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)?;
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;
}
#[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));
}
let tx = conn
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
.map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
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);
skip_through = sibling_version.min(latest_version);
applied_version = applied_version.max(migration.version);
continue;
}
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);
break;
}
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.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(),
})?;
}
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);
}
tx.commit().map_err(|e| SqliteError::Migration {
version: migration.version,
error: e.to_string(),
})?;
applied_version = migration.version;
}
validate_applied_migration_ledger(conn, applied_version)?;
Ok(applied_version)
}
#[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)]
#[path = "entity_version_migration_measurement.rs"]
mod entity_version_measurement;
#[cfg(test)]
#[path = "migrations_tests.rs"]
mod tests;