1use khive_storage::blob::ContentRef;
10use rusqlite::{Connection, OptionalExtension};
11use std::path::PathBuf;
12
13use crate::error::SqliteError;
14use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
15
16pub struct Migration {
22 pub id: &'static str,
24 pub up_sql: &'static str,
26 pub down_sql: Option<&'static str>,
28 pub is_already_applied: Option<fn(&Connection) -> bool>,
31}
32
33pub struct ServiceSchemaPlan {
35 pub service: &'static str,
37 pub sqlite: &'static [Migration],
39 pub postgres: &'static [Migration],
41}
42
43const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
44
45pub fn apply_schema_plan(conn: &Connection, plan: &ServiceSchemaPlan) -> Result<(), SqliteError> {
47 conn.execute_batch(SCHEMA_VERSION_TABLE)?;
48
49 for migration in plan.sqlite {
50 if let Some(check) = migration.is_already_applied {
52 if check(conn) {
53 continue;
54 }
55 }
56
57 let already: bool = conn.query_row(
59 "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
60 rusqlite::params![plan.service, migration.id],
61 |row| row.get(0),
62 )?;
63
64 if already {
65 continue;
66 }
67
68 let tx =
69 rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
70 tx.execute_batch(migration.up_sql)?;
71
72 tx.execute(
73 "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
74 rusqlite::params![
75 plan.service,
76 migration.id,
77 chrono::Utc::now().timestamp_micros(),
78 ],
79 )?;
80 tx.commit()?;
81 }
82
83 Ok(())
84}
85
86pub struct VersionedMigration {
96 pub version: u32,
98 pub name: &'static str,
100 pub up: &'static str,
103}
104
105const V1_UP: &str = include_str!("../sql/schema.sql");
108
109const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
110
111const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
112
113const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
114
115const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
116
117const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
118
119const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
120
121const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
122
123const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
124
125const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
126
127const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
128
129const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
130
131const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
132
133const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
134
135const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
136
137const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
138
139const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
140
141const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
142
143const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
144
145const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
146
147const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
148
149const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
150
151pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
153
154pub fn latest_schema_version() -> u32 {
159 MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
160}
161
162pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
170
171pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
176
177pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
184
185pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
191
192pub const MIGRATIONS: &[VersionedMigration] = &[
199 VersionedMigration {
200 version: 1,
201 name: "initial_schema",
202 up: V1_UP,
203 },
204 VersionedMigration {
205 version: 2,
206 name: "narrow_fts_sections_update_trigger",
207 up: V2_UP,
208 },
209 VersionedMigration {
210 version: 3,
211 name: "backfill_domain_mirror_atoms",
212 up: V3_UP,
213 },
214 VersionedMigration {
215 version: 4,
216 name: "fts_consolidation",
217 up: V4_UP,
218 },
219 VersionedMigration {
220 version: 5,
221 name: "unique_comm_message_external_id",
222 up: V5_UP,
223 },
224 VersionedMigration {
225 version: 6,
226 name: "brain_retune_driver",
227 up: V6_UP,
228 },
229 VersionedMigration {
230 version: 7,
231 name: "notes_seq",
232 up: V7_UP,
233 },
234 VersionedMigration {
235 version: 8,
236 name: "notes_seq_repair",
237 up: V8_UP,
238 },
239 VersionedMigration {
240 version: 9,
241 name: "entities_name_ci_index",
242 up: V9_UP,
243 },
244 VersionedMigration {
245 version: 10,
246 name: "entities_content_ref",
247 up: V10_UP,
248 },
249 VersionedMigration {
250 version: 11,
251 name: "ann_write_log",
252 up: V11_UP,
253 },
254 VersionedMigration {
255 version: 12,
256 name: "ann_write_log_model_seq_index",
257 up: V12_UP,
258 },
259 VersionedMigration {
260 version: 13,
261 name: "list_cursor_sequences",
262 up: V13_UP,
263 },
264 VersionedMigration {
265 version: 14,
266 name: "graph_edges_id_unique",
267 up: V14_UP,
268 },
269 VersionedMigration {
270 version: 15,
271 name: "serve_ledger_attribution",
272 up: V15_UP,
273 },
274 VersionedMigration {
275 version: 16,
276 name: "gtd_dependency_cycle_guards",
277 up: V16_UP,
278 },
279 VersionedMigration {
280 version: 17,
281 name: "agents_ddl",
282 up: V17_UP,
283 },
284 VersionedMigration {
285 version: 18,
286 name: "ann_consumer_pending",
287 up: V18_UP,
288 },
289 VersionedMigration {
290 version: 19,
291 name: "list_cursor_backfill_repair",
292 up: V19_UP,
293 },
294 VersionedMigration {
295 version: 20,
296 name: "blob_gc_claims",
297 up: V20_UP,
298 },
299 VersionedMigration {
300 version: ATTACHMENT_CUTOVER_VERSION,
301 name: "attachments_first_class",
302 up: V21_STAGE_UP,
306 },
307];
308
309#[derive(Clone, Copy, Debug, Eq, PartialEq)]
311pub enum AttachmentCutoverStatus {
312 Pending,
314 Incomplete,
316 Complete,
318}
319
320fn schema_object_exists(
321 conn: &Connection,
322 object_type: &str,
323 name: &str,
324) -> Result<bool, SqliteError> {
325 conn.query_row(
326 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
327 rusqlite::params![object_type, name],
328 |row| row.get(0),
329 )
330 .map_err(Into::into)
331}
332
333fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
334 conn.query_row(
335 "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
336 rusqlite::params![table, column],
337 |row| row.get(0),
338 )
339 .map_err(Into::into)
340}
341
342fn require_attachment_schema_objects(
343 conn: &Connection,
344 objects: &[(&str, &str)],
345 phase: &str,
346) -> Result<(), SqliteError> {
347 for (object_type, name) in objects {
348 if !schema_object_exists(conn, object_type, name)? {
349 return Err(SqliteError::InvalidData(format!(
350 "attachment cutover {phase} state is missing {object_type} {name:?}"
351 )));
352 }
353 }
354 Ok(())
355}
356
357fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
358 require_attachment_schema_objects(
359 conn,
360 &[
361 ("table", "attachments"),
362 ("index", "idx_attachments_content_ref"),
363 ],
364 "incomplete",
365 )?;
366 require_legacy_attachment_fences(conn)
367}
368
369fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
370 require_attachment_schema_objects(
371 conn,
372 &[
373 ("table", "attachments"),
374 ("table", "blob_gc_claims"),
375 ("index", "idx_attachments_content_ref"),
376 ("index", "idx_blob_gc_claims_content_ref"),
377 ("trigger", "attachments_reject_claimed_blob_insert"),
378 ("trigger", "attachments_reject_claimed_blob_update"),
379 ],
380 "complete",
381 )?;
382 if schema_column_exists(conn, "entities", "content_ref")? {
383 return Err(SqliteError::InvalidData(
384 "attachment cutover is complete but entities.content_ref still exists".into(),
385 ));
386 }
387 for (object_type, name) in [
388 ("index", "idx_entities_content_ref"),
389 ("trigger", "entities_reject_claimed_blob_insert"),
390 ("trigger", "entities_reject_claimed_blob_update"),
391 ] {
392 if schema_object_exists(conn, object_type, name)? {
393 return Err(SqliteError::InvalidData(format!(
394 "attachment cutover is complete but legacy {object_type} {name:?} still exists"
395 )));
396 }
397 }
398 Ok(())
399}
400
401pub fn attachment_cutover_status(
406 conn: &Connection,
407) -> Result<AttachmentCutoverStatus, SqliteError> {
408 let version = read_schema_version(conn)?;
409 let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
410 if !marker_table {
411 if version >= ATTACHMENT_CUTOVER_VERSION {
412 return Err(SqliteError::InvalidData(format!(
413 "migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
414 )));
415 }
416 if schema_object_exists(conn, "table", "attachments")? {
417 return Err(SqliteError::InvalidData(
418 "attachments table exists without the durable attachment cutover marker".into(),
419 ));
420 }
421 return Ok(AttachmentCutoverStatus::Pending);
422 }
423
424 let marker: Option<(String, Option<i64>)> = conn
425 .query_row(
426 "SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
427 [],
428 |row| Ok((row.get(0)?, row.get(1)?)),
429 )
430 .optional()?;
431 match marker {
432 Some((state, None)) if state == "incomplete" => {
433 if version >= ATTACHMENT_CUTOVER_VERSION {
434 Err(SqliteError::InvalidData(format!(
435 "attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
436 )))
437 } else {
438 validate_incomplete_attachment_schema(conn)?;
439 Ok(AttachmentCutoverStatus::Incomplete)
440 }
441 }
442 Some((state, Some(_))) if state == "complete" => {
443 if version >= ATTACHMENT_CUTOVER_VERSION {
447 validate_complete_attachment_schema(conn)?;
448 Ok(AttachmentCutoverStatus::Complete)
449 } else {
450 Err(SqliteError::InvalidData(format!(
451 "attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
452 )))
453 }
454 }
455 Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
456 "invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
457 ))),
458 None => Err(SqliteError::InvalidData(
459 "attachment cutover marker table exists without its singleton row".into(),
460 )),
461 }
462}
463
464fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
465 if !schema_column_exists(conn, "entities", "content_ref")? {
466 return Err(SqliteError::InvalidData(
467 "attachment cutover requires legacy entities.content_ref until finalization".into(),
468 ));
469 }
470 for (object_type, name) in [
471 ("table", "blob_gc_claims"),
472 ("index", "idx_blob_gc_claims_content_ref"),
473 ("index", "idx_entities_content_ref"),
474 ("trigger", "entities_reject_claimed_blob_insert"),
475 ("trigger", "entities_reject_claimed_blob_update"),
476 ] {
477 if !schema_object_exists(conn, object_type, name)? {
478 return Err(SqliteError::InvalidData(format!(
479 "attachment cutover requires legacy {object_type} {name:?} until finalization"
480 )));
481 }
482 }
483 Ok(())
484}
485
486fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
493 let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
494 if !(1..=4).contains(&width) {
495 return Err(SqliteError::InvalidData(format!(
496 "the text-encoding width probe returned {width}; refusing canonicality validation"
497 )));
498 }
499 Ok(width * 64)
500}
501
502fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
503 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
504 let invalid: Option<String> = conn
505 .query_row(
506 "SELECT id FROM entities \
507 WHERE content_ref IS NOT NULL \
508 AND (typeof(content_ref) <> 'text' \
509 OR length(content_ref) <> 64 \
510 OR length(CAST(content_ref AS BLOB)) <> ?1 \
511 OR content_ref GLOB '*[^0-9a-f]*') \
512 LIMIT 1",
513 [canonical_bytes],
514 |row| row.get(0),
515 )
516 .optional()?;
517 if let Some(id) = invalid {
518 return Err(SqliteError::InvalidData(format!(
519 "entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
520 )));
521 }
522 Ok(())
523}
524
525fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
526 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
527 for (table, identity) in [
528 ("attachments", "record_uuid"),
529 ("blob_gc_claims", "root_key"),
530 ] {
531 let sql = format!(
532 "SELECT {identity} FROM {table} \
533 WHERE typeof(content_ref) <> 'text' \
534 OR length(content_ref) <> 64 \
535 OR length(CAST(content_ref AS BLOB)) <> ?1 \
536 OR content_ref GLOB '*[^0-9a-f]*' \
537 LIMIT 1"
538 );
539 let invalid: Option<String> = conn
540 .query_row(&sql, [canonical_bytes], |row| row.get(0))
541 .optional()?;
542 if let Some(owner) = invalid {
543 return Err(SqliteError::InvalidData(format!(
544 "{table}.content_ref for {identity} {owner:?} is not canonical"
545 )));
546 }
547 }
548 Ok(())
549}
550
551fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
552 let dangling: Option<(String, String)> = conn
553 .query_row(
554 "SELECT record_uuid, substrate FROM attachments AS attachment \
555 WHERE (substrate = 'entity' AND NOT EXISTS ( \
556 SELECT 1 FROM entities WHERE id = attachment.record_uuid \
557 )) \
558 OR (substrate = 'note' AND NOT EXISTS ( \
559 SELECT 1 FROM notes WHERE id = attachment.record_uuid \
560 )) \
561 LIMIT 1",
562 [],
563 |row| Ok((row.get(0)?, row.get(1)?)),
564 )
565 .optional()?;
566 if let Some((record_uuid, substrate)) = dangling {
567 return Err(SqliteError::InvalidData(format!(
568 "attachment role references absent {substrate} record {record_uuid:?}"
569 )));
570 }
571 Ok(())
572}
573
574fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
575 let conflict: Option<String> = conn
576 .query_row(
577 "SELECT entity.id FROM entities AS entity \
578 LEFT JOIN attachments AS attachment \
579 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
580 WHERE entity.content_ref IS NOT NULL \
581 AND (attachment.record_uuid IS NULL \
582 OR attachment.substrate <> 'entity' \
583 OR attachment.content_ref <> entity.content_ref) \
584 LIMIT 1",
585 [],
586 |row| row.get(0),
587 )
588 .optional()?;
589 if let Some(record_uuid) = conflict {
590 return Err(SqliteError::InvalidData(format!(
591 "legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
592 )));
593 }
594 Ok(())
595}
596
597fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
598 require_legacy_attachment_fences(conn)?;
599 conn.execute_batch(V21_STAGE_UP)?;
600 validate_canonical_legacy_refs(conn)?;
601 validate_canonical_attachment_and_claim_refs(conn)?;
602
603 let conflict: Option<String> = conn
604 .query_row(
605 "SELECT entity.id FROM entities AS entity \
606 JOIN attachments AS attachment \
607 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
608 WHERE entity.content_ref IS NOT NULL \
609 AND (attachment.substrate <> 'entity' \
610 OR attachment.content_ref <> entity.content_ref) \
611 LIMIT 1",
612 [],
613 |row| row.get(0),
614 )
615 .optional()?;
616 if let Some(record_uuid) = conflict {
617 return Err(SqliteError::InvalidData(format!(
618 "existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
619 )));
620 }
621
622 conn.execute(
623 "INSERT INTO attachments \
624 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
625 SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
626 FROM entities WHERE content_ref IS NOT NULL \
627 ON CONFLICT(record_uuid, role) DO NOTHING",
628 [],
629 )?;
630 validate_legacy_content_backfill(conn)?;
631
632 conn.execute("DELETE FROM blob_gc_claims", [])?;
636 conn.execute(
637 "INSERT INTO attachment_cutover_state \
638 (singleton, state, started_at, completed_at) \
639 VALUES (1, 'incomplete', ?1, NULL) \
640 ON CONFLICT(singleton) DO NOTHING",
641 [now],
642 )?;
643 Ok(())
644}
645
646pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
653 match attachment_cutover_status(conn)? {
654 AttachmentCutoverStatus::Complete => return Ok(()),
655 AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
656 }
657 if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
658 return Err(SqliteError::InvalidData(format!(
659 "attachment cutover stage requires canonical V{} schema",
660 ATTACHMENT_CUTOVER_VERSION - 1
661 )));
662 }
663
664 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
665 let status = attachment_cutover_status(&tx)?;
666 if status == AttachmentCutoverStatus::Complete {
667 return Ok(());
668 }
669 stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
670 tx.commit()?;
671 Ok(())
672}
673
674#[allow(clippy::too_many_arguments)]
682pub fn apply_generic_verified_attachment(
683 conn: &Connection,
684 record_uuid: &str,
685 substrate: &str,
686 role: &str,
687 content_ref: &ContentRef,
688 media_type: Option<&str>,
689 size_bytes: Option<u64>,
690 created_at: i64,
691) -> Result<(), SqliteError> {
692 if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
693 return Err(SqliteError::InvalidData(
694 "verified application attachments may only be applied while V21 cutover is incomplete"
695 .into(),
696 ));
697 }
698 if role.is_empty() || role.chars().any(char::is_control) {
699 return Err(SqliteError::InvalidData(
700 "attachment role must be non-empty and contain no control characters".into(),
701 ));
702 }
703 let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
704 SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
705 })?;
706 let owner_table = match substrate {
707 "entity" => "entities",
708 "note" => "notes",
709 other => {
710 return Err(SqliteError::InvalidData(format!(
711 "attachment substrate must be 'entity' or 'note', got {other:?}"
712 )))
713 }
714 };
715 let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
716 let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
717 if !owner_exists {
718 return Err(SqliteError::InvalidData(format!(
719 "cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
720 )));
721 }
722 let claimed: bool = conn.query_row(
723 "SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
724 [content_ref.as_str()],
725 |row| row.get(0),
726 )?;
727 if claimed {
728 return Err(SqliteError::InvalidData(format!(
729 "cannot attach claimed content_ref {} during V21 cutover",
730 content_ref.as_str()
731 )));
732 }
733
734 let changed = conn.execute(
735 "INSERT INTO attachments \
736 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
737 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
738 ON CONFLICT(record_uuid, role) DO UPDATE SET \
739 media_type = excluded.media_type, \
740 size_bytes = excluded.size_bytes, \
741 created_at = excluded.created_at \
742 WHERE attachments.substrate = excluded.substrate \
743 AND attachments.content_ref = excluded.content_ref",
744 rusqlite::params![
745 record_uuid,
746 substrate,
747 role,
748 content_ref.as_str(),
749 media_type,
750 size_bytes,
751 created_at,
752 ],
753 )?;
754 if changed == 0 {
755 return Err(SqliteError::InvalidData(format!(
756 "attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
757 )));
758 }
759 Ok(())
760}
761
762fn finalize_attachment_cutover_on_connection(
763 conn: &Connection,
764 now: i64,
765) -> Result<(), SqliteError> {
766 require_legacy_attachment_fences(conn)?;
767 validate_canonical_legacy_refs(conn)?;
768 validate_canonical_attachment_and_claim_refs(conn)?;
769 validate_attachment_record_owners(conn)?;
770 validate_legacy_content_backfill(conn)?;
771
772 let remaining_claims: i64 =
773 conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
774 if remaining_claims != 0 {
775 return Err(SqliteError::InvalidData(format!(
776 "attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
777 )));
778 }
779
780 let uncovered_model: Option<String> = conn
781 .query_row(
782 "SELECT model.id FROM entities AS model \
783 WHERE model.entity_type = 'moodboard_model' \
784 AND model.content_ref IS NOT NULL \
785 AND NOT EXISTS ( \
786 SELECT 1 FROM attachments AS attachment \
787 WHERE attachment.record_uuid = model.id \
788 AND attachment.substrate = 'entity' \
789 AND attachment.role = 'fann-network' \
790 ) \
791 LIMIT 1",
792 [],
793 |row| row.get(0),
794 )
795 .optional()?;
796 if let Some(record_uuid) = uncovered_model {
797 return Err(SqliteError::InvalidData(format!(
798 "moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
799 )));
800 }
801
802 conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
803 conn.execute_batch(
804 "DROP TRIGGER entities_reject_claimed_blob_insert; \
805 DROP TRIGGER entities_reject_claimed_blob_update; \
806 DROP INDEX idx_entities_content_ref; \
807 ALTER TABLE entities DROP COLUMN content_ref;",
808 )?;
809 conn.execute(
810 "UPDATE attachment_cutover_state \
811 SET state = 'complete', completed_at = ?1 \
812 WHERE singleton = 1 AND state = 'incomplete'",
813 [now],
814 )?;
815 Ok(())
816}
817
818fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
819 let migration = MIGRATIONS
820 .iter()
821 .find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
822 .expect("V21 migration must be registered");
823 conn.execute(
824 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
825 rusqlite::params![migration.version, migration.name, now],
826 )?;
827 Ok(())
828}
829
830pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
837 if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
838 return Ok(());
839 }
840 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
841 match attachment_cutover_status(&tx)? {
842 AttachmentCutoverStatus::Complete => return Ok(()),
843 AttachmentCutoverStatus::Pending => {
844 return Err(SqliteError::InvalidData(
845 "attachment cutover must complete stage 1 before finalization".into(),
846 ))
847 }
848 AttachmentCutoverStatus::Incomplete => {}
849 }
850 let now = chrono::Utc::now().timestamp_micros();
851 finalize_attachment_cutover_on_connection(&tx, now)?;
852 record_attachment_cutover_migration(&tx, now)?;
853 tx.commit()?;
854 Ok(())
855}
856
857fn read_applied_migration_ledger(
859 conn: &Connection,
860 through_version: u32,
861) -> Result<Vec<(u32, String)>, SqliteError> {
862 let mut stmt = conn.prepare(
863 "SELECT version, name FROM _schema_migrations \
864 WHERE version <= ?1 ORDER BY version ASC",
865 )?;
866 let rows = stmt
867 .query_map([through_version], |row| {
868 Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
869 })?
870 .collect::<Result<Vec<_>, _>>()?;
871 Ok(rows)
872}
873
874fn validate_applied_migration_versions(
879 applied: &[(u32, String)],
880 through_version: u32,
881) -> Result<(), SqliteError> {
882 let expected: Vec<&VersionedMigration> = MIGRATIONS
883 .iter()
884 .filter(|migration| migration.version <= through_version)
885 .collect();
886 let mut applied_index = 0;
887
888 for migration in expected {
889 let Some((version, applied_name)) = applied.get(applied_index) else {
890 return Err(SqliteError::InvalidData(format!(
891 "migration history is missing version {} ('{}'); the applied ledger must be \
892 the exact contiguous canonical sequence through version {through_version}",
893 migration.version, migration.name,
894 )));
895 };
896 if *version < migration.version {
897 return Err(SqliteError::InvalidData(format!(
898 "migration history contains unknown version {version} recorded as \
899 '{applied_name}'; the applied ledger must contain only canonical versions"
900 )));
901 }
902 if *version > migration.version {
903 return Err(SqliteError::InvalidData(format!(
904 "migration history is missing version {} ('{}'); found version {version} \
905 next instead",
906 migration.version, migration.name,
907 )));
908 }
909 applied_index += 1;
910 }
911
912 if let Some((version, name)) = applied.get(applied_index) {
913 return Err(SqliteError::InvalidData(format!(
914 "migration history contains unknown version {version} recorded as '{name}'; \
915 the applied ledger must contain only canonical versions"
916 )));
917 }
918
919 Ok(())
920}
921
922fn validate_applied_migration_names(
923 applied: &[(u32, String)],
924 through_version: u32,
925 allow_known_v19_repairs: bool,
926) -> Result<(), SqliteError> {
927 for ((version, applied_name), migration) in applied.iter().zip(
928 MIGRATIONS
929 .iter()
930 .filter(|migration| migration.version <= through_version),
931 ) {
932 debug_assert_eq!(*version, migration.version);
933 if migration.name != applied_name.as_str() {
934 if allow_known_v19_repairs && matches!(*version, 13 | 14) {
935 continue;
936 }
937 return Err(SqliteError::InvalidData(format!(
938 "migration version {version} is recorded under name '{applied_name}', \
939 expected '{expected}'. This database's migration history does not match \
940 the current binary; recreate it from the current schema or repair the \
941 specific known divergence via a dedicated migration.",
942 expected = migration.name,
943 )));
944 }
945 }
946
947 Ok(())
948}
949
950fn validate_applied_migration_ledger(
954 conn: &Connection,
955 through_version: u32,
956) -> Result<(), SqliteError> {
957 let applied = read_applied_migration_ledger(conn, through_version)?;
958 validate_applied_migration_versions(&applied, through_version)?;
959 validate_applied_migration_names(&applied, through_version, false)
960}
961
962const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
963
964pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
970 match conn.query_row(
971 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
972 [],
973 |row| row.get(0),
974 ) {
975 Ok(version) => Ok(version),
976 Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
977 if msg.contains("no such table: _schema_migrations") =>
978 {
979 Ok(0)
980 }
981 Err(e) => Err(e.into()),
982 }
983}
984
985pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
990 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
991 read_schema_version(&conn)
992}
993
994pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
997 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
998 validate_schema_is_current(&conn)
999}
1000
1001pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
1009 let current_version = read_schema_version(conn)?;
1010 let latest_version = latest_schema_version();
1011
1012 if current_version < latest_version {
1013 return Err(SqliteError::InvalidData(format!(
1014 "read-only database schema version {current_version} is behind the latest known \
1015 migration {latest_version}; migrate a writable copy with this build before opening \
1016 the snapshot read-only"
1017 )));
1018 }
1019 if current_version > latest_version {
1020 return Err(SqliteError::InvalidData(format!(
1021 "read-only database schema version {current_version} is ahead of the latest known \
1022 migration {latest_version}; use a compatible newer build or recreate the snapshot"
1023 )));
1024 }
1025
1026 validate_applied_migration_ledger(conn, current_version)?;
1033 if current_version >= ATTACHMENT_CUTOVER_VERSION
1037 && attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
1038 {
1039 return Err(SqliteError::InvalidData(
1040 "read-only database has not completed the V21 attachment cutover".into(),
1041 ));
1042 }
1043
1044 Ok(current_version)
1045}
1046
1047#[cfg(test)]
1048pub(crate) mod test_sync {
1049 use std::sync::atomic::AtomicU32;
1050 use std::sync::{Arc, Barrier, Mutex};
1051
1052 pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
1056 pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
1058 pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
1062 std::sync::atomic::AtomicBool::new(false);
1063
1064 pub(crate) fn record_busy(_count: i32) -> bool {
1067 BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
1068 std::thread::sleep(std::time::Duration::from_millis(1));
1069 true
1070 }
1071
1072 pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
1075 std::sync::atomic::AtomicBool::new(false);
1076 pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
1081 std::sync::atomic::AtomicBool::new(false);
1082
1083 std::thread_local! {
1084 pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
1087 const { std::cell::Cell::new(false) };
1088 pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
1090 const { std::cell::Cell::new(false) };
1091 }
1092}
1093
1094fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
1103 let configured = conn.path().unwrap_or_default();
1104 let raw_path = if configured.is_empty() {
1105 conn.query_row(
1106 "SELECT file FROM pragma_database_list WHERE name = 'main'",
1107 [],
1108 |row| row.get::<_, String>(0),
1109 )?
1110 } else {
1111 configured.to_string()
1112 };
1113
1114 if raw_path.is_empty() {
1115 return Ok(None);
1116 }
1117 std::fs::canonicalize(&raw_path)
1118 .map(Some)
1119 .map_err(SqliteError::Io)
1120}
1121
1122fn validate_database_gc_owner(
1123 conn: &Connection,
1124 owner: &DatabaseGcOwnerGuard,
1125) -> Result<(), SqliteError> {
1126 let connection_path = canonical_connection_database_path(conn)?;
1127 if owner.database_path() != connection_path.as_deref() {
1128 return Err(SqliteError::InvalidData(format!(
1129 "database GC owner targets {:?}, but migration connection targets {:?}",
1130 owner.database_path(),
1131 connection_path.as_deref(),
1132 )));
1133 }
1134 Ok(())
1135}
1136
1137pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
1138 let database_path = canonical_connection_database_path(conn)?;
1139 if let Some(database_path) = database_path {
1140 let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
1145 SqliteError::InvalidData(format!(
1146 "failed to acquire database GC owner before schema migration: {error}"
1147 ))
1148 })?;
1149 return run_migrations_with_database_gc_owner(conn, &owner);
1150 }
1151
1152 run_migrations_with_busy_timeout(conn)
1155}
1156
1157pub(crate) fn run_migrations_with_database_gc_owner(
1158 conn: &mut Connection,
1159 owner: &DatabaseGcOwnerGuard,
1160) -> Result<u32, SqliteError> {
1161 validate_database_gc_owner(conn, owner)?;
1162 run_migrations_with_busy_timeout(conn)
1163}
1164
1165fn run_migrations_with_busy_timeout(conn: &mut Connection) -> Result<u32, SqliteError> {
1166 let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
1171 let raised = prior_busy_ms < 5_000;
1172 if raised {
1173 conn.busy_timeout(std::time::Duration::from_secs(5))?;
1174 }
1175 let result = run_migrations_locked(conn);
1176 if raised {
1177 let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
1178 }
1179 result
1180}
1181
1182fn run_migrations_locked(conn: &mut Connection) -> Result<u32, SqliteError> {
1183 conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
1184
1185 let current_version: u32 = read_schema_version(conn)?;
1186
1187 #[cfg(test)]
1191 if test_sync::PARTICIPATE.with(|p| p.get()) {
1192 conn.busy_handler(Some(test_sync::record_busy))?;
1195 let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
1196 if let Some(barrier) = barrier {
1197 barrier.wait();
1198 }
1199 }
1200
1201 let latest_version = latest_schema_version();
1207 if current_version > latest_version {
1208 return Err(SqliteError::InvalidData(format!(
1209 "database schema version {current_version} is ahead of the latest known migration \
1210 {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
1211 written by a newer build. Recreate it from the current schema; in-place downgrade is \
1212 not supported."
1213 )));
1214 }
1215
1216 let applied = read_applied_migration_ledger(conn, current_version)?;
1223 validate_applied_migration_versions(&applied, current_version)?;
1224 validate_applied_migration_names(&applied, current_version, current_version < 19)?;
1225
1226 let mut applied_version = current_version;
1227 let mut skip_through = current_version;
1231
1232 for migration in MIGRATIONS {
1233 if migration.version <= skip_through {
1234 applied_version = applied_version.max(migration.version);
1235 continue;
1236 }
1237
1238 #[cfg(test)]
1242 let instrumented_first_begin = test_sync::PARTICIPATE.with(|p| p.get())
1243 && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
1244 #[cfg(test)]
1245 if instrumented_first_begin {
1246 test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
1247 }
1248 let tx = conn
1249 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1250 .map_err(|e| SqliteError::Migration {
1251 version: migration.version,
1252 error: e.to_string(),
1253 })?;
1254
1255 let sibling_version: u32 = tx
1259 .query_row(
1260 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1261 [],
1262 |row| row.get(0),
1263 )
1264 .map_err(|e| SqliteError::Migration {
1265 version: migration.version,
1266 error: e.to_string(),
1267 })?;
1268 #[cfg(test)]
1269 if instrumented_first_begin {
1270 use std::sync::atomic::Ordering::SeqCst;
1271 if sibling_version == 0 {
1272 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1278 while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline
1279 {
1280 std::thread::yield_now();
1281 }
1282 } else {
1283 test_sync::LOSER_SAW_WINNER_COMMIT
1287 .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
1288 }
1289 }
1290
1291 if sibling_version > latest_version {
1296 return Err(SqliteError::InvalidData(format!(
1297 "database schema version {sibling_version} is ahead of the latest known \
1298 migration {latest_version} (committed by a concurrent process while this \
1299 one waited for the migration write lock). This build cannot run against \
1300 the newer schema; upgrade the binary or recreate the database."
1301 )));
1302 }
1303
1304 if sibling_version >= migration.version {
1305 #[cfg(test)]
1306 test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1307 skip_through = sibling_version.min(latest_version);
1308 applied_version = applied_version.max(migration.version);
1309 continue;
1310 }
1311
1312 if migration.version == ATTACHMENT_CUTOVER_VERSION {
1313 let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
1314 version: migration.version,
1315 error: e.to_string(),
1316 })?;
1317 let legacy_refs: i64 = tx
1318 .query_row(
1319 "SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
1320 [],
1321 |row| row.get(0),
1322 )
1323 .map_err(|e| SqliteError::Migration {
1324 version: migration.version,
1325 error: e.to_string(),
1326 })?;
1327
1328 if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
1332 drop(tx);
1333 break;
1334 }
1335 if status != AttachmentCutoverStatus::Pending {
1336 return Err(SqliteError::Migration {
1337 version: migration.version,
1338 error: format!("unexpected attachment cutover state {status:?}"),
1339 });
1340 }
1341
1342 let now = chrono::Utc::now().timestamp_micros();
1343 stage_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1344 SqliteError::Migration {
1345 version: migration.version,
1346 error: e.to_string(),
1347 }
1348 })?;
1349 finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1350 SqliteError::Migration {
1351 version: migration.version,
1352 error: e.to_string(),
1353 }
1354 })?;
1355 } else {
1356 tx.execute_batch(migration.up)
1357 .map_err(|e| SqliteError::Migration {
1358 version: migration.version,
1359 error: e.to_string(),
1360 })?;
1361 }
1362
1363 if migration.version == 19 {
1370 tx.execute_batch(
1371 "UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
1372 UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
1373 )
1374 .map_err(|e| SqliteError::Migration {
1375 version: migration.version,
1376 error: e.to_string(),
1377 })?;
1378 }
1379
1380 let now = chrono::Utc::now().timestamp_micros();
1381 tx.execute(
1382 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
1383 ON CONFLICT(version) DO NOTHING",
1384 rusqlite::params![migration.version, migration.name, now],
1385 )
1386 .map_err(|e| SqliteError::Migration {
1387 version: migration.version,
1388 error: e.to_string(),
1389 })?;
1390
1391 #[cfg(test)]
1392 if instrumented_first_begin {
1393 test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
1394 }
1395
1396 tx.commit().map_err(|e| SqliteError::Migration {
1397 version: migration.version,
1398 error: e.to_string(),
1399 })?;
1400
1401 applied_version = migration.version;
1402 }
1403
1404 validate_applied_migration_ledger(conn, applied_version)?;
1408
1409 Ok(applied_version)
1410}
1411
1412#[derive(Debug)]
1413pub struct EmbeddingModelRegistryRecord {
1414 pub engine_name: String,
1416 pub model_id: String,
1418 pub key_version: String,
1420 pub dimensions: u32,
1422 pub status: String,
1424 pub activated_at: Option<i64>,
1426 pub superseded_at: Option<i64>,
1428}
1429
1430pub fn query_embedding_models(
1436 db: Option<&std::path::Path>,
1437 engine_filter: Option<&str>,
1438) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1439 let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
1440 std::env::var("HOME")
1441 .map(std::path::PathBuf::from)
1442 .unwrap_or_else(|_| std::path::PathBuf::from("."))
1443 .join(".khive/khive.db")
1444 });
1445 if !path.exists() {
1446 return Ok(Vec::new());
1447 }
1448 let conn = Connection::open_with_flags(
1449 path,
1450 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
1451 | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
1452 | rusqlite::OpenFlags::SQLITE_OPEN_URI,
1453 )?;
1454 query_embedding_models_conn(&conn, engine_filter)
1455}
1456
1457pub(crate) fn query_embedding_models_conn(
1461 conn: &Connection,
1462 engine_filter: Option<&str>,
1463) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1464 let exists: bool = conn.query_row(
1465 "SELECT COUNT(*) > 0 FROM sqlite_master \
1466 WHERE type='table' AND name='_embedding_models'",
1467 [],
1468 |row| row.get(0),
1469 )?;
1470 if !exists {
1471 return Ok(Vec::new());
1472 }
1473
1474 let sql = if engine_filter.is_some() {
1475 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1476 FROM _embedding_models WHERE engine_name = ?1 \
1477 ORDER BY engine_name, activated_at IS NULL, activated_at"
1478 } else {
1479 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1480 FROM _embedding_models \
1481 ORDER BY engine_name, activated_at IS NULL, activated_at"
1482 };
1483 let mut stmt = conn.prepare(sql)?;
1484 let map_row = |row: &rusqlite::Row<'_>| {
1485 let dim_raw: i64 = row.get(3)?;
1486 let dimensions = u32::try_from(dim_raw).map_err(|_| {
1487 rusqlite::Error::FromSqlConversionFailure(
1488 3,
1489 rusqlite::types::Type::Integer,
1490 Box::new(std::io::Error::other(format!(
1491 "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
1492 u32::MAX,
1493 ))),
1494 )
1495 })?;
1496 Ok(EmbeddingModelRegistryRecord {
1497 engine_name: row.get(0)?,
1498 model_id: row.get(1)?,
1499 key_version: row.get(2)?,
1500 dimensions,
1501 status: row.get(4)?,
1502 activated_at: row.get(5)?,
1503 superseded_at: row.get(6)?,
1504 })
1505 };
1506
1507 if let Some(engine) = engine_filter {
1508 stmt.query_map([engine], map_row)?
1509 .collect::<Result<Vec<_>, _>>()
1510 .map_err(Into::into)
1511 } else {
1512 stmt.query_map([], map_row)?
1513 .collect::<Result<Vec<_>, _>>()
1514 .map_err(Into::into)
1515 }
1516}
1517
1518#[cfg(test)]
1523#[path = "migrations_tests.rs"]
1524mod tests;