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
16#[path = "session_identity_migration.rs"]
17mod session_identity_migration;
18
19pub struct Migration {
25 pub id: &'static str,
27 pub up_sql: &'static str,
29 pub down_sql: Option<&'static str>,
31 pub is_already_applied: Option<fn(&Connection) -> bool>,
35}
36
37pub struct ServiceSchemaPlan {
39 pub service: &'static str,
41 pub sqlite: &'static [Migration],
43 pub postgres: &'static [Migration],
45}
46
47const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
48
49pub fn apply_schema_plan(conn: &Connection, plan: &ServiceSchemaPlan) -> Result<(), SqliteError> {
51 conn.execute_batch(SCHEMA_VERSION_TABLE)?;
52
53 for migration in plan.sqlite {
54 let tx =
58 rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
59
60 if let Some(check) = migration.is_already_applied {
62 if check(&tx) {
63 continue;
64 }
65 }
66
67 let already: bool = tx.query_row(
69 "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
70 rusqlite::params![plan.service, migration.id],
71 |row| row.get(0),
72 )?;
73
74 if already {
75 continue;
76 }
77
78 tx.execute_batch(migration.up_sql)?;
79
80 tx.execute(
81 "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
82 rusqlite::params![
83 plan.service,
84 migration.id,
85 chrono::Utc::now().timestamp_micros(),
86 ],
87 )?;
88 tx.commit()?;
89 }
90
91 Ok(())
92}
93
94pub struct VersionedMigration {
104 pub version: u32,
106 pub name: &'static str,
108 pub up: &'static str,
111}
112
113const V1_UP: &str = include_str!("../sql/schema.sql");
116
117const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
118
119const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
120
121const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
122
123const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
124
125const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
126
127const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
128
129const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
130
131const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
132
133const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
134
135const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
136
137const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
138
139const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
140
141const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
142
143const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
144
145const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
146
147const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
148
149const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
150
151const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
152
153const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
154
155const V22_UP: &str = include_str!("../sql/022-notes-unread-probe-recipient.sql");
156
157const V23_UP: &str = include_str!("../sql/023-fts-record-kind.sql");
158
159const V24_UP: &str = include_str!("../sql/024-fts-rowid-map.sql");
160
161const V25_UP: &str = include_str!("../sql/025-notes-unread-probe-recipient-direction.sql");
162
163const V26_UP: &str = include_str!("../sql/026-knowledge-fts-repair.sql");
164
165const V27_UP: &str = include_str!("../sql/027-notes-hot-property-indexes.sql");
166const V28_UP: &str = include_str!("../sql/028-notes-key.sql");
167
168const V29_UP: &str = include_str!("../sql/029-note-streams.sql");
169const V30_UP: &str = include_str!("../sql/030-tool-source-mounts.sql");
170const V31_UP: &str = include_str!("../sql/031-note-versions.sql");
171const V32_UP: &str = include_str!("../sql/032-knowledge-count-indexes.sql");
172const V33_UP: &str = include_str!("../sql/033-notes-message-recipient-direction.sql");
173const V34_UP: &str = include_str!("../sql/034-notes-namespace-created.sql");
174const V35_UP: &str = include_str!("../sql/035-notes-unread-probe-recipient-type-direction.sql");
175const V36_UP: &str = include_str!("../sql/036-events-operation-attribution.sql");
176const V37_UP: &str = include_str!("../sql/037-entity-versions.sql");
177const V38_UP: &str = include_str!("../sql/038-entities-legacy-type-index.sql");
178const V39_UP: &str = include_str!("../sql/039-knowledge-cursor-indexes.sql");
179const SESSION_IDENTITY_UP: &str = include_str!("../sql/040-session-source-scope.sql");
180const SESSION_IDENTITY_MIGRATION_NAME: &str = "session_source_scoped_identity";
181const V41_UP: &str = include_str!("../sql/041-sender-transport.sql");
182
183const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
184
185const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
186
187pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
189
190pub fn latest_schema_version() -> u32 {
195 MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
196}
197
198pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
206
207pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
212
213pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
220
221pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
227
228pub const MIGRATIONS: &[VersionedMigration] = &[
235 VersionedMigration {
236 version: 1,
237 name: "initial_schema",
238 up: V1_UP,
239 },
240 VersionedMigration {
241 version: 2,
242 name: "narrow_fts_sections_update_trigger",
243 up: V2_UP,
244 },
245 VersionedMigration {
246 version: 3,
247 name: "backfill_domain_mirror_atoms",
248 up: V3_UP,
249 },
250 VersionedMigration {
251 version: 4,
252 name: "fts_consolidation",
253 up: V4_UP,
254 },
255 VersionedMigration {
256 version: 5,
257 name: "unique_comm_message_external_id",
258 up: V5_UP,
259 },
260 VersionedMigration {
261 version: 6,
262 name: "brain_retune_driver",
263 up: V6_UP,
264 },
265 VersionedMigration {
266 version: 7,
267 name: "notes_seq",
268 up: V7_UP,
269 },
270 VersionedMigration {
271 version: 8,
272 name: "notes_seq_repair",
273 up: V8_UP,
274 },
275 VersionedMigration {
276 version: 9,
277 name: "entities_name_ci_index",
278 up: V9_UP,
279 },
280 VersionedMigration {
281 version: 10,
282 name: "entities_content_ref",
283 up: V10_UP,
284 },
285 VersionedMigration {
286 version: 11,
287 name: "ann_write_log",
288 up: V11_UP,
289 },
290 VersionedMigration {
291 version: 12,
292 name: "ann_write_log_model_seq_index",
293 up: V12_UP,
294 },
295 VersionedMigration {
296 version: 13,
297 name: "list_cursor_sequences",
298 up: V13_UP,
299 },
300 VersionedMigration {
301 version: 14,
302 name: "graph_edges_id_unique",
303 up: V14_UP,
304 },
305 VersionedMigration {
306 version: 15,
307 name: "serve_ledger_attribution",
308 up: V15_UP,
309 },
310 VersionedMigration {
311 version: 16,
312 name: "gtd_dependency_cycle_guards",
313 up: V16_UP,
314 },
315 VersionedMigration {
316 version: 17,
317 name: "agents_ddl",
318 up: V17_UP,
319 },
320 VersionedMigration {
321 version: 18,
322 name: "ann_consumer_pending",
323 up: V18_UP,
324 },
325 VersionedMigration {
326 version: 19,
327 name: "list_cursor_backfill_repair",
328 up: V19_UP,
329 },
330 VersionedMigration {
331 version: 20,
332 name: "blob_gc_claims",
333 up: V20_UP,
334 },
335 VersionedMigration {
336 version: ATTACHMENT_CUTOVER_VERSION,
337 name: "attachments_first_class",
338 up: V21_STAGE_UP,
342 },
343 VersionedMigration {
344 version: 22,
345 name: "notes_unread_probe_recipient",
346 up: V22_UP,
347 },
348 VersionedMigration {
349 version: 23,
350 name: "fts_record_kind",
351 up: V23_UP,
352 },
353 VersionedMigration {
354 version: 24,
355 name: "fts_rowid_map",
356 up: V24_UP,
357 },
358 VersionedMigration {
359 version: 25,
360 name: "notes_unread_probe_recipient_direction",
361 up: V25_UP,
362 },
363 VersionedMigration {
364 version: 26,
365 name: "knowledge_fts_repair",
366 up: V26_UP,
367 },
368 VersionedMigration {
369 version: 27,
370 name: "notes_hot_property_indexes",
371 up: V27_UP,
372 },
373 VersionedMigration {
374 version: 28,
375 name: "notes_key",
376 up: V28_UP,
377 },
378 VersionedMigration {
379 version: 29,
380 name: "note_streams",
381 up: V29_UP,
382 },
383 VersionedMigration {
384 version: 30,
385 name: "tool_source_mounts",
386 up: V30_UP,
387 },
388 VersionedMigration {
389 version: 31,
390 name: "note_versions",
391 up: V31_UP,
392 },
393 VersionedMigration {
394 version: 32,
395 name: "knowledge_count_indexes",
396 up: V32_UP,
397 },
398 VersionedMigration {
399 version: 33,
400 name: "notes_message_recipient_direction",
401 up: V33_UP,
402 },
403 VersionedMigration {
404 version: 34,
405 name: "notes_namespace_created",
406 up: V34_UP,
407 },
408 VersionedMigration {
409 version: 35,
410 name: "notes_unread_probe_recipient_type_direction",
411 up: V35_UP,
412 },
413 VersionedMigration {
414 version: 36,
415 name: "events_operation_attribution",
416 up: V36_UP,
417 },
418 VersionedMigration {
419 version: 37,
420 name: "entity_versions",
421 up: V37_UP,
422 },
423 VersionedMigration {
424 version: 38,
425 name: "entities_legacy_type_index",
426 up: V38_UP,
427 },
428 VersionedMigration {
429 version: 39,
430 name: "knowledge_cursor_indexes",
431 up: V39_UP,
432 },
433 VersionedMigration {
434 version: 40,
435 name: SESSION_IDENTITY_MIGRATION_NAME,
436 up: SESSION_IDENTITY_UP,
437 },
438 VersionedMigration {
439 version: 41,
440 name: "sender_transport",
441 up: V41_UP,
442 },
443];
444
445#[derive(Clone, Copy, Debug, Eq, PartialEq)]
447pub enum AttachmentCutoverStatus {
448 Pending,
450 Incomplete,
452 Complete,
454}
455
456fn schema_object_exists(
457 conn: &Connection,
458 object_type: &str,
459 name: &str,
460) -> Result<bool, SqliteError> {
461 conn.query_row(
462 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
463 rusqlite::params![object_type, name],
464 |row| row.get(0),
465 )
466 .map_err(Into::into)
467}
468
469fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
470 conn.query_row(
471 "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
472 rusqlite::params![table, column],
473 |row| row.get(0),
474 )
475 .map_err(Into::into)
476}
477
478fn require_attachment_schema_objects(
479 conn: &Connection,
480 objects: &[(&str, &str)],
481 phase: &str,
482) -> Result<(), SqliteError> {
483 for (object_type, name) in objects {
484 if !schema_object_exists(conn, object_type, name)? {
485 return Err(SqliteError::InvalidData(format!(
486 "attachment cutover {phase} state is missing {object_type} {name:?}"
487 )));
488 }
489 }
490 Ok(())
491}
492
493fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
494 require_attachment_schema_objects(
495 conn,
496 &[
497 ("table", "attachments"),
498 ("index", "idx_attachments_content_ref"),
499 ],
500 "incomplete",
501 )?;
502 require_legacy_attachment_fences(conn)
503}
504
505fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
506 require_attachment_schema_objects(
507 conn,
508 &[
509 ("table", "attachments"),
510 ("table", "blob_gc_claims"),
511 ("index", "idx_attachments_content_ref"),
512 ("index", "idx_blob_gc_claims_content_ref"),
513 ("trigger", "attachments_reject_claimed_blob_insert"),
514 ("trigger", "attachments_reject_claimed_blob_update"),
515 ],
516 "complete",
517 )?;
518 if schema_column_exists(conn, "entities", "content_ref")? {
519 return Err(SqliteError::InvalidData(
520 "attachment cutover is complete but entities.content_ref still exists".into(),
521 ));
522 }
523 for (object_type, name) in [
524 ("index", "idx_entities_content_ref"),
525 ("trigger", "entities_reject_claimed_blob_insert"),
526 ("trigger", "entities_reject_claimed_blob_update"),
527 ] {
528 if schema_object_exists(conn, object_type, name)? {
529 return Err(SqliteError::InvalidData(format!(
530 "attachment cutover is complete but legacy {object_type} {name:?} still exists"
531 )));
532 }
533 }
534 Ok(())
535}
536
537pub fn attachment_cutover_status(
542 conn: &Connection,
543) -> Result<AttachmentCutoverStatus, SqliteError> {
544 let version = read_schema_version(conn)?;
545 let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
546 if !marker_table {
547 if version >= ATTACHMENT_CUTOVER_VERSION {
548 return Err(SqliteError::InvalidData(format!(
549 "migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
550 )));
551 }
552 if schema_object_exists(conn, "table", "attachments")? {
553 return Err(SqliteError::InvalidData(
554 "attachments table exists without the durable attachment cutover marker".into(),
555 ));
556 }
557 return Ok(AttachmentCutoverStatus::Pending);
558 }
559
560 let marker: Option<(String, Option<i64>)> = conn
561 .query_row(
562 "SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
563 [],
564 |row| Ok((row.get(0)?, row.get(1)?)),
565 )
566 .optional()?;
567 match marker {
568 Some((state, None)) if state == "incomplete" => {
569 if version >= ATTACHMENT_CUTOVER_VERSION {
570 Err(SqliteError::InvalidData(format!(
571 "attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
572 )))
573 } else {
574 validate_incomplete_attachment_schema(conn)?;
575 Ok(AttachmentCutoverStatus::Incomplete)
576 }
577 }
578 Some((state, Some(_))) if state == "complete" => {
579 if version >= ATTACHMENT_CUTOVER_VERSION {
583 validate_complete_attachment_schema(conn)?;
584 Ok(AttachmentCutoverStatus::Complete)
585 } else {
586 Err(SqliteError::InvalidData(format!(
587 "attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
588 )))
589 }
590 }
591 Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
592 "invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
593 ))),
594 None => Err(SqliteError::InvalidData(
595 "attachment cutover marker table exists without its singleton row".into(),
596 )),
597 }
598}
599
600fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
601 if !schema_column_exists(conn, "entities", "content_ref")? {
602 return Err(SqliteError::InvalidData(
603 "attachment cutover requires legacy entities.content_ref until finalization".into(),
604 ));
605 }
606 for (object_type, name) in [
607 ("table", "blob_gc_claims"),
608 ("index", "idx_blob_gc_claims_content_ref"),
609 ("index", "idx_entities_content_ref"),
610 ("trigger", "entities_reject_claimed_blob_insert"),
611 ("trigger", "entities_reject_claimed_blob_update"),
612 ] {
613 if !schema_object_exists(conn, object_type, name)? {
614 return Err(SqliteError::InvalidData(format!(
615 "attachment cutover requires legacy {object_type} {name:?} until finalization"
616 )));
617 }
618 }
619 Ok(())
620}
621
622fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
629 let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
630 if !(1..=4).contains(&width) {
631 return Err(SqliteError::InvalidData(format!(
632 "the text-encoding width probe returned {width}; refusing canonicality validation"
633 )));
634 }
635 Ok(width * 64)
636}
637
638fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
639 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
640 let invalid: Option<String> = conn
641 .query_row(
642 "SELECT id FROM entities \
643 WHERE content_ref IS NOT NULL \
644 AND (typeof(content_ref) <> 'text' \
645 OR length(content_ref) <> 64 \
646 OR length(CAST(content_ref AS BLOB)) <> ?1 \
647 OR content_ref GLOB '*[^0-9a-f]*') \
648 LIMIT 1",
649 [canonical_bytes],
650 |row| row.get(0),
651 )
652 .optional()?;
653 if let Some(id) = invalid {
654 return Err(SqliteError::InvalidData(format!(
655 "entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
656 )));
657 }
658 Ok(())
659}
660
661fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
662 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
663 for (table, identity) in [
664 ("attachments", "record_uuid"),
665 ("blob_gc_claims", "root_key"),
666 ] {
667 let sql = format!(
668 "SELECT {identity} FROM {table} \
669 WHERE typeof(content_ref) <> 'text' \
670 OR length(content_ref) <> 64 \
671 OR length(CAST(content_ref AS BLOB)) <> ?1 \
672 OR content_ref GLOB '*[^0-9a-f]*' \
673 LIMIT 1"
674 );
675 let invalid: Option<String> = conn
676 .query_row(&sql, [canonical_bytes], |row| row.get(0))
677 .optional()?;
678 if let Some(owner) = invalid {
679 return Err(SqliteError::InvalidData(format!(
680 "{table}.content_ref for {identity} {owner:?} is not canonical"
681 )));
682 }
683 }
684 Ok(())
685}
686
687fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
688 let dangling: Option<(String, String)> = conn
689 .query_row(
690 "SELECT record_uuid, substrate FROM attachments AS attachment \
691 WHERE (substrate = 'entity' AND NOT EXISTS ( \
692 SELECT 1 FROM entities WHERE id = attachment.record_uuid \
693 )) \
694 OR (substrate = 'note' AND NOT EXISTS ( \
695 SELECT 1 FROM notes WHERE id = attachment.record_uuid \
696 )) \
697 LIMIT 1",
698 [],
699 |row| Ok((row.get(0)?, row.get(1)?)),
700 )
701 .optional()?;
702 if let Some((record_uuid, substrate)) = dangling {
703 return Err(SqliteError::InvalidData(format!(
704 "attachment role references absent {substrate} record {record_uuid:?}"
705 )));
706 }
707 Ok(())
708}
709
710fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
711 let conflict: Option<String> = conn
712 .query_row(
713 "SELECT entity.id FROM entities AS entity \
714 LEFT JOIN attachments AS attachment \
715 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
716 WHERE entity.content_ref IS NOT NULL \
717 AND (attachment.record_uuid IS NULL \
718 OR attachment.substrate <> 'entity' \
719 OR attachment.content_ref <> entity.content_ref) \
720 LIMIT 1",
721 [],
722 |row| row.get(0),
723 )
724 .optional()?;
725 if let Some(record_uuid) = conflict {
726 return Err(SqliteError::InvalidData(format!(
727 "legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
728 )));
729 }
730 Ok(())
731}
732
733fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
734 require_legacy_attachment_fences(conn)?;
735 conn.execute_batch(V21_STAGE_UP)?;
736 validate_canonical_legacy_refs(conn)?;
737 validate_canonical_attachment_and_claim_refs(conn)?;
738
739 let conflict: Option<String> = conn
740 .query_row(
741 "SELECT entity.id FROM entities AS entity \
742 JOIN attachments AS attachment \
743 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
744 WHERE entity.content_ref IS NOT NULL \
745 AND (attachment.substrate <> 'entity' \
746 OR attachment.content_ref <> entity.content_ref) \
747 LIMIT 1",
748 [],
749 |row| row.get(0),
750 )
751 .optional()?;
752 if let Some(record_uuid) = conflict {
753 return Err(SqliteError::InvalidData(format!(
754 "existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
755 )));
756 }
757
758 conn.execute(
759 "INSERT INTO attachments \
760 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
761 SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
762 FROM entities WHERE content_ref IS NOT NULL \
763 ON CONFLICT(record_uuid, role) DO NOTHING",
764 [],
765 )?;
766 validate_legacy_content_backfill(conn)?;
767
768 conn.execute("DELETE FROM blob_gc_claims", [])?;
772 conn.execute(
773 "INSERT INTO attachment_cutover_state \
774 (singleton, state, started_at, completed_at) \
775 VALUES (1, 'incomplete', ?1, NULL) \
776 ON CONFLICT(singleton) DO NOTHING",
777 [now],
778 )?;
779 Ok(())
780}
781
782pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
789 match attachment_cutover_status(conn)? {
790 AttachmentCutoverStatus::Complete => return Ok(()),
791 AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
792 }
793 if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
794 return Err(SqliteError::InvalidData(format!(
795 "attachment cutover stage requires canonical V{} schema",
796 ATTACHMENT_CUTOVER_VERSION - 1
797 )));
798 }
799
800 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
801 let status = attachment_cutover_status(&tx)?;
802 if status == AttachmentCutoverStatus::Complete {
803 return Ok(());
804 }
805 stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
806 tx.commit()?;
807 Ok(())
808}
809
810#[allow(clippy::too_many_arguments)]
818pub fn apply_generic_verified_attachment(
819 conn: &Connection,
820 record_uuid: &str,
821 substrate: &str,
822 role: &str,
823 content_ref: &ContentRef,
824 media_type: Option<&str>,
825 size_bytes: Option<u64>,
826 created_at: i64,
827) -> Result<(), SqliteError> {
828 if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
829 return Err(SqliteError::InvalidData(
830 "verified application attachments may only be applied while V21 cutover is incomplete"
831 .into(),
832 ));
833 }
834 if role.is_empty() || role.chars().any(char::is_control) {
835 return Err(SqliteError::InvalidData(
836 "attachment role must be non-empty and contain no control characters".into(),
837 ));
838 }
839 let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
840 SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
841 })?;
842 let owner_table = match substrate {
843 "entity" => "entities",
844 "note" => "notes",
845 other => {
846 return Err(SqliteError::InvalidData(format!(
847 "attachment substrate must be 'entity' or 'note', got {other:?}"
848 )))
849 }
850 };
851 let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
852 let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
853 if !owner_exists {
854 return Err(SqliteError::InvalidData(format!(
855 "cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
856 )));
857 }
858 let claimed: bool = conn.query_row(
859 "SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
860 [content_ref.as_str()],
861 |row| row.get(0),
862 )?;
863 if claimed {
864 return Err(SqliteError::InvalidData(format!(
865 "cannot attach claimed content_ref {} during V21 cutover",
866 content_ref.as_str()
867 )));
868 }
869
870 let changed = conn.execute(
871 "INSERT INTO attachments \
872 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
873 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
874 ON CONFLICT(record_uuid, role) DO UPDATE SET \
875 media_type = excluded.media_type, \
876 size_bytes = excluded.size_bytes, \
877 created_at = excluded.created_at \
878 WHERE attachments.substrate = excluded.substrate \
879 AND attachments.content_ref = excluded.content_ref",
880 rusqlite::params![
881 record_uuid,
882 substrate,
883 role,
884 content_ref.as_str(),
885 media_type,
886 size_bytes,
887 created_at,
888 ],
889 )?;
890 if changed == 0 {
891 return Err(SqliteError::InvalidData(format!(
892 "attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
893 )));
894 }
895 Ok(())
896}
897
898fn finalize_attachment_cutover_on_connection(
899 conn: &Connection,
900 now: i64,
901) -> Result<(), SqliteError> {
902 require_legacy_attachment_fences(conn)?;
903 validate_canonical_legacy_refs(conn)?;
904 validate_canonical_attachment_and_claim_refs(conn)?;
905 validate_attachment_record_owners(conn)?;
906 validate_legacy_content_backfill(conn)?;
907
908 let remaining_claims: i64 =
909 conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
910 if remaining_claims != 0 {
911 return Err(SqliteError::InvalidData(format!(
912 "attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
913 )));
914 }
915
916 let uncovered_model: Option<String> = conn
917 .query_row(
918 "SELECT model.id FROM entities AS model \
919 WHERE model.entity_type = 'moodboard_model' \
920 AND model.content_ref IS NOT NULL \
921 AND NOT EXISTS ( \
922 SELECT 1 FROM attachments AS attachment \
923 WHERE attachment.record_uuid = model.id \
924 AND attachment.substrate = 'entity' \
925 AND attachment.role = 'fann-network' \
926 ) \
927 LIMIT 1",
928 [],
929 |row| row.get(0),
930 )
931 .optional()?;
932 if let Some(record_uuid) = uncovered_model {
933 return Err(SqliteError::InvalidData(format!(
934 "moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
935 )));
936 }
937
938 conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
939 conn.execute_batch(
940 "DROP TRIGGER entities_reject_claimed_blob_insert; \
941 DROP TRIGGER entities_reject_claimed_blob_update; \
942 DROP INDEX idx_entities_content_ref; \
943 ALTER TABLE entities DROP COLUMN content_ref;",
944 )?;
945 conn.execute(
946 "UPDATE attachment_cutover_state \
947 SET state = 'complete', completed_at = ?1 \
948 WHERE singleton = 1 AND state = 'incomplete'",
949 [now],
950 )?;
951 Ok(())
952}
953
954fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
955 let migration = MIGRATIONS
956 .iter()
957 .find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
958 .expect("V21 migration must be registered");
959 conn.execute(
960 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
961 rusqlite::params![migration.version, migration.name, now],
962 )?;
963 Ok(())
964}
965
966pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
973 if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
974 return Ok(());
975 }
976 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
977 match attachment_cutover_status(&tx)? {
978 AttachmentCutoverStatus::Complete => return Ok(()),
979 AttachmentCutoverStatus::Pending => {
980 return Err(SqliteError::InvalidData(
981 "attachment cutover must complete stage 1 before finalization".into(),
982 ))
983 }
984 AttachmentCutoverStatus::Incomplete => {}
985 }
986 let now = chrono::Utc::now().timestamp_micros();
987 finalize_attachment_cutover_on_connection(&tx, now)?;
988 record_attachment_cutover_migration(&tx, now)?;
989 tx.commit()?;
990 Ok(())
991}
992
993fn read_applied_migration_ledger(
995 conn: &Connection,
996 through_version: u32,
997) -> Result<Vec<(u32, String)>, SqliteError> {
998 let mut stmt = conn.prepare(
999 "SELECT version, name FROM _schema_migrations \
1000 WHERE version <= ?1 ORDER BY version ASC",
1001 )?;
1002 let rows = stmt
1003 .query_map([through_version], |row| {
1004 Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
1005 })?
1006 .collect::<Result<Vec<_>, _>>()?;
1007 Ok(rows)
1008}
1009
1010fn validate_applied_migration_versions(
1015 applied: &[(u32, String)],
1016 through_version: u32,
1017) -> Result<(), SqliteError> {
1018 let expected: Vec<&VersionedMigration> = MIGRATIONS
1019 .iter()
1020 .filter(|migration| migration.version <= through_version)
1021 .collect();
1022 let mut applied_index = 0;
1023
1024 for migration in expected {
1025 let Some((version, applied_name)) = applied.get(applied_index) else {
1026 return Err(SqliteError::InvalidData(format!(
1027 "migration history is missing version {} ('{}'); the applied ledger must be \
1028 the exact contiguous canonical sequence through version {through_version}",
1029 migration.version, migration.name,
1030 )));
1031 };
1032 if *version < migration.version {
1033 return Err(SqliteError::InvalidData(format!(
1034 "migration history contains unknown version {version} recorded as \
1035 '{applied_name}'; the applied ledger must contain only canonical versions"
1036 )));
1037 }
1038 if *version > migration.version {
1039 return Err(SqliteError::InvalidData(format!(
1040 "migration history is missing version {} ('{}'); found version {version} \
1041 next instead",
1042 migration.version, migration.name,
1043 )));
1044 }
1045 applied_index += 1;
1046 }
1047
1048 if let Some((version, name)) = applied.get(applied_index) {
1049 return Err(SqliteError::InvalidData(format!(
1050 "migration history contains unknown version {version} recorded as '{name}'; \
1051 the applied ledger must contain only canonical versions"
1052 )));
1053 }
1054
1055 Ok(())
1056}
1057
1058fn validate_applied_migration_names(
1059 applied: &[(u32, String)],
1060 through_version: u32,
1061 allow_known_v19_repairs: bool,
1062) -> Result<(), SqliteError> {
1063 for ((version, applied_name), migration) in applied.iter().zip(
1064 MIGRATIONS
1065 .iter()
1066 .filter(|migration| migration.version <= through_version),
1067 ) {
1068 debug_assert_eq!(*version, migration.version);
1069 if migration.name != applied_name.as_str() {
1070 if allow_known_v19_repairs && matches!(*version, 13 | 14) {
1071 continue;
1072 }
1073 return Err(SqliteError::InvalidData(format!(
1074 "migration version {version} is recorded under name '{applied_name}', \
1075 expected '{expected}'. This database's migration history does not match \
1076 the current binary; recreate it from the current schema or repair the \
1077 specific known divergence via a dedicated migration.",
1078 expected = migration.name,
1079 )));
1080 }
1081 }
1082
1083 Ok(())
1084}
1085
1086fn validate_applied_migration_ledger(
1090 conn: &Connection,
1091 through_version: u32,
1092) -> Result<(), SqliteError> {
1093 let applied = read_applied_migration_ledger(conn, through_version)?;
1094 validate_applied_migration_versions(&applied, through_version)?;
1095 validate_applied_migration_names(&applied, through_version, false)
1096}
1097
1098const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
1099
1100pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
1106 match conn.query_row(
1107 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1108 [],
1109 |row| row.get(0),
1110 ) {
1111 Ok(version) => Ok(version),
1112 Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
1113 if msg.contains("no such table: _schema_migrations") =>
1114 {
1115 Ok(0)
1116 }
1117 Err(e) => Err(e.into()),
1118 }
1119}
1120
1121pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
1126 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1127 read_schema_version(&conn)
1128}
1129
1130pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
1133 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1134 validate_schema_is_current(&conn)
1135}
1136
1137pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
1145 let current_version = read_schema_version(conn)?;
1146 let latest_version = latest_schema_version();
1147
1148 if current_version < latest_version {
1149 return Err(SqliteError::InvalidData(format!(
1150 "read-only database schema version {current_version} is behind the latest known \
1151 migration {latest_version}; migrate a writable copy with this build before opening \
1152 the snapshot read-only"
1153 )));
1154 }
1155 if current_version > latest_version {
1156 return Err(SqliteError::InvalidData(format!(
1157 "read-only database schema version {current_version} is ahead of the latest known \
1158 migration {latest_version}; use a compatible newer build or recreate the snapshot"
1159 )));
1160 }
1161
1162 validate_applied_migration_ledger(conn, current_version)?;
1169 if current_version >= ATTACHMENT_CUTOVER_VERSION
1173 && attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
1174 {
1175 return Err(SqliteError::InvalidData(
1176 "read-only database has not completed the V21 attachment cutover".into(),
1177 ));
1178 }
1179
1180 Ok(current_version)
1181}
1182
1183#[cfg(test)]
1184pub(crate) mod test_sync {
1185 use std::sync::atomic::AtomicU32;
1186 use std::sync::{Arc, Barrier, Mutex};
1187
1188 pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
1192 pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
1194 pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
1198 std::sync::atomic::AtomicBool::new(false);
1199
1200 pub(crate) fn record_busy(_count: i32) -> bool {
1203 BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
1204 std::thread::sleep(std::time::Duration::from_millis(1));
1205 true
1206 }
1207
1208 pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
1211 std::sync::atomic::AtomicBool::new(false);
1212 pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
1217 std::sync::atomic::AtomicBool::new(false);
1218
1219 std::thread_local! {
1220 pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
1223 const { std::cell::Cell::new(false) };
1224 pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
1226 const { std::cell::Cell::new(false) };
1227 }
1228}
1229
1230fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
1239 let configured = conn.path().unwrap_or_default();
1240 let raw_path = if configured.is_empty() {
1241 conn.query_row(
1242 "SELECT file FROM pragma_database_list WHERE name = 'main'",
1243 [],
1244 |row| row.get::<_, String>(0),
1245 )?
1246 } else {
1247 configured.to_string()
1248 };
1249
1250 if raw_path.is_empty() {
1251 return Ok(None);
1252 }
1253 std::fs::canonicalize(&raw_path)
1254 .map(Some)
1255 .map_err(SqliteError::Io)
1256}
1257
1258fn validate_database_gc_owner(
1259 conn: &Connection,
1260 owner: &DatabaseGcOwnerGuard,
1261) -> Result<(), SqliteError> {
1262 let connection_path = canonical_connection_database_path(conn)?;
1263 if owner.database_path() != connection_path.as_deref() {
1264 return Err(SqliteError::InvalidData(format!(
1265 "database GC owner targets {:?}, but migration connection targets {:?}",
1266 owner.database_path(),
1267 connection_path.as_deref(),
1268 )));
1269 }
1270 Ok(())
1271}
1272
1273pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
1274 let database_path = canonical_connection_database_path(conn)?;
1275 if let Some(database_path) = database_path {
1276 let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
1281 SqliteError::InvalidData(format!(
1282 "failed to acquire database GC owner before schema migration: {error}"
1283 ))
1284 })?;
1285 return run_migrations_with_database_gc_owner(conn, &owner);
1286 }
1287
1288 run_migrations_with_busy_timeout(conn)
1291}
1292
1293pub(crate) fn run_migrations_with_database_gc_owner(
1294 conn: &mut Connection,
1295 owner: &DatabaseGcOwnerGuard,
1296) -> Result<u32, SqliteError> {
1297 validate_database_gc_owner(conn, owner)?;
1298 run_migrations_with_busy_timeout(conn)
1299}
1300
1301fn run_migrations_with_busy_timeout(conn: &mut Connection) -> Result<u32, SqliteError> {
1302 let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
1307 let raised = prior_busy_ms < 5_000;
1308 if raised {
1309 conn.busy_timeout(std::time::Duration::from_secs(5))?;
1310 }
1311 let result = run_migrations_locked(conn);
1312 if raised {
1313 let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
1314 }
1315 result
1316}
1317
1318fn run_migrations_locked(conn: &mut Connection) -> Result<u32, SqliteError> {
1319 conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
1320
1321 let current_version: u32 = read_schema_version(conn)?;
1322
1323 #[cfg(test)]
1327 if test_sync::PARTICIPATE.with(|p| p.get()) {
1328 conn.busy_handler(Some(test_sync::record_busy))?;
1331 let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
1332 if let Some(barrier) = barrier {
1333 barrier.wait();
1334 }
1335 }
1336
1337 let latest_version = latest_schema_version();
1343 if current_version > latest_version {
1344 return Err(SqliteError::InvalidData(format!(
1345 "database schema version {current_version} is ahead of the latest known migration \
1346 {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
1347 written by a newer build. Recreate it from the current schema; in-place downgrade is \
1348 not supported."
1349 )));
1350 }
1351
1352 let applied = read_applied_migration_ledger(conn, current_version)?;
1359 validate_applied_migration_versions(&applied, current_version)?;
1360 validate_applied_migration_names(&applied, current_version, current_version < 19)?;
1361
1362 let mut applied_version = current_version;
1363 let mut skip_through = current_version;
1367
1368 for migration in MIGRATIONS {
1369 if migration.version <= skip_through {
1370 applied_version = applied_version.max(migration.version);
1371 continue;
1372 }
1373
1374 #[cfg(test)]
1378 let instrumented_first_begin = test_sync::PARTICIPATE.with(|p| p.get())
1379 && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
1380 #[cfg(test)]
1381 if instrumented_first_begin {
1382 test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
1383 }
1384 let tx = conn
1385 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1386 .map_err(|e| SqliteError::Migration {
1387 version: migration.version,
1388 error: e.to_string(),
1389 })?;
1390
1391 let sibling_version: u32 = tx
1395 .query_row(
1396 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1397 [],
1398 |row| row.get(0),
1399 )
1400 .map_err(|e| SqliteError::Migration {
1401 version: migration.version,
1402 error: e.to_string(),
1403 })?;
1404 #[cfg(test)]
1405 if instrumented_first_begin {
1406 use std::sync::atomic::Ordering::SeqCst;
1407 if sibling_version == 0 {
1408 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1414 while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline
1415 {
1416 std::thread::yield_now();
1417 }
1418 } else {
1419 test_sync::LOSER_SAW_WINNER_COMMIT
1423 .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
1424 }
1425 }
1426
1427 if sibling_version > latest_version {
1432 return Err(SqliteError::InvalidData(format!(
1433 "database schema version {sibling_version} is ahead of the latest known \
1434 migration {latest_version} (committed by a concurrent process while this \
1435 one waited for the migration write lock). This build cannot run against \
1436 the newer schema; upgrade the binary or recreate the database."
1437 )));
1438 }
1439
1440 if sibling_version >= migration.version {
1441 #[cfg(test)]
1442 test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1443 skip_through = sibling_version.min(latest_version);
1444 applied_version = applied_version.max(migration.version);
1445 continue;
1446 }
1447
1448 if migration.version == ATTACHMENT_CUTOVER_VERSION {
1449 let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
1450 version: migration.version,
1451 error: e.to_string(),
1452 })?;
1453 let legacy_refs: i64 = tx
1454 .query_row(
1455 "SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
1456 [],
1457 |row| row.get(0),
1458 )
1459 .map_err(|e| SqliteError::Migration {
1460 version: migration.version,
1461 error: e.to_string(),
1462 })?;
1463
1464 if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
1468 drop(tx);
1469 break;
1470 }
1471 if status != AttachmentCutoverStatus::Pending {
1472 return Err(SqliteError::Migration {
1473 version: migration.version,
1474 error: format!("unexpected attachment cutover state {status:?}"),
1475 });
1476 }
1477
1478 let now = chrono::Utc::now().timestamp_micros();
1479 stage_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1480 SqliteError::Migration {
1481 version: migration.version,
1482 error: e.to_string(),
1483 }
1484 })?;
1485 finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1486 SqliteError::Migration {
1487 version: migration.version,
1488 error: e.to_string(),
1489 }
1490 })?;
1491 } else if migration.name == SESSION_IDENTITY_MIGRATION_NAME {
1492 tx.execute_batch(migration.up)
1493 .map_err(|error| SqliteError::Migration {
1494 version: migration.version,
1495 error: error.to_string(),
1496 })?;
1497 session_identity_migration::apply(&tx).map_err(|error| SqliteError::Migration {
1498 version: migration.version,
1499 error: error.to_string(),
1500 })?;
1501 } else {
1502 tx.execute_batch(migration.up)
1503 .map_err(|e| SqliteError::Migration {
1504 version: migration.version,
1505 error: e.to_string(),
1506 })?;
1507 }
1508
1509 if migration.version == 19 {
1516 tx.execute_batch(
1517 "UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
1518 UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
1519 )
1520 .map_err(|e| SqliteError::Migration {
1521 version: migration.version,
1522 error: e.to_string(),
1523 })?;
1524 }
1525
1526 let now = chrono::Utc::now().timestamp_micros();
1527 tx.execute(
1528 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
1529 ON CONFLICT(version) DO NOTHING",
1530 rusqlite::params![migration.version, migration.name, now],
1531 )
1532 .map_err(|e| SqliteError::Migration {
1533 version: migration.version,
1534 error: e.to_string(),
1535 })?;
1536
1537 #[cfg(test)]
1538 if instrumented_first_begin {
1539 test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
1540 }
1541
1542 tx.commit().map_err(|e| SqliteError::Migration {
1543 version: migration.version,
1544 error: e.to_string(),
1545 })?;
1546
1547 applied_version = migration.version;
1548 }
1549
1550 validate_applied_migration_ledger(conn, applied_version)?;
1554
1555 Ok(applied_version)
1556}
1557
1558#[derive(Debug)]
1559pub struct EmbeddingModelRegistryRecord {
1560 pub engine_name: String,
1562 pub model_id: String,
1564 pub key_version: String,
1566 pub dimensions: u32,
1568 pub status: String,
1570 pub activated_at: Option<i64>,
1572 pub superseded_at: Option<i64>,
1574}
1575
1576pub fn query_embedding_models(
1582 db: Option<&std::path::Path>,
1583 engine_filter: Option<&str>,
1584) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1585 let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
1586 std::env::var("HOME")
1587 .map(std::path::PathBuf::from)
1588 .unwrap_or_else(|_| std::path::PathBuf::from("."))
1589 .join(".khive/khive.db")
1590 });
1591 if !path.exists() {
1592 return Ok(Vec::new());
1593 }
1594 let conn = Connection::open_with_flags(
1595 path,
1596 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
1597 | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
1598 | rusqlite::OpenFlags::SQLITE_OPEN_URI,
1599 )?;
1600 query_embedding_models_conn(&conn, engine_filter)
1601}
1602
1603pub(crate) fn query_embedding_models_conn(
1607 conn: &Connection,
1608 engine_filter: Option<&str>,
1609) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
1610 let exists: bool = conn.query_row(
1611 "SELECT COUNT(*) > 0 FROM sqlite_master \
1612 WHERE type='table' AND name='_embedding_models'",
1613 [],
1614 |row| row.get(0),
1615 )?;
1616 if !exists {
1617 return Ok(Vec::new());
1618 }
1619
1620 let sql = if engine_filter.is_some() {
1621 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1622 FROM _embedding_models WHERE engine_name = ?1 \
1623 ORDER BY engine_name, activated_at IS NULL, activated_at"
1624 } else {
1625 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
1626 FROM _embedding_models \
1627 ORDER BY engine_name, activated_at IS NULL, activated_at"
1628 };
1629 let mut stmt = conn.prepare(sql)?;
1630 let map_row = |row: &rusqlite::Row<'_>| {
1631 let dim_raw: i64 = row.get(3)?;
1632 let dimensions = u32::try_from(dim_raw).map_err(|_| {
1633 rusqlite::Error::FromSqlConversionFailure(
1634 3,
1635 rusqlite::types::Type::Integer,
1636 Box::new(std::io::Error::other(format!(
1637 "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
1638 u32::MAX,
1639 ))),
1640 )
1641 })?;
1642 Ok(EmbeddingModelRegistryRecord {
1643 engine_name: row.get(0)?,
1644 model_id: row.get(1)?,
1645 key_version: row.get(2)?,
1646 dimensions,
1647 status: row.get(4)?,
1648 activated_at: row.get(5)?,
1649 superseded_at: row.get(6)?,
1650 })
1651 };
1652
1653 if let Some(engine) = engine_filter {
1654 stmt.query_map([engine], map_row)?
1655 .collect::<Result<Vec<_>, _>>()
1656 .map_err(Into::into)
1657 } else {
1658 stmt.query_map([], map_row)?
1659 .collect::<Result<Vec<_>, _>>()
1660 .map_err(Into::into)
1661 }
1662}
1663
1664#[cfg(test)]
1669#[path = "entity_version_migration_measurement.rs"]
1670mod entity_version_measurement;
1671
1672#[cfg(test)]
1673#[path = "migrations_tests.rs"]
1674mod tests;