1use khive_storage::blob::ContentRef;
16use rusqlite::{Connection, OptionalExtension};
17use std::path::{Path, PathBuf};
18
19use crate::error::SqliteError;
20use crate::pool::WriteAdmission;
21use crate::stores::blob::{try_acquire_database_gc_owner_for_path, DatabaseGcOwnerGuard};
22
23#[path = "raw_migration_settlement.rs"]
24mod raw_migration_settlement;
25use raw_migration_settlement::{RawMigrationTransactions, RawMigrationWriteUnit};
26
27pub(crate) trait MigrationTransactions {
34 fn admitted<T>(
37 &mut self,
38 operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
39 ) -> Result<T, SqliteError>;
40}
41
42#[derive(Clone, Debug)]
48pub struct MigrationWritePolicy {
49 disk_guard: crate::EffectiveDiskGuardConfig,
50 volume_lock_dir: PathBuf,
51}
52
53impl MigrationWritePolicy {
54 pub fn new(
55 disk_guard: crate::EffectiveDiskGuardConfig,
56 volume_lock_dir: impl Into<PathBuf>,
57 ) -> Result<Self, SqliteError> {
58 disk_guard.validate()?;
59 let volume_lock_dir = volume_lock_dir.into();
60 if !volume_lock_dir.is_absolute() {
61 return Err(SqliteError::InvalidConfig(
62 "migration volume-lock directory must be absolute".to_string(),
63 ));
64 }
65 Ok(Self {
66 disk_guard,
67 volume_lock_dir,
68 })
69 }
70
71 pub fn from_environment() -> Result<Self, SqliteError> {
75 let disk_guard = crate::DiskGuardEnvironment::capture().resolve(None, None)?;
76 Self::new(disk_guard, crate::default_volume_lock_dir()?)
77 }
78
79 pub fn disk_guard_config(&self) -> crate::EffectiveDiskGuardConfig {
80 self.disk_guard
81 }
82
83 pub fn volume_lock_dir(&self) -> &Path {
84 &self.volume_lock_dir
85 }
86}
87
88#[path = "session_identity_migration.rs"]
89mod session_identity_migration;
90
91mod memory_visibility;
92
93pub struct Migration {
99 pub id: &'static str,
101 pub up_sql: &'static str,
103 pub down_sql: Option<&'static str>,
105 pub is_already_applied: Option<fn(&Connection) -> bool>,
109}
110
111pub struct ServiceSchemaPlan {
113 pub service: &'static str,
115 pub sqlite: &'static [Migration],
117 pub postgres: &'static [Migration],
119}
120
121const SCHEMA_VERSION_TABLE: &str = include_str!("../sql/schema-version-table.sql");
122
123pub fn apply_schema_plan(
126 conn: &mut Connection,
127 plan: &ServiceSchemaPlan,
128) -> Result<(), SqliteError> {
129 let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
130 apply_schema_plan_with_admission(
131 &mut RawMigrationTransactions::new(conn, &admission),
132 plan,
133 &admission,
134 )
135}
136
137pub fn apply_schema_plan_with_policy(
139 conn: &mut Connection,
140 plan: &ServiceSchemaPlan,
141 policy: &MigrationWritePolicy,
142) -> Result<(), SqliteError> {
143 let admission =
144 WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
145 apply_schema_plan_with_admission(
146 &mut RawMigrationTransactions::new(conn, &admission),
147 plan,
148 &admission,
149 )
150}
151
152pub(crate) fn apply_schema_plan_with_admission(
155 writes: &mut impl MigrationTransactions,
156 plan: &ServiceSchemaPlan,
157 admission: &WriteAdmission,
158) -> Result<(), SqliteError> {
159 writes.admitted(|conn| {
160 admission.check()?;
161 conn.execute_batch(SCHEMA_VERSION_TABLE)?;
162 require_autocommit(conn, "schema-version bootstrap")
163 })?;
164
165 for migration in plan.sqlite {
166 writes.admitted(|conn| apply_service_migration(conn, plan, migration, admission))?;
167 }
168
169 Ok(())
170}
171
172fn apply_service_migration(
175 conn: &Connection,
176 plan: &ServiceSchemaPlan,
177 migration: &Migration,
178 admission: &WriteAdmission,
179) -> Result<(), SqliteError> {
180 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
184 if let Err(error) = admission.check() {
185 let rollback = tx.rollback();
186 return Err(capacity_refusal_after_rollback(
187 conn,
188 rollback,
189 error,
190 "service schema migration",
191 ));
192 }
193
194 if let Some(check) = migration.is_already_applied {
196 if check(&tx) {
197 return Ok(());
198 }
199 }
200
201 let already: bool = tx.query_row(
203 "SELECT COUNT(*) > 0 FROM _schema_versions WHERE service = ?1 AND migration_id = ?2",
204 rusqlite::params![plan.service, migration.id],
205 |row| row.get(0),
206 )?;
207
208 if already {
209 return Ok(());
210 }
211
212 tx.execute_batch(migration.up_sql)?;
213
214 tx.execute(
215 "INSERT INTO _schema_versions (service, migration_id, applied_at) VALUES (?1, ?2, ?3)",
216 rusqlite::params![
217 plan.service,
218 migration.id,
219 chrono::Utc::now().timestamp_micros(),
220 ],
221 )?;
222 tx.commit()?;
223 Ok(())
224}
225
226fn require_autocommit(conn: &Connection, operation: &str) -> Result<(), SqliteError> {
227 if conn.is_autocommit() {
228 Ok(())
229 } else {
230 Err(SqliteError::InvalidData(format!(
231 "{operation} did not return the SQLite connection to autocommit"
232 )))
233 }
234}
235
236pub(crate) fn capacity_refusal_after_rollback(
237 conn: &Connection,
238 rollback: rusqlite::Result<()>,
239 refusal: SqliteError,
240 operation: &str,
241) -> SqliteError {
242 if let Err(error) = rollback {
243 return SqliteError::InvalidData(format!(
244 "{operation} capacity refusal could not roll back: {error}; \
245 initial refusal: {refusal}"
246 ));
247 }
248 if !conn.is_autocommit() {
249 return SqliteError::InvalidData(format!(
250 "{operation} capacity refusal rolled back without restoring autocommit; \
251 initial refusal: {refusal}"
252 ));
253 }
254 refusal
255}
256
257pub struct VersionedMigration {
267 pub version: u32,
269 pub name: &'static str,
271 pub up: &'static str,
274}
275
276const V1_UP: &str = include_str!("../sql/schema.sql");
279
280const V2_UP: &str = include_str!("../sql/002-narrow-fts-sections-update-trigger.sql");
281
282const V3_UP: &str = include_str!("../sql/003-backfill-domain-mirror-atoms.sql");
283
284const V4_UP: &str = include_str!("../sql/004-fts-consolidation.sql");
285
286const V5_UP: &str = include_str!("../sql/005-unique-comm-external-id.sql");
287
288const V6_UP: &str = include_str!("../sql/006-brain-retune-driver.sql");
289
290const V7_UP: &str = include_str!("../sql/007-notes-seq.sql");
291
292const V8_UP: &str = include_str!("../sql/008-notes-seq-repair.sql");
293
294const V9_UP: &str = include_str!("../sql/009-entities-name-ci-index.sql");
295
296const V10_UP: &str = include_str!("../sql/010-entities-content-ref.sql");
297
298const V11_UP: &str = include_str!("../sql/011-ann-write-log.sql");
299
300const V12_UP: &str = include_str!("../sql/012-ann-write-log-model-seq-index.sql");
301
302const V13_UP: &str = include_str!("../sql/013-list-cursor-sequences.sql");
303
304const V14_UP: &str = include_str!("../sql/014-graph-edges-id-unique.sql");
305
306const V15_UP: &str = include_str!("../sql/015-serve-ledger-attribution.sql");
307
308const V16_UP: &str = include_str!("../sql/016-gtd-dependency-cycle-guards.sql");
309
310const V17_UP: &str = include_str!("../sql/017-agents-ddl.sql");
311
312const V18_UP: &str = include_str!("../sql/018-ann-consumer-pending.sql");
313
314const V19_UP: &str = include_str!("../sql/019-list-cursor-backfill-repair.sql");
315
316const V20_UP: &str = include_str!("../sql/020-blob-gc-claims.sql");
317
318const V22_UP: &str = include_str!("../sql/022-notes-unread-probe-recipient.sql");
319
320const V23_UP: &str = include_str!("../sql/023-fts-record-kind.sql");
321
322const V24_UP: &str = include_str!("../sql/024-fts-rowid-map.sql");
323
324const V25_UP: &str = include_str!("../sql/025-notes-unread-probe-recipient-direction.sql");
325
326const V26_UP: &str = include_str!("../sql/026-knowledge-fts-repair.sql");
327
328const V27_UP: &str = include_str!("../sql/027-notes-hot-property-indexes.sql");
329const V28_UP: &str = include_str!("../sql/028-notes-key.sql");
330
331const V29_UP: &str = include_str!("../sql/029-note-streams.sql");
332const V30_UP: &str = include_str!("../sql/030-tool-source-mounts.sql");
333const V31_UP: &str = include_str!("../sql/031-note-versions.sql");
334const V32_UP: &str = include_str!("../sql/032-knowledge-count-indexes.sql");
335const V33_UP: &str = include_str!("../sql/033-notes-message-recipient-direction.sql");
336const V34_UP: &str = include_str!("../sql/034-notes-namespace-created.sql");
337const V35_UP: &str = include_str!("../sql/035-notes-unread-probe-recipient-type-direction.sql");
338const V36_UP: &str = include_str!("../sql/036-events-operation-attribution.sql");
339const V37_UP: &str = include_str!("../sql/037-entity-versions.sql");
340const V38_UP: &str = include_str!("../sql/038-entities-legacy-type-index.sql");
341const V39_UP: &str = include_str!("../sql/039-knowledge-cursor-indexes.sql");
342const SESSION_IDENTITY_UP: &str = include_str!("../sql/040-session-source-scope.sql");
343const SESSION_IDENTITY_MIGRATION_NAME: &str = "session_source_scoped_identity";
344const V41_UP: &str = include_str!("../sql/041-sender-transport.sql");
345const V42_UP: &str = include_str!("../sql/042-comm-external-id-channel-scope.sql");
346const V43_UP: &str = include_str!("../sql/043-vector-provenance.sql");
347const V44_COLUMNS: &str = include_str!("../sql/044-comm-outbound-due-a-columns.sql");
348const V44_UP: &str = include_str!("../sql/044-comm-outbound-due-b-index.sql");
349
350pub(crate) fn migrate_outbound_due_key(tx: &rusqlite::Transaction<'_>) -> rusqlite::Result<()> {
354 let has_column = |name: &str| -> rusqlite::Result<bool> {
355 tx.query_row(
356 "SELECT EXISTS(SELECT 1 FROM pragma_table_info('notes') WHERE name = ?1)",
357 [name],
358 |row| row.get(0),
359 )
360 };
361 match (has_column("strict_due_key")?, has_column("due_source")?) {
362 (false, false) => tx.execute_batch(V44_COLUMNS)?,
363 (true, true) => {}
364 _ => return Err(rusqlite::Error::InvalidQuery),
365 }
366
367 let mut after_id = String::new();
368 let mut first_page = true;
369 loop {
370 let rows: Vec<(String, String)> = {
371 let comparator = if first_page { ">=" } else { ">" };
372 let mut stmt = tx.prepare(&format!(
373 "SELECT id, json_extract(properties, '$.next_attempt_at') FROM notes \
374 WHERE id {comparator} ?1 AND json_type(properties, '$.next_attempt_at') = 'text' \
375 ORDER BY id LIMIT 500"
376 ))?;
377 let collected = stmt
378 .query_map([&after_id], |row| Ok((row.get(0)?, row.get(1)?)))?
379 .collect::<rusqlite::Result<_>>()?;
380 collected
381 };
382 if rows.is_empty() {
383 break;
384 }
385 after_id = rows.last().expect("nonempty V44 page").0.clone();
386 first_page = false;
387 for (id, source) in rows {
388 if let Some(key) = crate::pool::strict_rfc3339_key(&source) {
389 tx.execute(
390 "UPDATE notes SET strict_due_key = ?1, due_source = ?2 \
391 WHERE id = ?3 AND (strict_due_key IS NOT ?1 OR due_source IS NOT ?2)",
392 rusqlite::params![key, source, id],
393 )?;
394 }
395 }
396 }
397 tx.execute_batch(V44_UP)
398}
399
400const RECIPIENT_TRANSPORT_VERSION: u32 = 45;
401const V45_UP: &str = include_str!("../sql/045-recipient-transport.sql");
402const V46_UP: &str = include_str!("../sql/046-memory-visibility-receipts.sql");
403const V47_UP: &str = include_str!("../sql/047-attachment-role-quarantine.sql");
404const V49_UP: &str = include_str!("../sql/049-git-note-property-indexes.sql");
405
406const V50_UP: &str = include_str!("../sql/050-entity-list-plans.sql");
407const V51_UP: &str = include_str!("../sql/051-schedule-core-indexes.sql");
408const V52_UP: &str = include_str!("../sql/052-comm-core-indexes.sql");
409const V53_UP: &str = include_str!("../sql/053-entity-kind-list-order.sql");
410const MEMORY_VISIBILITY_CUTOVER_VERSION: u32 = 54;
411const V54_UP: &str = include_str!("../sql/054-memory-visibility-epochs.sql");
412const V48_UP: &str = include_str!("../sql/048-acknowledgement-journal-a-table.sql");
413const ACKNOWLEDGEMENT_JOURNAL_INDEX: &str =
414 include_str!("../sql/048-acknowledgement-journal-b-index.sql");
415
416pub(crate) fn migrate_acknowledgement_journal(
418 tx: &rusqlite::Transaction<'_>,
419) -> rusqlite::Result<()> {
420 let columns: i64 = tx.query_row(
421 "SELECT count(*) FROM pragma_table_info('comm_ack_work') \
422 WHERE name IN ('attempt_count','not_before','retirement_reason')",
423 [],
424 |row| row.get(0),
425 )?;
426 match columns {
427 0 => tx.execute_batch(V48_UP)?,
428 3 => {}
429 _ => return Err(rusqlite::Error::InvalidQuery),
430 }
431 tx.execute_batch(ACKNOWLEDGEMENT_JOURNAL_INDEX)
432}
433
434const V21_STAGE_UP: &str = include_str!("../sql/021-attachments-a-stage.sql");
435
436const V21_ATTACHMENT_FENCES_UP: &str = include_str!("../sql/021-attachments-b-claim-fences.sql");
437
438pub const ATTACHMENT_CUTOVER_VERSION: u32 = 21;
440
441pub fn latest_schema_version() -> u32 {
446 MIGRATIONS.last().map(|m| m.version).unwrap_or(0)
447}
448
449pub const ANN_WRITE_LOG_DDL: &str = V11_UP;
457
458pub const ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL: &str = V12_UP;
463
464pub const ANN_CONSUMER_PENDING_DDL: &str = include_str!("../sql/ann-consumer-pending-ddl.sql");
471
472pub const VECTOR_PROVENANCE_DDL: &str = V43_UP;
475
476pub const EMBEDDING_MODELS_DDL: &str = include_str!("../sql/embedding-models-ddl.sql");
482
483pub const MIGRATIONS: &[VersionedMigration] = &[
490 VersionedMigration {
491 version: 1,
492 name: "initial_schema",
493 up: V1_UP,
494 },
495 VersionedMigration {
496 version: 2,
497 name: "narrow_fts_sections_update_trigger",
498 up: V2_UP,
499 },
500 VersionedMigration {
501 version: 3,
502 name: "backfill_domain_mirror_atoms",
503 up: V3_UP,
504 },
505 VersionedMigration {
506 version: 4,
507 name: "fts_consolidation",
508 up: V4_UP,
509 },
510 VersionedMigration {
511 version: 5,
512 name: "unique_comm_message_external_id",
513 up: V5_UP,
514 },
515 VersionedMigration {
516 version: 6,
517 name: "brain_retune_driver",
518 up: V6_UP,
519 },
520 VersionedMigration {
521 version: 7,
522 name: "notes_seq",
523 up: V7_UP,
524 },
525 VersionedMigration {
526 version: 8,
527 name: "notes_seq_repair",
528 up: V8_UP,
529 },
530 VersionedMigration {
531 version: 9,
532 name: "entities_name_ci_index",
533 up: V9_UP,
534 },
535 VersionedMigration {
536 version: 10,
537 name: "entities_content_ref",
538 up: V10_UP,
539 },
540 VersionedMigration {
541 version: 11,
542 name: "ann_write_log",
543 up: V11_UP,
544 },
545 VersionedMigration {
546 version: 12,
547 name: "ann_write_log_model_seq_index",
548 up: V12_UP,
549 },
550 VersionedMigration {
551 version: 13,
552 name: "list_cursor_sequences",
553 up: V13_UP,
554 },
555 VersionedMigration {
556 version: 14,
557 name: "graph_edges_id_unique",
558 up: V14_UP,
559 },
560 VersionedMigration {
561 version: 15,
562 name: "serve_ledger_attribution",
563 up: V15_UP,
564 },
565 VersionedMigration {
566 version: 16,
567 name: "gtd_dependency_cycle_guards",
568 up: V16_UP,
569 },
570 VersionedMigration {
571 version: 17,
572 name: "agents_ddl",
573 up: V17_UP,
574 },
575 VersionedMigration {
576 version: 18,
577 name: "ann_consumer_pending",
578 up: V18_UP,
579 },
580 VersionedMigration {
581 version: 19,
582 name: "list_cursor_backfill_repair",
583 up: V19_UP,
584 },
585 VersionedMigration {
586 version: 20,
587 name: "blob_gc_claims",
588 up: V20_UP,
589 },
590 VersionedMigration {
591 version: ATTACHMENT_CUTOVER_VERSION,
592 name: "attachments_first_class",
593 up: V21_STAGE_UP,
597 },
598 VersionedMigration {
599 version: 22,
600 name: "notes_unread_probe_recipient",
601 up: V22_UP,
602 },
603 VersionedMigration {
604 version: 23,
605 name: "fts_record_kind",
606 up: V23_UP,
607 },
608 VersionedMigration {
609 version: 24,
610 name: "fts_rowid_map",
611 up: V24_UP,
612 },
613 VersionedMigration {
614 version: 25,
615 name: "notes_unread_probe_recipient_direction",
616 up: V25_UP,
617 },
618 VersionedMigration {
619 version: 26,
620 name: "knowledge_fts_repair",
621 up: V26_UP,
622 },
623 VersionedMigration {
624 version: 27,
625 name: "notes_hot_property_indexes",
626 up: V27_UP,
627 },
628 VersionedMigration {
629 version: 28,
630 name: "notes_key",
631 up: V28_UP,
632 },
633 VersionedMigration {
634 version: 29,
635 name: "note_streams",
636 up: V29_UP,
637 },
638 VersionedMigration {
639 version: 30,
640 name: "tool_source_mounts",
641 up: V30_UP,
642 },
643 VersionedMigration {
644 version: 31,
645 name: "note_versions",
646 up: V31_UP,
647 },
648 VersionedMigration {
649 version: 32,
650 name: "knowledge_count_indexes",
651 up: V32_UP,
652 },
653 VersionedMigration {
654 version: 33,
655 name: "notes_message_recipient_direction",
656 up: V33_UP,
657 },
658 VersionedMigration {
659 version: 34,
660 name: "notes_namespace_created",
661 up: V34_UP,
662 },
663 VersionedMigration {
664 version: 35,
665 name: "notes_unread_probe_recipient_type_direction",
666 up: V35_UP,
667 },
668 VersionedMigration {
669 version: 36,
670 name: "events_operation_attribution",
671 up: V36_UP,
672 },
673 VersionedMigration {
674 version: 37,
675 name: "entity_versions",
676 up: V37_UP,
677 },
678 VersionedMigration {
679 version: 38,
680 name: "entities_legacy_type_index",
681 up: V38_UP,
682 },
683 VersionedMigration {
684 version: 39,
685 name: "knowledge_cursor_indexes",
686 up: V39_UP,
687 },
688 VersionedMigration {
689 version: 40,
690 name: SESSION_IDENTITY_MIGRATION_NAME,
691 up: SESSION_IDENTITY_UP,
692 },
693 VersionedMigration {
694 version: 41,
695 name: "sender_transport",
696 up: V41_UP,
697 },
698 VersionedMigration {
699 version: 42,
700 name: "comm_external_id_channel_scope",
701 up: V42_UP,
702 },
703 VersionedMigration {
704 version: 43,
705 name: "vector_provenance",
706 up: V43_UP,
707 },
708 VersionedMigration {
709 version: 44,
710 name: "comm_outbound_due",
711 up: V44_UP,
712 },
713 VersionedMigration {
714 version: RECIPIENT_TRANSPORT_VERSION,
715 name: "recipient_transport",
716 up: V45_UP,
717 },
718 VersionedMigration {
719 version: 46,
720 name: "memory_visibility_receipts",
721 up: V46_UP,
722 },
723 VersionedMigration {
724 version: 47,
725 name: "attachment_role_quarantine",
726 up: V47_UP,
727 },
728 VersionedMigration {
729 version: 48,
730 name: "acknowledgement_journal",
731 up: V48_UP,
732 },
733 VersionedMigration {
734 version: 49,
735 name: "git_note_property_indexes",
736 up: V49_UP,
737 },
738 VersionedMigration {
739 version: 50,
740 name: "entity_list_plans",
741 up: V50_UP,
742 },
743 VersionedMigration {
744 version: 51,
745 name: "schedule_core_indexes",
746 up: V51_UP,
747 },
748 VersionedMigration {
749 version: 52,
750 name: "comm_core_indexes",
751 up: V52_UP,
752 },
753 VersionedMigration {
754 version: 53,
755 name: "entity_kind_list_order",
756 up: V53_UP,
757 },
758 VersionedMigration {
759 version: MEMORY_VISIBILITY_CUTOVER_VERSION,
760 name: "memory_visibility_epochs",
761 up: V54_UP,
762 },
763];
764
765#[derive(Clone, Copy, Debug, Eq, PartialEq)]
767pub enum AttachmentCutoverStatus {
768 Pending,
770 Incomplete,
772 Complete,
774}
775
776fn schema_object_exists(
777 conn: &Connection,
778 object_type: &str,
779 name: &str,
780) -> Result<bool, SqliteError> {
781 conn.query_row(
782 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type = ?1 AND name = ?2",
783 rusqlite::params![object_type, name],
784 |row| row.get(0),
785 )
786 .map_err(Into::into)
787}
788
789fn schema_column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool, SqliteError> {
790 conn.query_row(
791 "SELECT COUNT(*) > 0 FROM pragma_table_info(?1) WHERE name = ?2",
792 rusqlite::params![table, column],
793 |row| row.get(0),
794 )
795 .map_err(Into::into)
796}
797
798fn require_attachment_schema_objects(
799 conn: &Connection,
800 objects: &[(&str, &str)],
801 phase: &str,
802) -> Result<(), SqliteError> {
803 for (object_type, name) in objects {
804 if !schema_object_exists(conn, object_type, name)? {
805 return Err(SqliteError::InvalidData(format!(
806 "attachment cutover {phase} state is missing {object_type} {name:?}"
807 )));
808 }
809 }
810 Ok(())
811}
812
813fn validate_incomplete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
814 require_attachment_schema_objects(
815 conn,
816 &[
817 ("table", "attachments"),
818 ("index", "idx_attachments_content_ref"),
819 ],
820 "incomplete",
821 )?;
822 require_legacy_attachment_fences(conn)
823}
824
825fn validate_complete_attachment_schema(conn: &Connection) -> Result<(), SqliteError> {
826 require_attachment_schema_objects(
827 conn,
828 &[
829 ("table", "attachments"),
830 ("table", "blob_gc_claims"),
831 ("index", "idx_attachments_content_ref"),
832 ("index", "idx_blob_gc_claims_content_ref"),
833 ("trigger", "attachments_reject_claimed_blob_insert"),
834 ("trigger", "attachments_reject_claimed_blob_update"),
835 ],
836 "complete",
837 )?;
838 if schema_column_exists(conn, "entities", "content_ref")? {
839 return Err(SqliteError::InvalidData(
840 "attachment cutover is complete but entities.content_ref still exists".into(),
841 ));
842 }
843 for (object_type, name) in [
844 ("index", "idx_entities_content_ref"),
845 ("trigger", "entities_reject_claimed_blob_insert"),
846 ("trigger", "entities_reject_claimed_blob_update"),
847 ] {
848 if schema_object_exists(conn, object_type, name)? {
849 return Err(SqliteError::InvalidData(format!(
850 "attachment cutover is complete but legacy {object_type} {name:?} still exists"
851 )));
852 }
853 }
854 Ok(())
855}
856
857pub fn attachment_cutover_status(
862 conn: &Connection,
863) -> Result<AttachmentCutoverStatus, SqliteError> {
864 let version = read_schema_version(conn)?;
865 let marker_table = schema_object_exists(conn, "table", "attachment_cutover_state")?;
866 if !marker_table {
867 if version >= ATTACHMENT_CUTOVER_VERSION {
868 return Err(SqliteError::InvalidData(format!(
869 "migration V{ATTACHMENT_CUTOVER_VERSION} is recorded but its attachment cutover marker is absent"
870 )));
871 }
872 if schema_object_exists(conn, "table", "attachments")? {
873 return Err(SqliteError::InvalidData(
874 "attachments table exists without the durable attachment cutover marker".into(),
875 ));
876 }
877 return Ok(AttachmentCutoverStatus::Pending);
878 }
879
880 let marker: Option<(String, Option<i64>)> = conn
881 .query_row(
882 "SELECT state, completed_at FROM attachment_cutover_state WHERE singleton = 1",
883 [],
884 |row| Ok((row.get(0)?, row.get(1)?)),
885 )
886 .optional()?;
887 match marker {
888 Some((state, None)) if state == "incomplete" => {
889 if version >= ATTACHMENT_CUTOVER_VERSION {
890 Err(SqliteError::InvalidData(format!(
891 "attachment cutover is incomplete but migration V{ATTACHMENT_CUTOVER_VERSION} is already recorded"
892 )))
893 } else {
894 validate_incomplete_attachment_schema(conn)?;
895 Ok(AttachmentCutoverStatus::Incomplete)
896 }
897 }
898 Some((state, Some(_))) if state == "complete" => {
899 if version >= ATTACHMENT_CUTOVER_VERSION {
903 validate_complete_attachment_schema(conn)?;
904 Ok(AttachmentCutoverStatus::Complete)
905 } else {
906 Err(SqliteError::InvalidData(format!(
907 "attachment cutover is complete but schema ledger is at V{version}, below V{ATTACHMENT_CUTOVER_VERSION}"
908 )))
909 }
910 }
911 Some((state, completed_at)) => Err(SqliteError::InvalidData(format!(
912 "invalid attachment cutover marker state {state:?} with completed_at={completed_at:?}"
913 ))),
914 None => Err(SqliteError::InvalidData(
915 "attachment cutover marker table exists without its singleton row".into(),
916 )),
917 }
918}
919
920fn require_legacy_attachment_fences(conn: &Connection) -> Result<(), SqliteError> {
921 if !schema_column_exists(conn, "entities", "content_ref")? {
922 return Err(SqliteError::InvalidData(
923 "attachment cutover requires legacy entities.content_ref until finalization".into(),
924 ));
925 }
926 for (object_type, name) in [
927 ("table", "blob_gc_claims"),
928 ("index", "idx_blob_gc_claims_content_ref"),
929 ("index", "idx_entities_content_ref"),
930 ("trigger", "entities_reject_claimed_blob_insert"),
931 ("trigger", "entities_reject_claimed_blob_update"),
932 ] {
933 if !schema_object_exists(conn, object_type, name)? {
934 return Err(SqliteError::InvalidData(format!(
935 "attachment cutover requires legacy {object_type} {name:?} until finalization"
936 )));
937 }
938 }
939 Ok(())
940}
941
942fn canonical_content_ref_byte_width(conn: &Connection) -> Result<i64, SqliteError> {
949 let width: i64 = conn.query_row("SELECT length(CAST('x' AS BLOB))", [], |row| row.get(0))?;
950 if !(1..=4).contains(&width) {
951 return Err(SqliteError::InvalidData(format!(
952 "the text-encoding width probe returned {width}; refusing canonicality validation"
953 )));
954 }
955 Ok(width * 64)
956}
957
958fn validate_canonical_legacy_refs(conn: &Connection) -> Result<(), SqliteError> {
959 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
960 let invalid: Option<String> = conn
961 .query_row(
962 "SELECT id FROM entities \
963 WHERE content_ref IS NOT NULL \
964 AND (typeof(content_ref) <> 'text' \
965 OR length(content_ref) <> 64 \
966 OR length(CAST(content_ref AS BLOB)) <> ?1 \
967 OR content_ref GLOB '*[^0-9a-f]*') \
968 LIMIT 1",
969 [canonical_bytes],
970 |row| row.get(0),
971 )
972 .optional()?;
973 if let Some(id) = invalid {
974 return Err(SqliteError::InvalidData(format!(
975 "entities.content_ref for record {id:?} is not a canonical 64-character lowercase hexadecimal ContentRef"
976 )));
977 }
978 Ok(())
979}
980
981fn validate_canonical_attachment_and_claim_refs(conn: &Connection) -> Result<(), SqliteError> {
982 let canonical_bytes = canonical_content_ref_byte_width(conn)?;
983 for (table, identity) in [
984 ("attachments", "record_uuid"),
985 ("blob_gc_claims", "root_key"),
986 ] {
987 let sql = format!(
988 "SELECT {identity} FROM {table} \
989 WHERE typeof(content_ref) <> 'text' \
990 OR length(content_ref) <> 64 \
991 OR length(CAST(content_ref AS BLOB)) <> ?1 \
992 OR content_ref GLOB '*[^0-9a-f]*' \
993 LIMIT 1"
994 );
995 let invalid: Option<String> = conn
996 .query_row(&sql, [canonical_bytes], |row| row.get(0))
997 .optional()?;
998 if let Some(owner) = invalid {
999 return Err(SqliteError::InvalidData(format!(
1000 "{table}.content_ref for {identity} {owner:?} is not canonical"
1001 )));
1002 }
1003 }
1004 Ok(())
1005}
1006
1007fn validate_attachment_record_owners(conn: &Connection) -> Result<(), SqliteError> {
1008 let dangling: Option<(String, String)> = conn
1009 .query_row(
1010 "SELECT record_uuid, substrate FROM attachments AS attachment \
1011 WHERE (substrate = 'entity' AND NOT EXISTS ( \
1012 SELECT 1 FROM entities WHERE id = attachment.record_uuid \
1013 )) \
1014 OR (substrate = 'note' AND NOT EXISTS ( \
1015 SELECT 1 FROM notes WHERE id = attachment.record_uuid \
1016 )) \
1017 LIMIT 1",
1018 [],
1019 |row| Ok((row.get(0)?, row.get(1)?)),
1020 )
1021 .optional()?;
1022 if let Some((record_uuid, substrate)) = dangling {
1023 return Err(SqliteError::InvalidData(format!(
1024 "attachment role references absent {substrate} record {record_uuid:?}"
1025 )));
1026 }
1027 Ok(())
1028}
1029
1030fn validate_legacy_content_backfill(conn: &Connection) -> Result<(), SqliteError> {
1031 let conflict: Option<String> = conn
1032 .query_row(
1033 "SELECT entity.id FROM entities AS entity \
1034 LEFT JOIN attachments AS attachment \
1035 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
1036 WHERE entity.content_ref IS NOT NULL \
1037 AND (attachment.record_uuid IS NULL \
1038 OR attachment.substrate <> 'entity' \
1039 OR attachment.content_ref <> entity.content_ref) \
1040 LIMIT 1",
1041 [],
1042 |row| row.get(0),
1043 )
1044 .optional()?;
1045 if let Some(record_uuid) = conflict {
1046 return Err(SqliteError::InvalidData(format!(
1047 "legacy content attachment for entity {record_uuid:?} is missing or conflicts with entities.content_ref"
1048 )));
1049 }
1050 Ok(())
1051}
1052
1053fn stage_attachment_cutover_on_connection(conn: &Connection, now: i64) -> Result<(), SqliteError> {
1054 require_legacy_attachment_fences(conn)?;
1055 conn.execute_batch(V21_STAGE_UP)?;
1056 validate_canonical_legacy_refs(conn)?;
1057 validate_canonical_attachment_and_claim_refs(conn)?;
1058
1059 let conflict: Option<String> = conn
1060 .query_row(
1061 "SELECT entity.id FROM entities AS entity \
1062 JOIN attachments AS attachment \
1063 ON attachment.record_uuid = entity.id AND attachment.role = 'content' \
1064 WHERE entity.content_ref IS NOT NULL \
1065 AND (attachment.substrate <> 'entity' \
1066 OR attachment.content_ref <> entity.content_ref) \
1067 LIMIT 1",
1068 [],
1069 |row| row.get(0),
1070 )
1071 .optional()?;
1072 if let Some(record_uuid) = conflict {
1073 return Err(SqliteError::InvalidData(format!(
1074 "existing content attachment for entity {record_uuid:?} conflicts with entities.content_ref"
1075 )));
1076 }
1077
1078 conn.execute(
1079 "INSERT INTO attachments \
1080 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
1081 SELECT id, 'entity', 'content', content_ref, NULL, NULL, created_at \
1082 FROM entities WHERE content_ref IS NOT NULL \
1083 ON CONFLICT(record_uuid, role) DO NOTHING",
1084 [],
1085 )?;
1086 validate_legacy_content_backfill(conn)?;
1087
1088 conn.execute("DELETE FROM blob_gc_claims", [])?;
1092 conn.execute(
1093 "INSERT INTO attachment_cutover_state \
1094 (singleton, state, started_at, completed_at) \
1095 VALUES (1, 'incomplete', ?1, NULL) \
1096 ON CONFLICT(singleton) DO NOTHING",
1097 [now],
1098 )?;
1099 Ok(())
1100}
1101
1102pub fn stage_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
1109 let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
1110 RawMigrationWriteUnit::new(conn, &admission)?
1111 .run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
1112}
1113
1114pub fn stage_attachment_cutover_with_policy(
1116 conn: &mut Connection,
1117 policy: &MigrationWritePolicy,
1118) -> Result<(), SqliteError> {
1119 let admission =
1120 WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
1121 RawMigrationWriteUnit::new(conn, &admission)?
1122 .run(|conn| stage_attachment_cutover_with_admission(conn, &admission))
1123}
1124
1125pub(crate) fn stage_attachment_cutover_with_admission(
1127 conn: &mut Connection,
1128 admission: &WriteAdmission,
1129) -> Result<(), SqliteError> {
1130 match attachment_cutover_status(conn)? {
1131 AttachmentCutoverStatus::Complete => return Ok(()),
1132 AttachmentCutoverStatus::Pending | AttachmentCutoverStatus::Incomplete => {}
1133 }
1134 if read_schema_version(conn)? != ATTACHMENT_CUTOVER_VERSION - 1 {
1135 return Err(SqliteError::InvalidData(format!(
1136 "attachment cutover stage requires canonical V{} schema",
1137 ATTACHMENT_CUTOVER_VERSION - 1
1138 )));
1139 }
1140
1141 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1142 if let Err(error) = admission.check() {
1143 let rollback = tx.rollback();
1144 return Err(capacity_refusal_after_rollback(
1145 conn,
1146 rollback,
1147 error,
1148 "attachment cutover stage",
1149 ));
1150 }
1151 let status = attachment_cutover_status(&tx)?;
1152 if status == AttachmentCutoverStatus::Complete {
1153 return Ok(());
1154 }
1155 stage_attachment_cutover_on_connection(&tx, chrono::Utc::now().timestamp_micros())?;
1156 tx.commit()?;
1157 Ok(())
1158}
1159
1160#[allow(clippy::too_many_arguments)]
1168pub fn apply_generic_verified_attachment(
1169 conn: &Connection,
1170 record_uuid: &str,
1171 substrate: &str,
1172 role: &str,
1173 content_ref: &ContentRef,
1174 media_type: Option<&str>,
1175 size_bytes: Option<u64>,
1176 created_at: i64,
1177) -> Result<(), SqliteError> {
1178 if attachment_cutover_status(conn)? != AttachmentCutoverStatus::Incomplete {
1179 return Err(SqliteError::InvalidData(
1180 "verified application attachments may only be applied while V21 cutover is incomplete"
1181 .into(),
1182 ));
1183 }
1184 if role.is_empty() || role.chars().any(char::is_control) {
1185 return Err(SqliteError::InvalidData(
1186 "attachment role must be non-empty and contain no control characters".into(),
1187 ));
1188 }
1189 let size_bytes = size_bytes.map(i64::try_from).transpose().map_err(|_| {
1190 SqliteError::InvalidData("attachment size_bytes exceeds SQLite INTEGER".into())
1191 })?;
1192 let owner_table = match substrate {
1193 "entity" => "entities",
1194 "note" => "notes",
1195 other => {
1196 return Err(SqliteError::InvalidData(format!(
1197 "attachment substrate must be 'entity' or 'note', got {other:?}"
1198 )))
1199 }
1200 };
1201 let owner_sql = format!("SELECT COUNT(*) > 0 FROM {owner_table} WHERE id = ?1");
1202 let owner_exists: bool = conn.query_row(&owner_sql, [record_uuid], |row| row.get(0))?;
1203 if !owner_exists {
1204 return Err(SqliteError::InvalidData(format!(
1205 "cannot attach role {role:?}: {substrate} record {record_uuid:?} does not exist"
1206 )));
1207 }
1208 let claimed: bool = conn.query_row(
1209 "SELECT COUNT(*) > 0 FROM blob_gc_claims WHERE content_ref = ?1",
1210 [content_ref.as_str()],
1211 |row| row.get(0),
1212 )?;
1213 if claimed {
1214 return Err(SqliteError::InvalidData(format!(
1215 "cannot attach claimed content_ref {} during V21 cutover",
1216 content_ref.as_str()
1217 )));
1218 }
1219
1220 let changed = conn.execute(
1221 "INSERT INTO attachments \
1222 (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) \
1223 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) \
1224 ON CONFLICT(record_uuid, role) DO UPDATE SET \
1225 media_type = excluded.media_type, \
1226 size_bytes = excluded.size_bytes, \
1227 created_at = excluded.created_at \
1228 WHERE attachments.substrate = excluded.substrate \
1229 AND attachments.content_ref = excluded.content_ref",
1230 rusqlite::params![
1231 record_uuid,
1232 substrate,
1233 role,
1234 content_ref.as_str(),
1235 media_type,
1236 size_bytes,
1237 created_at,
1238 ],
1239 )?;
1240 if changed == 0 {
1241 return Err(SqliteError::InvalidData(format!(
1242 "attachment role {role:?} for record {record_uuid:?} conflicts with an existing substrate or content_ref"
1243 )));
1244 }
1245 Ok(())
1246}
1247
1248fn finalize_attachment_cutover_on_connection(
1249 conn: &Connection,
1250 now: i64,
1251) -> Result<(), SqliteError> {
1252 require_legacy_attachment_fences(conn)?;
1253 validate_canonical_legacy_refs(conn)?;
1254 validate_canonical_attachment_and_claim_refs(conn)?;
1255 validate_attachment_record_owners(conn)?;
1256 validate_legacy_content_backfill(conn)?;
1257
1258 let remaining_claims: i64 =
1259 conn.query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))?;
1260 if remaining_claims != 0 {
1261 return Err(SqliteError::InvalidData(format!(
1262 "attachment cutover cannot finalize while {remaining_claims} blob GC claim rows remain"
1263 )));
1264 }
1265
1266 let uncovered_model: Option<String> = conn
1267 .query_row(
1268 "SELECT model.id FROM entities AS model \
1269 WHERE model.entity_type = 'moodboard_model' \
1270 AND model.content_ref IS NOT NULL \
1271 AND NOT EXISTS ( \
1272 SELECT 1 FROM attachments AS attachment \
1273 WHERE attachment.record_uuid = model.id \
1274 AND attachment.substrate = 'entity' \
1275 AND attachment.role = 'fann-network' \
1276 ) \
1277 LIMIT 1",
1278 [],
1279 |row| row.get(0),
1280 )
1281 .optional()?;
1282 if let Some(record_uuid) = uncovered_model {
1283 return Err(SqliteError::InvalidData(format!(
1284 "moodboard_model {record_uuid:?} has legacy content but no verified 'fann-network' attachment"
1285 )));
1286 }
1287
1288 conn.execute_batch(V21_ATTACHMENT_FENCES_UP)?;
1289 conn.execute_batch(
1290 "DROP TRIGGER entities_reject_claimed_blob_insert; \
1291 DROP TRIGGER entities_reject_claimed_blob_update; \
1292 DROP INDEX idx_entities_content_ref; \
1293 ALTER TABLE entities DROP COLUMN content_ref;",
1294 )?;
1295 conn.execute(
1296 "UPDATE attachment_cutover_state \
1297 SET state = 'complete', completed_at = ?1 \
1298 WHERE singleton = 1 AND state = 'incomplete'",
1299 [now],
1300 )?;
1301 Ok(())
1302}
1303
1304fn record_attachment_cutover_migration(conn: &Connection, now: i64) -> Result<(), SqliteError> {
1305 let migration = MIGRATIONS
1306 .iter()
1307 .find(|migration| migration.version == ATTACHMENT_CUTOVER_VERSION)
1308 .expect("V21 migration must be registered");
1309 conn.execute(
1310 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3)",
1311 rusqlite::params![migration.version, migration.name, now],
1312 )?;
1313 Ok(())
1314}
1315
1316pub fn finalize_attachment_cutover(conn: &mut Connection) -> Result<(), SqliteError> {
1323 let admission = WriteAdmission::for_canonical_path(canonical_connection_database_path(conn)?)?;
1324 RawMigrationWriteUnit::new(conn, &admission)?
1325 .run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
1326}
1327
1328pub fn finalize_attachment_cutover_with_policy(
1330 conn: &mut Connection,
1331 policy: &MigrationWritePolicy,
1332) -> Result<(), SqliteError> {
1333 let admission =
1334 WriteAdmission::for_migration_policy(canonical_connection_database_path(conn)?, policy)?;
1335 RawMigrationWriteUnit::new(conn, &admission)?
1336 .run(|conn| finalize_attachment_cutover_with_admission(conn, &admission))
1337}
1338
1339pub(crate) fn finalize_attachment_cutover_with_admission(
1341 conn: &mut Connection,
1342 admission: &WriteAdmission,
1343) -> Result<(), SqliteError> {
1344 if attachment_cutover_status(conn)? == AttachmentCutoverStatus::Complete {
1345 return Ok(());
1346 }
1347 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Exclusive)?;
1348 if let Err(error) = admission.check() {
1349 let rollback = tx.rollback();
1350 return Err(capacity_refusal_after_rollback(
1351 conn,
1352 rollback,
1353 error,
1354 "attachment cutover finalization",
1355 ));
1356 }
1357 match attachment_cutover_status(&tx)? {
1358 AttachmentCutoverStatus::Complete => return Ok(()),
1359 AttachmentCutoverStatus::Pending => {
1360 return Err(SqliteError::InvalidData(
1361 "attachment cutover must complete stage 1 before finalization".into(),
1362 ))
1363 }
1364 AttachmentCutoverStatus::Incomplete => {}
1365 }
1366 let now = chrono::Utc::now().timestamp_micros();
1367 finalize_attachment_cutover_on_connection(&tx, now)?;
1368 record_attachment_cutover_migration(&tx, now)?;
1369 tx.commit()?;
1370 Ok(())
1371}
1372
1373fn read_applied_migration_ledger(
1375 conn: &Connection,
1376 through_version: u32,
1377) -> Result<Vec<(u32, String)>, SqliteError> {
1378 let mut stmt = conn.prepare(
1379 "SELECT version, name FROM _schema_migrations \
1380 WHERE version <= ?1 ORDER BY version ASC",
1381 )?;
1382 let rows = stmt
1383 .query_map([through_version], |row| {
1384 Ok((row.get::<_, u32>(0)?, row.get::<_, String>(1)?))
1385 })?
1386 .collect::<Result<Vec<_>, _>>()?;
1387 Ok(rows)
1388}
1389
1390fn validate_applied_migration_versions(
1395 applied: &[(u32, String)],
1396 through_version: u32,
1397) -> Result<(), SqliteError> {
1398 let expected: Vec<&VersionedMigration> = MIGRATIONS
1399 .iter()
1400 .filter(|migration| migration.version <= through_version)
1401 .collect();
1402 let mut applied_index = 0;
1403
1404 for migration in expected {
1405 let Some((version, applied_name)) = applied.get(applied_index) else {
1406 return Err(SqliteError::InvalidData(format!(
1407 "migration history is missing version {} ('{}'); the applied ledger must be \
1408 the exact contiguous canonical sequence through version {through_version}",
1409 migration.version, migration.name,
1410 )));
1411 };
1412 if *version < migration.version {
1413 return Err(SqliteError::InvalidData(format!(
1414 "migration history contains unknown version {version} recorded as \
1415 '{applied_name}'; the applied ledger must contain only canonical versions"
1416 )));
1417 }
1418 if *version > migration.version {
1419 return Err(SqliteError::InvalidData(format!(
1420 "migration history is missing version {} ('{}'); found version {version} \
1421 next instead",
1422 migration.version, migration.name,
1423 )));
1424 }
1425 applied_index += 1;
1426 }
1427
1428 if let Some((version, name)) = applied.get(applied_index) {
1429 return Err(SqliteError::InvalidData(format!(
1430 "migration history contains unknown version {version} recorded as '{name}'; \
1431 the applied ledger must contain only canonical versions"
1432 )));
1433 }
1434
1435 Ok(())
1436}
1437
1438fn validate_applied_migration_names(
1439 applied: &[(u32, String)],
1440 through_version: u32,
1441 allow_known_v19_repairs: bool,
1442) -> Result<(), SqliteError> {
1443 for ((version, applied_name), migration) in applied.iter().zip(
1444 MIGRATIONS
1445 .iter()
1446 .filter(|migration| migration.version <= through_version),
1447 ) {
1448 debug_assert_eq!(*version, migration.version);
1449 if migration.name != applied_name.as_str() {
1450 if allow_known_v19_repairs && matches!(*version, 13 | 14) {
1451 continue;
1452 }
1453 return Err(SqliteError::InvalidData(format!(
1454 "migration version {version} is recorded under name '{applied_name}', \
1455 expected '{expected}'. This database's migration history does not match \
1456 the current binary; recreate it from the current schema or repair the \
1457 specific known divergence via a dedicated migration.",
1458 expected = migration.name,
1459 )));
1460 }
1461 }
1462
1463 Ok(())
1464}
1465
1466fn validate_applied_migration_ledger(
1470 conn: &Connection,
1471 through_version: u32,
1472) -> Result<(), SqliteError> {
1473 let applied = read_applied_migration_ledger(conn, through_version)?;
1474 validate_applied_migration_versions(&applied, through_version)?;
1475 validate_applied_migration_names(&applied, through_version, false)
1476}
1477
1478const MIGRATION_TRACKING_TABLE: &str = include_str!("../sql/schema-migrations-table.sql");
1479
1480pub fn read_schema_version(conn: &Connection) -> Result<u32, SqliteError> {
1486 match conn.query_row(
1487 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1488 [],
1489 |row| row.get(0),
1490 ) {
1491 Ok(version) => Ok(version),
1492 Err(rusqlite::Error::SqliteFailure(_, Some(ref msg)))
1493 if msg.contains("no such table: _schema_migrations") =>
1494 {
1495 Ok(0)
1496 }
1497 Err(e) => Err(e.into()),
1498 }
1499}
1500
1501pub fn inspect_schema_version(path: &std::path::Path) -> Result<u32, SqliteError> {
1506 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1507 read_schema_version(&conn)
1508}
1509
1510pub fn inspect_schema_is_current(path: &std::path::Path) -> Result<u32, SqliteError> {
1513 let conn = crate::pool::open_read_only_snapshot_connection(path)?;
1514 validate_schema_is_current(&conn)
1515}
1516
1517pub fn validate_schema_is_current(conn: &Connection) -> Result<u32, SqliteError> {
1525 let current_version = read_schema_version(conn)?;
1526 let latest_version = latest_schema_version();
1527
1528 if current_version < latest_version {
1529 return Err(SqliteError::InvalidData(format!(
1530 "read-only database schema version {current_version} is behind the latest known \
1531 migration {latest_version}; migrate a writable copy with this build before opening \
1532 the snapshot read-only"
1533 )));
1534 }
1535 if current_version > latest_version {
1536 return Err(SqliteError::InvalidData(format!(
1537 "read-only database schema version {current_version} is ahead of the latest known \
1538 migration {latest_version}; use a compatible newer build or recreate the snapshot"
1539 )));
1540 }
1541
1542 validate_applied_migration_ledger(conn, current_version)?;
1549 if current_version >= ATTACHMENT_CUTOVER_VERSION
1553 && attachment_cutover_status(conn)? != AttachmentCutoverStatus::Complete
1554 {
1555 return Err(SqliteError::InvalidData(
1556 "read-only database has not completed the V21 attachment cutover".into(),
1557 ));
1558 }
1559
1560 if current_version >= MEMORY_VISIBILITY_CUTOVER_VERSION {
1561 memory_visibility::validate_cutover(conn)?;
1562 }
1563
1564 Ok(current_version)
1565}
1566
1567pub fn validate_memory_visibility_cutover(conn: &Connection) -> Result<(), SqliteError> {
1570 validate_schema_is_current(conn).map(|_| ())
1571}
1572
1573#[cfg(test)]
1574pub(crate) mod test_sync {
1575 use std::sync::atomic::AtomicU32;
1576 use std::sync::{Arc, Barrier, Mutex};
1577
1578 pub(crate) static STALE_READ_BARRIER: Mutex<Option<Arc<Barrier>>> = Mutex::new(None);
1582 pub(crate) static LOCKED_FAST_FORWARDS: AtomicU32 = AtomicU32::new(0);
1584 pub(crate) static BUSY_OBSERVED: std::sync::atomic::AtomicBool =
1588 std::sync::atomic::AtomicBool::new(false);
1589
1590 pub(crate) fn record_busy(_count: i32) -> bool {
1593 BUSY_OBSERVED.store(true, std::sync::atomic::Ordering::SeqCst);
1594 std::thread::sleep(std::time::Duration::from_millis(1));
1595 true
1596 }
1597
1598 pub(crate) static WINNER_COMMITTED: std::sync::atomic::AtomicBool =
1601 std::sync::atomic::AtomicBool::new(false);
1602 pub(crate) static LOSER_SAW_WINNER_COMMIT: std::sync::atomic::AtomicBool =
1607 std::sync::atomic::AtomicBool::new(false);
1608
1609 std::thread_local! {
1610 pub(crate) static PARTICIPATE: std::cell::Cell<bool> =
1613 const { std::cell::Cell::new(false) };
1614 pub(crate) static FIRST_BEGIN_DONE: std::cell::Cell<bool> =
1616 const { std::cell::Cell::new(false) };
1617 }
1618}
1619
1620fn canonical_connection_database_path(conn: &Connection) -> Result<Option<PathBuf>, SqliteError> {
1629 let configured = conn.path().unwrap_or_default();
1630 let raw_path = if configured.is_empty() {
1631 conn.query_row(
1632 "SELECT file FROM pragma_database_list WHERE name = 'main'",
1633 [],
1634 |row| row.get::<_, String>(0),
1635 )?
1636 } else {
1637 configured.to_string()
1638 };
1639
1640 if raw_path.is_empty() {
1641 return Ok(None);
1642 }
1643 let canonical = std::fs::canonicalize(&raw_path).map_err(SqliteError::Io)?;
1644 #[cfg(any(unix, windows))]
1645 crate::pool::opened_sqlite_file_identity(conn, &canonical)?;
1646 Ok(Some(canonical))
1647}
1648
1649fn validate_database_gc_owner(
1650 conn: &Connection,
1651 owner: &DatabaseGcOwnerGuard,
1652) -> Result<(), SqliteError> {
1653 let connection_path = canonical_connection_database_path(conn)?;
1654 if owner.database_path() != connection_path.as_deref() {
1655 return Err(SqliteError::InvalidData(format!(
1656 "database GC owner targets {:?}, but migration connection targets {:?}",
1657 owner.database_path(),
1658 connection_path.as_deref(),
1659 )));
1660 }
1661 Ok(())
1662}
1663
1664pub fn run_migrations(conn: &mut Connection) -> Result<u32, SqliteError> {
1665 let database_path = canonical_connection_database_path(conn)?;
1666 let admission = WriteAdmission::for_canonical_path(database_path.clone())?;
1667 run_raw_migrations_with_admission(conn, database_path, &admission)
1668}
1669
1670pub fn run_migrations_with_policy(
1672 conn: &mut Connection,
1673 policy: &MigrationWritePolicy,
1674) -> Result<u32, SqliteError> {
1675 let database_path = canonical_connection_database_path(conn)?;
1676 let admission = WriteAdmission::for_migration_policy(database_path.clone(), policy)?;
1677 run_raw_migrations_with_admission(conn, database_path, &admission)
1678}
1679
1680fn run_raw_migrations_with_admission(
1681 conn: &mut Connection,
1682 database_path: Option<PathBuf>,
1683 admission: &WriteAdmission,
1684) -> Result<u32, SqliteError> {
1685 if let Some(database_path) = database_path {
1686 let owner = try_acquire_database_gc_owner_for_path(database_path).map_err(|error| {
1691 SqliteError::InvalidData(format!(
1692 "failed to acquire database GC owner before schema migration: {error}"
1693 ))
1694 })?;
1695 return run_migrations_with_database_gc_owner(
1696 &mut RawMigrationTransactions::new(conn, admission),
1697 &owner,
1698 admission,
1699 );
1700 }
1701
1702 run_versioned_migrations(
1705 &mut RawMigrationTransactions::new(conn, admission),
1706 None,
1707 admission,
1708 )
1709}
1710
1711pub(crate) fn run_migrations_with_database_gc_owner(
1712 writes: &mut impl MigrationTransactions,
1713 owner: &DatabaseGcOwnerGuard,
1714 admission: &WriteAdmission,
1715) -> Result<u32, SqliteError> {
1716 run_versioned_migrations(writes, Some(owner), admission)
1717}
1718
1719fn with_migration_busy_timeout<T>(
1725 conn: &mut Connection,
1726 operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
1727) -> Result<T, SqliteError> {
1728 let prior_busy_ms: i64 = conn.query_row("PRAGMA busy_timeout", [], |row| row.get(0))?;
1729 let raised = prior_busy_ms < 5_000;
1730 if raised {
1731 conn.busy_timeout(std::time::Duration::from_secs(5))?;
1732 }
1733 let result = operation(conn);
1734 if raised {
1735 let _ = conn.busy_timeout(std::time::Duration::from_millis(prior_busy_ms.max(0) as u64));
1736 }
1737 result
1738}
1739
1740enum MigrationStep {
1742 Applied,
1744 AppliedBySibling(u32),
1747 Deferred,
1749}
1750
1751fn run_versioned_migrations(
1758 writes: &mut impl MigrationTransactions,
1759 owner: Option<&DatabaseGcOwnerGuard>,
1760 admission: &WriteAdmission,
1761) -> Result<u32, SqliteError> {
1762 let current_version = writes.admitted(|conn| {
1763 if let Some(owner) = owner {
1764 validate_database_gc_owner(conn, owner)?;
1765 }
1766 with_migration_busy_timeout(conn, |conn| bootstrap_migration_ledger(conn, admission))
1767 })?;
1768
1769 #[cfg(test)]
1774 if test_sync::PARTICIPATE.with(|p| p.get()) {
1775 let barrier = test_sync::STALE_READ_BARRIER.lock().unwrap().clone();
1776 if let Some(barrier) = barrier {
1777 barrier.wait();
1778 }
1779 }
1780 let latest_version = latest_schema_version();
1781
1782 let mut applied_version = current_version;
1783 let mut skip_through = current_version;
1787
1788 for migration in MIGRATIONS {
1789 if migration.version <= skip_through {
1790 applied_version = applied_version.max(migration.version);
1791 continue;
1792 }
1793 let step = writes.admitted(|conn| {
1794 with_migration_busy_timeout(conn, |conn| {
1795 apply_versioned_migration(conn, migration, admission, latest_version)
1796 })
1797 })?;
1798 match step {
1799 MigrationStep::Applied => applied_version = migration.version,
1800 MigrationStep::AppliedBySibling(sibling_version) => {
1801 skip_through = sibling_version.min(latest_version);
1802 applied_version = applied_version.max(migration.version);
1803 }
1804 MigrationStep::Deferred => break,
1805 }
1806 }
1807
1808 writes.admitted(|conn| validate_applied_migration_ledger(conn, applied_version))?;
1812
1813 Ok(applied_version)
1814}
1815
1816fn bootstrap_migration_ledger(
1820 conn: &mut Connection,
1821 admission: &WriteAdmission,
1822) -> Result<u32, SqliteError> {
1823 admission.check()?;
1824 conn.execute_batch(MIGRATION_TRACKING_TABLE)?;
1825 require_autocommit(conn, "migration-tracking bootstrap")?;
1826
1827 let current_version: u32 = read_schema_version(conn)?;
1828
1829 let latest_version = latest_schema_version();
1835 if current_version > latest_version {
1836 return Err(SqliteError::InvalidData(format!(
1837 "database schema version {current_version} is ahead of the latest known migration \
1838 {latest_version}. This database predates the consolidated baseline (ADR-015) or was \
1839 written by a newer build. Recreate it from the current schema; in-place downgrade is \
1840 not supported."
1841 )));
1842 }
1843
1844 let applied = read_applied_migration_ledger(conn, current_version)?;
1851 validate_applied_migration_versions(&applied, current_version)?;
1852 validate_applied_migration_names(&applied, current_version, current_version < 19)?;
1853
1854 Ok(current_version)
1855}
1856
1857fn apply_versioned_migration(
1860 conn: &mut Connection,
1861 migration: &VersionedMigration,
1862 admission: &WriteAdmission,
1863 latest_version: u32,
1864) -> Result<MigrationStep, SqliteError> {
1865 #[cfg(test)]
1869 let instrumented_first_begin =
1870 test_sync::PARTICIPATE.with(|p| p.get()) && !test_sync::FIRST_BEGIN_DONE.with(|f| f.get());
1871 #[cfg(test)]
1872 if instrumented_first_begin {
1873 test_sync::FIRST_BEGIN_DONE.with(|f| f.set(true));
1874 }
1875 #[cfg(test)]
1878 if test_sync::PARTICIPATE.with(|p| p.get()) {
1879 conn.busy_handler(Some(test_sync::record_busy))?;
1880 }
1881 let tx = conn
1882 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
1883 .map_err(|e| SqliteError::Migration {
1884 version: migration.version,
1885 error: e.to_string(),
1886 })?;
1887 if let Err(error) = admission.check() {
1888 let rollback = tx.rollback();
1889 return Err(capacity_refusal_after_rollback(
1890 conn,
1891 rollback,
1892 error,
1893 "core schema migration",
1894 ));
1895 }
1896
1897 let sibling_version: u32 = tx
1901 .query_row(
1902 "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations",
1903 [],
1904 |row| row.get(0),
1905 )
1906 .map_err(|e| SqliteError::Migration {
1907 version: migration.version,
1908 error: e.to_string(),
1909 })?;
1910 #[cfg(test)]
1911 if instrumented_first_begin {
1912 use std::sync::atomic::Ordering::SeqCst;
1913 if sibling_version == 0 {
1914 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1920 while !test_sync::BUSY_OBSERVED.load(SeqCst) && std::time::Instant::now() < deadline {
1921 std::thread::yield_now();
1922 }
1923 } else {
1924 test_sync::LOSER_SAW_WINNER_COMMIT
1928 .store(test_sync::WINNER_COMMITTED.load(SeqCst), SeqCst);
1929 }
1930 }
1931
1932 if sibling_version > latest_version {
1937 return Err(SqliteError::InvalidData(format!(
1938 "database schema version {sibling_version} is ahead of the latest known \
1939 migration {latest_version} (committed by a concurrent process while this \
1940 one waited for the migration write lock). This build cannot run against \
1941 the newer schema; upgrade the binary or recreate the database."
1942 )));
1943 }
1944
1945 if sibling_version >= migration.version {
1946 #[cfg(test)]
1947 test_sync::LOCKED_FAST_FORWARDS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1948 return Ok(MigrationStep::AppliedBySibling(sibling_version));
1949 }
1950
1951 if migration.version == 46 {
1952 memory_visibility::capture_pre_v46(&tx).map_err(|error| SqliteError::Migration {
1953 version: migration.version,
1954 error: error.to_string(),
1955 })?;
1956 #[cfg(test)]
1957 memory_visibility::test_state::stop_at(memory_visibility::test_state::Stop::AfterCapture)?;
1958 }
1959
1960 if migration.version == ATTACHMENT_CUTOVER_VERSION {
1961 let status = attachment_cutover_status(&tx).map_err(|e| SqliteError::Migration {
1962 version: migration.version,
1963 error: e.to_string(),
1964 })?;
1965 let legacy_refs: i64 = tx
1966 .query_row(
1967 "SELECT COUNT(*) FROM entities WHERE content_ref IS NOT NULL",
1968 [],
1969 |row| row.get(0),
1970 )
1971 .map_err(|e| SqliteError::Migration {
1972 version: migration.version,
1973 error: e.to_string(),
1974 })?;
1975
1976 if status == AttachmentCutoverStatus::Incomplete || legacy_refs != 0 {
1980 drop(tx);
1981 return Ok(MigrationStep::Deferred);
1982 }
1983 if status != AttachmentCutoverStatus::Pending {
1984 return Err(SqliteError::Migration {
1985 version: migration.version,
1986 error: format!("unexpected attachment cutover state {status:?}"),
1987 });
1988 }
1989
1990 let now = chrono::Utc::now().timestamp_micros();
1991 stage_attachment_cutover_on_connection(&tx, now).map_err(|e| SqliteError::Migration {
1992 version: migration.version,
1993 error: e.to_string(),
1994 })?;
1995 finalize_attachment_cutover_on_connection(&tx, now).map_err(|e| {
1996 SqliteError::Migration {
1997 version: migration.version,
1998 error: e.to_string(),
1999 }
2000 })?;
2001 } else if migration.version == 36 {
2002 crate::stores::event::ensure_operation_attribution_columns(&tx).map_err(|error| {
2004 SqliteError::Migration {
2005 version: migration.version,
2006 error: error.to_string(),
2007 }
2008 })?;
2009 } else if migration.version == 44 {
2010 migrate_outbound_due_key(&tx).map_err(|error| SqliteError::Migration {
2011 version: migration.version,
2012 error: error.to_string(),
2013 })?;
2014 } else if migration.version == 48 {
2015 migrate_acknowledgement_journal(&tx).map_err(|error| SqliteError::Migration {
2016 version: migration.version,
2017 error: error.to_string(),
2018 })?;
2019 } else if migration.name == SESSION_IDENTITY_MIGRATION_NAME {
2020 tx.execute_batch(migration.up)
2021 .map_err(|error| SqliteError::Migration {
2022 version: migration.version,
2023 error: error.to_string(),
2024 })?;
2025 session_identity_migration::apply(&tx).map_err(|error| SqliteError::Migration {
2026 version: migration.version,
2027 error: error.to_string(),
2028 })?;
2029 } else {
2030 tx.execute_batch(migration.up)
2031 .map_err(|e| SqliteError::Migration {
2032 version: migration.version,
2033 error: e.to_string(),
2034 })?;
2035 }
2036
2037 let visibility_counts = if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
2038 Some(
2039 memory_visibility::cutover_counts(&tx).map_err(|error| SqliteError::Migration {
2040 version: migration.version,
2041 error: error.to_string(),
2042 })?,
2043 )
2044 } else {
2045 None
2046 };
2047
2048 if migration.version == 19 {
2055 tx.execute_batch(
2056 "UPDATE _schema_migrations SET name = 'list_cursor_sequences' WHERE version = 13;\n\
2057 UPDATE _schema_migrations SET name = 'graph_edges_id_unique' WHERE version = 14;",
2058 )
2059 .map_err(|e| SqliteError::Migration {
2060 version: migration.version,
2061 error: e.to_string(),
2062 })?;
2063 }
2064
2065 let now = chrono::Utc::now().timestamp_micros();
2066 tx.execute(
2067 "INSERT INTO _schema_migrations (version, name, applied_at) VALUES (?1, ?2, ?3) \
2068 ON CONFLICT(version) DO NOTHING",
2069 rusqlite::params![migration.version, migration.name, now],
2070 )
2071 .map_err(|e| SqliteError::Migration {
2072 version: migration.version,
2073 error: e.to_string(),
2074 })?;
2075
2076 #[cfg(test)]
2077 if instrumented_first_begin {
2078 test_sync::WINNER_COMMITTED.store(true, std::sync::atomic::Ordering::SeqCst);
2079 }
2080
2081 #[cfg(test)]
2082 if migration.version == MEMORY_VISIBILITY_CUTOVER_VERSION {
2083 memory_visibility::test_state::stop_at(
2084 memory_visibility::test_state::Stop::BeforeCutoverCommit,
2085 )?;
2086 }
2087 tx.commit().map_err(|e| SqliteError::Migration {
2088 version: migration.version,
2089 error: e.to_string(),
2090 })?;
2091 if let Some(counts) = visibility_counts {
2092 memory_visibility::log_counts(&counts, conn.path().unwrap_or(":memory:"));
2093 }
2094 #[cfg(test)]
2095 if migration.version == 46 {
2096 memory_visibility::test_state::stop_at(
2097 memory_visibility::test_state::Stop::AfterV46Commit,
2098 )?;
2099 }
2100
2101 Ok(MigrationStep::Applied)
2102}
2103
2104#[derive(Debug)]
2105pub struct EmbeddingModelRegistryRecord {
2106 pub engine_name: String,
2108 pub model_id: String,
2110 pub key_version: String,
2112 pub dimensions: u32,
2114 pub status: String,
2116 pub activated_at: Option<i64>,
2118 pub superseded_at: Option<i64>,
2120}
2121
2122pub fn query_embedding_models(
2128 db: Option<&std::path::Path>,
2129 engine_filter: Option<&str>,
2130) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
2131 let path = db.map(std::path::Path::to_path_buf).unwrap_or_else(|| {
2132 std::env::var("HOME")
2133 .map(std::path::PathBuf::from)
2134 .unwrap_or_else(|_| std::path::PathBuf::from("."))
2135 .join(".khive/khive.db")
2136 });
2137 if !path.exists() {
2138 return Ok(Vec::new());
2139 }
2140 let conn = Connection::open_with_flags(
2141 path,
2142 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
2143 | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
2144 | rusqlite::OpenFlags::SQLITE_OPEN_URI,
2145 )?;
2146 query_embedding_models_conn(&conn, engine_filter)
2147}
2148
2149pub(crate) fn query_embedding_models_conn(
2153 conn: &Connection,
2154 engine_filter: Option<&str>,
2155) -> Result<Vec<EmbeddingModelRegistryRecord>, SqliteError> {
2156 let exists: bool = conn.query_row(
2157 "SELECT COUNT(*) > 0 FROM sqlite_master \
2158 WHERE type='table' AND name='_embedding_models'",
2159 [],
2160 |row| row.get(0),
2161 )?;
2162 if !exists {
2163 return Ok(Vec::new());
2164 }
2165
2166 let sql = if engine_filter.is_some() {
2167 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
2168 FROM _embedding_models WHERE engine_name = ?1 \
2169 ORDER BY engine_name, activated_at IS NULL, activated_at"
2170 } else {
2171 "SELECT engine_name, model_id, key_version, dim, status, activated_at, superseded_at \
2172 FROM _embedding_models \
2173 ORDER BY engine_name, activated_at IS NULL, activated_at"
2174 };
2175 let mut stmt = conn.prepare(sql)?;
2176 let map_row = |row: &rusqlite::Row<'_>| {
2177 let dim_raw: i64 = row.get(3)?;
2178 let dimensions = u32::try_from(dim_raw).map_err(|_| {
2179 rusqlite::Error::FromSqlConversionFailure(
2180 3,
2181 rusqlite::types::Type::Integer,
2182 Box::new(std::io::Error::other(format!(
2183 "_embedding_models.dim value {dim_raw} is outside the valid u32 range [0, {}]",
2184 u32::MAX,
2185 ))),
2186 )
2187 })?;
2188 Ok(EmbeddingModelRegistryRecord {
2189 engine_name: row.get(0)?,
2190 model_id: row.get(1)?,
2191 key_version: row.get(2)?,
2192 dimensions,
2193 status: row.get(4)?,
2194 activated_at: row.get(5)?,
2195 superseded_at: row.get(6)?,
2196 })
2197 };
2198
2199 if let Some(engine) = engine_filter {
2200 stmt.query_map([engine], map_row)?
2201 .collect::<Result<Vec<_>, _>>()
2202 .map_err(Into::into)
2203 } else {
2204 stmt.query_map([], map_row)?
2205 .collect::<Result<Vec<_>, _>>()
2206 .map_err(Into::into)
2207 }
2208}
2209
2210#[cfg(test)]
2212pub(crate) fn migration_test_policy() -> MigrationWritePolicy {
2213 MigrationWritePolicy::new(
2214 crate::DiskGuardEnvironment::capture()
2215 .resolve(None, None)
2216 .expect("test disk policy"),
2217 crate::PoolConfig::for_test()
2218 .volume_lock_dir
2219 .expect("test volume-lock directory"),
2220 )
2221 .expect("valid test migration policy")
2222}
2223
2224#[cfg(test)]
2225pub(crate) fn run_migrations_for_test(conn: &mut Connection) -> Result<u32, SqliteError> {
2226 run_migrations_with_policy(conn, &migration_test_policy())
2227}
2228
2229#[cfg(test)]
2230fn apply_schema_plan_for_test(
2231 conn: &mut Connection,
2232 plan: &ServiceSchemaPlan,
2233) -> Result<(), SqliteError> {
2234 apply_schema_plan_with_policy(conn, plan, &migration_test_policy())
2235}
2236
2237#[cfg(test)]
2238fn stage_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
2239 stage_attachment_cutover_with_policy(conn, &migration_test_policy())
2240}
2241
2242#[cfg(test)]
2243fn finalize_attachment_cutover_for_test(conn: &mut Connection) -> Result<(), SqliteError> {
2244 finalize_attachment_cutover_with_policy(conn, &migration_test_policy())
2245}
2246
2247#[cfg(test)]
2252#[path = "entity_version_migration_measurement.rs"]
2253mod entity_version_measurement;
2254
2255#[cfg(test)]
2256#[path = "migrations_tests.rs"]
2257mod tests;
2258
2259#[cfg(test)]
2260#[path = "raw_migration_settlement_tests.rs"]
2261mod raw_settlement_tests;
2262
2263#[cfg(test)]
2264#[path = "git_note_index_migration_tests.rs"]
2265mod git_note_indexes;
2266
2267#[cfg(test)]
2268#[path = "entity_list_index_migration_tests.rs"]
2269mod entity_list_indexes;
2270
2271#[cfg(test)]
2272#[path = "schedule_core_index_migration_tests.rs"]
2273mod schedule_core_index_migration_tests;
2274
2275#[cfg(test)]
2276#[path = "comm_core_index_migration_tests.rs"]
2277mod comm_core_index_migration_tests;