1pub mod lifecycle;
2mod pcache2;
3
4use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
5use std::error::Error;
6use std::fmt::{Display, Formatter};
7use std::fs::{File, OpenOptions};
8use std::io::{Seek, SeekFrom, Write};
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
11use std::sync::mpsc::{self, Receiver, SyncSender};
12use std::sync::Once;
13use std::sync::{Arc, Condvar, Mutex};
14use std::thread::{self, JoinHandle};
15use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
16
17use fathomdb_embedder::EmbedderEvent;
18#[cfg(feature = "operator")]
20use fathomdb_embedder::MeanRecomputeTrigger;
21use fathomdb_embedder_api::{Embedder, EmbedderError as RuntimeEmbedderError, EmbedderIdentity};
22use fathomdb_query::compile_text_query;
23use fathomdb_schema::{
24 migrate_with_event_sink, MigrationError as SchemaMigrationError, MigrationStepReport,
25 LOCK_SUFFIX, MIGRATIONS, SCHEMA_VERSION,
26};
27#[cfg(feature = "operator")]
29use fathomdb_schema::CANONICAL_TABLES;
30use jsonschema::JSONSchema;
31use rusqlite::{params, Connection, OptionalExtension};
32use serde_json::Value;
33#[cfg(feature = "operator")]
35use sha2::Digest;
36use sqlite_vec::sqlite3_vec_init;
37
38#[cfg(unix)]
39use std::os::unix::fs::OpenOptionsExt;
40
41const DEFAULT_EMBEDDER_NAME: &str = "fathomdb-bge-small-en-v1.5";
47const DEFAULT_EMBEDDER_REVISION: &str = "5c38ec7c405ec4b44b94cc5a9bb96e735b38267a";
48const DEFAULT_EMBEDDER_DIMENSION: u32 = 384;
49
50const BGE_SMALL_EMBEDDER_NAME: &str = "fathomdb-bge-small-en-v1.5";
60
61const DEFAULT_SLOW_THRESHOLD_MS: u64 = 100;
64const DEFAULT_VECTOR_PROFILE: &str = "default";
65const DEFAULT_VECTOR_PARTITION: &str = "vector_default";
66#[cfg(feature = "operator")]
72const REBUILD_DRAIN_TIMEOUT_MS: u64 = 30_000;
73const SEARCH_INDEX_TOKENIZER_SCHEMA_VERSION: u32 = 11;
81const SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY: &str =
92 "search_index_tokenizer_reproject_complete";
93const DEFAULT_PROVENANCE_ROW_CAP: u64 = 1_000_000;
94const PROJECTION_CURSOR_KEY: &str = "projection_cursor";
95const PROJECTION_WORKERS: usize = 2;
96const DEFAULT_EMBED_TIMEOUT_MS: u64 = 30_000;
102const DEFAULT_EMBED_CIRCUIT_THRESHOLD: u64 = 8;
110const PROJECTION_COMMIT_BATCH: usize = 16;
111const PROJECTION_INFLIGHT_LIMIT: usize = PROJECTION_WORKERS * PROJECTION_COMMIT_BATCH;
115const PROJECTION_SCAN_FETCH: usize = PROJECTION_INFLIGHT_LIMIT;
118const DEFAULT_PROJECTION_RETRY_DELAYS_MS: [u64; 3] = [1_000, 4_000, 16_000];
119
120const READER_POOL_SIZE: usize = 8;
124
125const READER_LOOKASIDE_SLOT_SIZE: std::os::raw::c_int = 1200;
133
134const READER_LOOKASIDE_SLOT_COUNT: std::os::raw::c_int = 500;
139
140pub struct Engine {
141 path: PathBuf,
142 next_cursor: AtomicU64,
143 closed: AtomicBool,
144 lock: Mutex<Option<File>>,
145 connection: Mutex<Option<Connection>>,
146 reader_pool: ReaderWorkerPool,
147 counters: lifecycle::Counters,
148 subscribers: Arc<lifecycle::SubscriberRegistry>,
149 profiling_enabled: Arc<AtomicBool>,
150 slow_threshold_ms: Arc<AtomicU64>,
151 runtime_embedder: Option<Arc<dyn Embedder>>,
152 runtime_embedder_identity: EmbedderIdentity,
153 projection_runtime: ProjectionRuntime,
154 provenance_row_cap: AtomicU64,
155 #[allow(clippy::vec_box)]
167 profile_contexts: Mutex<Vec<Box<ProfileContext>>>,
168 #[allow(dead_code)]
175 reader_lookaside_rcs: Vec<i32>,
176 #[cfg(debug_assertions)]
177 force_next_commit_failure: AtomicBool,
178}
179
180#[derive(Clone, Debug)]
181struct ProjectionJob {
182 cursor: u64,
183 kind: String,
184 body: String,
185}
186
187#[derive(Debug, Default)]
188struct ProjectionRuntimeState {
189 active_jobs: usize,
190 queued_jobs: usize,
191 frozen: bool,
192 pending_scan: bool,
193 stopping: bool,
194 in_flight: BTreeSet<u64>,
195}
196
197struct ProjectionRuntimeShared {
198 path: PathBuf,
199 embedder: Option<Arc<dyn Embedder>>,
200 embedder_identity: EmbedderIdentity,
201 state: Mutex<ProjectionRuntimeState>,
202 state_cvar: Condvar,
203 queue: Mutex<VecDeque<ProjectionJob>>,
204 queue_cvar: Condvar,
205 retry_delays_ms: Mutex<Vec<u64>>,
206 embed_timeout_ms: AtomicU64,
214 embed_serialize: Mutex<()>,
256 live_embed_threads: Arc<AtomicU64>,
272 embed_circuit_open: AtomicBool,
273 embed_circuit_threshold: AtomicU64,
274 mean_accumulator: Mutex<Option<MeanAccumulator>>,
280 pending_events: Mutex<Vec<EmbedderEvent>>,
285 commit_gate: Mutex<()>,
293 search_limit_override: AtomicUsize,
301 recency_reweight_enabled: AtomicBool,
306 vector_stage_only_for_test: AtomicBool,
318 #[cfg(debug_assertions)]
323 force_recompute_failure: AtomicBool,
324}
325
326impl std::fmt::Debug for ProjectionRuntimeShared {
327 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
328 f.debug_struct("ProjectionRuntimeShared")
329 .field("path", &self.path)
330 .field("embedder_identity", &self.embedder_identity)
331 .finish_non_exhaustive()
332 }
333}
334
335#[derive(Debug)]
336struct ProjectionRuntime {
337 shared: Arc<ProjectionRuntimeShared>,
338 dispatcher: Mutex<Option<JoinHandle<()>>>,
339 workers: Mutex<Vec<JoinHandle<()>>>,
340}
341
342impl std::fmt::Debug for Engine {
343 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
344 f.debug_struct("Engine")
345 .field("path", &self.path)
346 .field("closed", &self.closed.load(Ordering::SeqCst))
347 .field("runtime_embedder_identity", &self.runtime_embedder_identity)
348 .finish_non_exhaustive()
349 }
350}
351
352#[derive(Debug)]
361struct ProfileContext {
362 subscribers: Arc<lifecycle::SubscriberRegistry>,
363 profiling_enabled: Arc<AtomicBool>,
364 slow_threshold_ms: Arc<AtomicU64>,
365}
366
367struct ReaderWorkerPool {
376 senders: Vec<SyncSender<ReaderRequest>>,
377 handles: Mutex<Option<Vec<JoinHandle<()>>>>,
378 next: AtomicUsize,
379 shutdown: AtomicBool,
380 live_workers: Arc<AtomicUsize>,
381}
382
383impl std::fmt::Debug for ReaderWorkerPool {
384 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
385 f.debug_struct("ReaderWorkerPool")
386 .field("worker_count", &self.senders.len())
387 .field("live_workers", &self.live_workers.load(Ordering::Relaxed))
388 .field("shutdown", &self.shutdown.load(Ordering::Relaxed))
389 .finish()
390 }
391}
392
393enum ReaderRequest {
397 Search {
398 compiled: fathomdb_query::CompiledQuery,
399 query_vector: Option<String>,
402 query_vector_bin: Option<String>,
406 search_limit: usize,
411 filter: Option<Box<SearchFilter>>,
417 recency_enabled: bool,
420 vector_stage_only: bool,
424 respond: SyncSender<ReaderResponse>,
425 },
426 GetById {
431 logical_ids: Vec<String>,
432 respond: SyncSender<rusqlite::Result<Vec<Option<NodeRecord>>>>,
433 },
434 ReadCollection {
438 collection: String,
439 after_id: Option<i64>,
440 limit: usize,
441 respond: SyncSender<rusqlite::Result<Vec<OpStoreRow>>>,
442 },
443 Shutdown,
444 #[cfg(debug_assertions)]
450 LookasideStatus {
451 respond: SyncSender<i32>,
452 },
453 #[cfg(debug_assertions)]
459 CacheStatus {
460 snapshot_label: String,
461 respond: SyncSender<(String, i32, i32, i32)>,
462 },
463}
464
465type ReaderResponse = rusqlite::Result<(u64, Option<SoftFallback>, Vec<SearchHit>)>;
466
467#[cfg(debug_assertions)]
471#[doc(hidden)]
472#[derive(Clone, Debug)]
473pub struct CacheStatusReply {
474 pub worker_idx: usize,
475 pub snapshot_label: String,
476 pub cache_hit: i32,
477 pub cache_miss: i32,
478 pub cache_used_bytes: i32,
479}
480
481const READER_WORKER_CHANNEL_CAPACITY: usize = 4;
485
486impl ReaderWorkerPool {
487 fn new(connections: Vec<Connection>) -> Self {
488 let live_workers = Arc::new(AtomicUsize::new(0));
489 let mut senders = Vec::with_capacity(connections.len());
490 let mut handles = Vec::with_capacity(connections.len());
491 for (idx, connection) in connections.into_iter().enumerate() {
492 let (tx, rx) = mpsc::sync_channel::<ReaderRequest>(READER_WORKER_CHANNEL_CAPACITY);
493 let live = Arc::clone(&live_workers);
494 let handle = thread::Builder::new()
495 .name(format!("fathomdb-reader-{idx}"))
496 .spawn(move || reader_worker_loop(connection, rx, live))
497 .expect("spawn reader worker");
498 senders.push(tx);
499 handles.push(handle);
500 }
501 Self {
502 senders,
503 handles: Mutex::new(Some(handles)),
504 next: AtomicUsize::new(0),
505 shutdown: AtomicBool::new(false),
506 live_workers,
507 }
508 }
509
510 fn worker_count(&self) -> usize {
511 self.senders.len()
512 }
513
514 fn live_count(&self) -> usize {
515 self.live_workers.load(Ordering::SeqCst)
516 }
517
518 #[cfg(debug_assertions)]
523 fn lookaside_used_per_worker(&self) -> Vec<i32> {
524 let mut results = Vec::with_capacity(self.senders.len());
525 for sender in &self.senders {
526 let (tx, rx) = mpsc::sync_channel::<i32>(1);
527 if sender.send(ReaderRequest::LookasideStatus { respond: tx }).is_ok() {
528 results.push(rx.recv().unwrap_or(-1));
529 } else {
530 results.push(-1);
531 }
532 }
533 results
534 }
535
536 #[cfg(debug_assertions)]
542 fn cache_status_per_worker(&self, snapshot_label: &str) -> Vec<CacheStatusReply> {
543 let mut results = Vec::with_capacity(self.senders.len());
544 for (idx, sender) in self.senders.iter().enumerate() {
545 let (tx, rx) = mpsc::sync_channel::<(String, i32, i32, i32)>(1);
546 let request = ReaderRequest::CacheStatus {
547 snapshot_label: snapshot_label.to_string(),
548 respond: tx,
549 };
550 if sender.send(request).is_ok() {
551 if let Ok((label, hit, miss, used)) = rx.recv() {
552 results.push(CacheStatusReply {
553 worker_idx: idx,
554 snapshot_label: label,
555 cache_hit: hit,
556 cache_miss: miss,
557 cache_used_bytes: used,
558 });
559 continue;
560 }
561 }
562 results.push(CacheStatusReply {
563 worker_idx: idx,
564 snapshot_label: snapshot_label.to_string(),
565 cache_hit: -1,
566 cache_miss: -1,
567 cache_used_bytes: -1,
568 });
569 }
570 results
571 }
572
573 fn dispatch(&self, request: ReaderRequest) -> Result<(), ReaderRequest> {
577 if self.shutdown.load(Ordering::Relaxed) {
578 return Err(request);
579 }
580 let n = self.senders.len();
581 if n == 0 {
582 return Err(request);
583 }
584 let idx = self.next.fetch_add(1, Ordering::Relaxed) % n;
585 self.senders[idx].send(request).map_err(|err| err.0)
586 }
587
588 fn shutdown(&self) {
592 if self.shutdown.swap(true, Ordering::SeqCst) {
593 return;
594 }
595 for sender in &self.senders {
596 let _ = sender.send(ReaderRequest::Shutdown);
597 }
598 if let Ok(mut slot) = self.handles.lock() {
599 if let Some(handles) = slot.take() {
600 for handle in handles {
601 let _ = handle.join();
602 }
603 }
604 }
605 }
606}
607
608impl Drop for ReaderWorkerPool {
609 fn drop(&mut self) {
610 self.shutdown();
611 }
612}
613
614fn reader_worker_loop(
615 mut connection: Connection,
616 rx: Receiver<ReaderRequest>,
617 live_workers: Arc<AtomicUsize>,
618) {
619 live_workers.fetch_add(1, Ordering::SeqCst);
620 struct LiveGuard(Arc<AtomicUsize>);
622 impl Drop for LiveGuard {
623 fn drop(&mut self) {
624 self.0.fetch_sub(1, Ordering::SeqCst);
625 }
626 }
627 let _guard = LiveGuard(live_workers);
628
629 while let Ok(request) = rx.recv() {
630 match request {
631 ReaderRequest::Shutdown => break,
632 ReaderRequest::Search {
633 compiled,
634 query_vector,
635 query_vector_bin,
636 search_limit,
637 filter,
638 recency_enabled,
639 vector_stage_only,
640 respond,
641 } => {
642 let result = read_search_in_tx(
643 &mut connection,
644 &compiled,
645 query_vector.as_deref(),
646 query_vector_bin.as_deref(),
647 search_limit,
648 filter.as_deref(),
649 recency_enabled,
650 vector_stage_only,
651 );
652 let _ = respond.send(result);
655 }
656 ReaderRequest::GetById { logical_ids, respond } => {
657 let result = read_get_by_id_in_tx(&mut connection, &logical_ids);
658 let _ = respond.send(result);
659 }
660 ReaderRequest::ReadCollection { collection, after_id, limit, respond } => {
661 let result = read_collection_in_tx(&mut connection, &collection, after_id, limit);
662 let _ = respond.send(result);
663 }
664 #[cfg(debug_assertions)]
665 ReaderRequest::LookasideStatus { respond } => {
666 let _ = respond.send(read_lookaside_used_hiwtr(&connection));
667 }
668 #[cfg(debug_assertions)]
669 ReaderRequest::CacheStatus { snapshot_label, respond } => {
670 let (hit, miss, used) = read_cache_status(&connection);
671 let _ = respond.send((snapshot_label, hit, miss, used));
672 }
673 }
674 }
675
676 uninstall_profile_callback(&connection);
681 drop(connection);
682}
683
684impl ProjectionRuntime {
685 fn new(
686 path: PathBuf,
687 embedder: Option<Arc<dyn Embedder>>,
688 embedder_identity: EmbedderIdentity,
689 mean_already_pinned: bool,
690 ) -> Self {
691 let mc_required = identity_requires_mean_centering(&embedder_identity);
698 let mean_accumulator = if mc_required && !mean_already_pinned {
699 Some(MeanAccumulator::new(embedder_identity.dimension as usize))
700 } else {
701 None
702 };
703 let shared = Arc::new(ProjectionRuntimeShared {
704 path,
705 embedder,
706 embedder_identity,
707 state: Mutex::new(ProjectionRuntimeState::default()),
708 state_cvar: Condvar::new(),
709 queue: Mutex::new(VecDeque::new()),
710 queue_cvar: Condvar::new(),
711 retry_delays_ms: Mutex::new(DEFAULT_PROJECTION_RETRY_DELAYS_MS.to_vec()),
712 embed_timeout_ms: AtomicU64::new(DEFAULT_EMBED_TIMEOUT_MS),
713 embed_serialize: Mutex::new(()),
714 live_embed_threads: Arc::new(AtomicU64::new(0)),
715 embed_circuit_open: AtomicBool::new(false),
716 embed_circuit_threshold: AtomicU64::new(DEFAULT_EMBED_CIRCUIT_THRESHOLD),
717 mean_accumulator: Mutex::new(mean_accumulator),
718 pending_events: Mutex::new(Vec::new()),
719 commit_gate: Mutex::new(()),
720 search_limit_override: AtomicUsize::new(SEARCH_RERANK_LIMIT),
721 recency_reweight_enabled: AtomicBool::new(false),
722 vector_stage_only_for_test: AtomicBool::new(false),
723 #[cfg(debug_assertions)]
724 force_recompute_failure: AtomicBool::new(false),
725 });
726
727 let dispatcher_shared = Arc::clone(&shared);
728 let dispatcher = thread::spawn(move || projection_dispatcher_loop(dispatcher_shared));
729
730 let mut workers = Vec::with_capacity(PROJECTION_WORKERS);
731 for _ in 0..PROJECTION_WORKERS {
732 let worker_shared = Arc::clone(&shared);
733 workers.push(thread::spawn(move || projection_worker_loop(worker_shared)));
734 }
735
736 Self { shared, dispatcher: Mutex::new(Some(dispatcher)), workers: Mutex::new(workers) }
737 }
738
739 fn notify_new_work(&self) {
740 if let Ok(mut state) = self.shared.state.lock() {
741 state.pending_scan = true;
742 self.shared.state_cvar.notify_all();
743 }
744 }
745
746 fn set_frozen(&self, frozen: bool) {
747 if let Ok(mut state) = self.shared.state.lock() {
748 state.frozen = frozen;
749 if !frozen {
750 state.pending_scan = true;
751 }
752 self.shared.state_cvar.notify_all();
753 }
754 }
755
756 fn wait_for_idle(&self, timeout_ms: u64) -> bool {
757 let deadline = Instant::now() + Duration::from_millis(timeout_ms);
758 let mut state = match self.shared.state.lock() {
759 Ok(state) => state,
760 Err(_) => return false,
761 };
762 loop {
763 if state.active_jobs == 0 && state.queued_jobs == 0 {
764 drop(state);
765 if !database_has_pending_projection_work(&self.shared.path).unwrap_or(true) {
766 return true;
767 }
768 state = match self.shared.state.lock() {
769 Ok(state) => state,
770 Err(_) => return false,
771 };
772 }
773 let now = Instant::now();
774 if now >= deadline {
775 return false;
776 }
777 let wait = deadline.saturating_duration_since(now);
778 let Ok((next_state, _)) = self.shared.state_cvar.wait_timeout(state, wait) else {
779 return false;
780 };
781 state = next_state;
782 }
783 }
784
785 fn set_retry_delays_for_test(&self, delays_ms: &[u64]) {
786 if let Ok(mut delays) = self.shared.retry_delays_ms.lock() {
787 *delays = delays_ms.to_vec();
788 }
789 }
790
791 fn set_embed_timeout_ms_for_test(&self, timeout_ms: u64) {
792 self.shared.embed_timeout_ms.store(timeout_ms, Ordering::Relaxed);
793 }
794
795 fn set_embed_circuit_threshold_for_test(&self, threshold: u64) {
796 self.shared.embed_circuit_threshold.store(threshold, Ordering::Relaxed);
797 }
798
799 fn embed_circuit_open_for_test(&self) -> bool {
800 self.shared.embed_circuit_open.load(Ordering::Relaxed)
801 }
802
803 fn stop(&self) {
804 if let Ok(mut state) = self.shared.state.lock() {
805 if state.stopping {
806 return;
807 }
808 state.stopping = true;
809 state.pending_scan = false;
810 self.shared.state_cvar.notify_all();
811 }
812 if let Ok(mut queue) = self.shared.queue.lock() {
813 queue.clear();
814 self.shared.queue_cvar.notify_all();
815 }
816
817 if let Ok(mut dispatcher) = self.dispatcher.lock() {
818 if let Some(handle) = dispatcher.take() {
819 let _ = handle.join();
820 }
821 }
822 if let Ok(mut workers) = self.workers.lock() {
823 for handle in workers.drain(..) {
824 let _ = handle.join();
825 }
826 }
827 }
828}
829
830#[derive(Clone, Debug, Eq, PartialEq)]
831pub struct OpenReport {
832 pub schema_version_before: u32,
833 pub schema_version_after: u32,
834 pub migration_steps: Vec<MigrationStepReport>,
835 pub embedder_warmup_ms: u64,
836 pub query_backend: &'static str,
837 pub default_embedder: EmbedderIdentity,
838 pub embedder_download_ms: Option<u64>,
852 pub embedder_events: Vec<EmbedderEvent>,
856 pub embedder_mean_centering_required: bool,
863 pub embedder_mean_vec_pinned: bool,
869}
870
871#[derive(Debug)]
872pub struct OpenedEngine {
873 pub engine: Engine,
874 pub report: OpenReport,
875}
876
877#[derive(Clone, Debug)]
880struct LoaderInfo {
881 download_ms: Option<u64>,
882 events: Vec<EmbedderEvent>,
883}
884
885#[derive(Clone, Debug, Eq, PartialEq)]
886pub struct WriteReceipt {
887 pub cursor: u64,
890 pub row_cursors: Vec<u64>,
895 pub dangling_edge_endpoints: u64,
902}
903
904#[derive(Clone, Debug, Eq, PartialEq)]
910pub struct SoftFallback {
911 pub branch: SoftFallbackBranch,
912}
913
914#[derive(Clone, Copy, Debug, Eq, PartialEq)]
920pub enum SoftFallbackBranch {
921 Vector,
922 Text,
923}
924
925#[derive(Clone, Debug, PartialEq)]
941pub struct SearchHit {
942 pub id: u64,
943 pub kind: String,
944 pub body: String,
945 pub score: f64,
946 pub branch: SoftFallbackBranch,
947}
948
949#[derive(Clone, Debug, Eq, PartialEq)]
957pub struct NodeRecord {
958 pub logical_id: String,
959 pub kind: String,
960 pub body: String,
961 pub write_cursor: u64,
962}
963
964#[derive(Clone, Debug, Eq, PartialEq)]
967pub struct OpStoreRow {
968 pub id: i64,
969 pub collection: String,
970 pub record_key: String,
971 pub op_kind: String,
972 pub payload: String,
973 pub schema_id: Option<String>,
974 pub write_cursor: u64,
975}
976
977#[derive(Clone, Debug, PartialEq)]
981pub struct SearchResult {
982 pub projection_cursor: u64,
983 pub soft_fallback: Option<SoftFallback>,
984 pub results: Vec<SearchHit>,
985}
986
987#[derive(Clone, Debug, Default, Eq, PartialEq)]
1003pub struct SearchFilter {
1004 pub source_type: Option<String>,
1005 pub kind: Option<String>,
1006 pub created_after: Option<i64>,
1007 pub status: Option<String>,
1008}
1009
1010impl SearchFilter {
1011 fn is_unfiltered(&self) -> bool {
1015 self.source_type.is_none()
1016 && self.kind.is_none()
1017 && self.created_after.is_none()
1018 && self.status.is_none()
1019 }
1020}
1021
1022#[non_exhaustive]
1028#[derive(Clone, Debug, Eq, PartialEq)]
1029pub enum PreparedWrite {
1030 Node {
1031 kind: String,
1032 body: String,
1033 source_id: Option<String>,
1038 logical_id: Option<String>,
1044 },
1045 Edge {
1046 kind: String,
1047 from: String,
1048 to: String,
1049 source_id: Option<String>,
1051 logical_id: Option<String>,
1054 },
1055 OpStore {
1056 collection: String,
1057 record_key: String,
1058 schema_id: Option<String>,
1059 body: String,
1060 },
1061 AdminSchema {
1062 name: String,
1063 kind: String,
1064 schema_json: String,
1065 retention_json: String,
1066 },
1067}
1068
1069#[derive(Clone, Debug, Default, Eq, PartialEq)]
1075pub struct CounterSnapshot {
1076 pub queries: u64,
1077 pub writes: u64,
1078 pub write_rows: u64,
1079 pub errors_by_code: BTreeMap<String, u64>,
1080 pub admin_ops: u64,
1081 pub cache_hit: u64,
1082 pub cache_miss: u64,
1083}
1084
1085pub use lifecycle::Subscription;
1086
1087#[derive(Clone, Debug, Eq, PartialEq)]
1092pub struct CorruptionDetail {
1093 pub kind: CorruptionKind,
1094 pub stage: OpenStage,
1095 pub locator: CorruptionLocator,
1096 pub recovery_hint: RecoveryHint,
1097}
1098
1099#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1105pub enum CorruptionKind {
1106 WalReplayFailure,
1107 HeaderMalformed,
1108 SchemaInconsistent,
1109 EmbedderIdentityDrift,
1110}
1111
1112#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1118pub enum OpenStage {
1119 WalReplay,
1120 HeaderProbe,
1121 SchemaProbe,
1122 EmbedderIdentity,
1123}
1124
1125#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1131pub enum CorruptionLocator {
1132 FileOffset { offset: u64 },
1133 PageId { page: u32 },
1134 TableRow { table: &'static str, rowid: i64 },
1135 Vec0ShadowRow { partition: &'static str, rowid: i64 },
1136 MigrationStep { from: u32, to: u32 },
1137 OpaqueSqliteError { sqlite_extended_code: i32 },
1138}
1139
1140#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1146pub struct RecoveryHint {
1147 pub code: &'static str,
1148 pub doc_anchor: &'static str,
1149}
1150
1151#[derive(Clone, Debug, Eq, PartialEq)]
1152pub enum EngineOpenError {
1153 DatabaseLocked {
1154 holder_pid: Option<u32>,
1155 },
1156 Corruption(CorruptionDetail),
1157 IncompatibleSchemaVersion {
1158 seen: u32,
1159 supported: u32,
1160 },
1161 MigrationError {
1162 schema_version_before: u32,
1163 schema_version_current: u32,
1164 step_id: u32,
1165 },
1166 EmbedderIdentityMismatch {
1167 stored: EmbedderIdentity,
1168 supplied: EmbedderIdentity,
1169 },
1170 EmbedderDimensionMismatch {
1171 stored: u32,
1172 supplied: u32,
1173 },
1174 Embedder(RuntimeEmbedderError),
1176 Io {
1177 message: String,
1178 },
1179}
1180
1181#[derive(Clone)]
1184pub enum EmbedderChoice {
1185 Default,
1194 Caller(Arc<dyn Embedder>),
1197 None,
1201}
1202
1203impl Display for EngineOpenError {
1204 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1205 match self {
1206 Self::DatabaseLocked { holder_pid } => match holder_pid {
1207 Some(pid) => write!(f, "database is locked by process {pid}"),
1208 None => write!(f, "database is locked by another engine instance"),
1209 },
1210 Self::Corruption(detail) => {
1211 write!(
1212 f,
1213 "engine corruption at {:?} stage: {}",
1214 detail.stage, detail.recovery_hint.code
1215 )
1216 }
1217 Self::IncompatibleSchemaVersion { seen, supported } => write!(
1218 f,
1219 "database schema version {seen} is incompatible with supported version {supported}"
1220 ),
1221 Self::MigrationError {
1222 schema_version_before,
1223 schema_version_current,
1224 step_id,
1225 } => write!(
1226 f,
1227 "schema migration failed at step {step_id}; schema version remained between {schema_version_before} and {schema_version_current}"
1228 ),
1229 Self::EmbedderIdentityMismatch { stored, supplied } => write!(
1230 f,
1231 "embedder identity mismatch: stored {}@{}, supplied {}@{}",
1232 stored.name, stored.revision, supplied.name, supplied.revision,
1233 ),
1234 Self::EmbedderDimensionMismatch { stored, supplied } => write!(
1235 f,
1236 "embedder vector dimension mismatch: stored {stored}, supplied {supplied}",
1237 ),
1238 Self::Embedder(err) => match err {
1239 RuntimeEmbedderError::Timeout => write!(f, "embedder timeout during open"),
1240 RuntimeEmbedderError::Failed { message } => {
1241 write!(f, "embedder failure during open: {message}")
1242 }
1243 },
1244 Self::Io { message } => write!(f, "database I/O error: {message}"),
1245 }
1246 }
1247}
1248
1249impl Error for EngineOpenError {}
1250
1251#[derive(Clone, Debug, Eq, PartialEq)]
1252pub enum EngineError {
1253 Storage,
1254 Projection,
1255 Vector,
1256 Embedder,
1257 EmbedderNotConfigured,
1258 KindNotVectorIndexed,
1259 EmbedderDimensionMismatch { expected: u32, actual: u32 },
1260 Scheduler,
1261 OpStore,
1262 WriteValidation,
1263 SchemaValidation,
1264 Overloaded,
1265 Closing,
1266}
1267
1268impl Display for EngineError {
1269 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1270 match self {
1271 Self::Storage => write!(f, "storage error"),
1272 Self::Projection => write!(f, "projection error"),
1273 Self::Vector => write!(f, "vector error"),
1274 Self::Embedder => write!(f, "embedder error"),
1275 Self::EmbedderNotConfigured => write!(f, "embedder is not configured"),
1276 Self::KindNotVectorIndexed => write!(f, "kind is not configured for vector indexing"),
1277 Self::EmbedderDimensionMismatch { expected, actual } => {
1278 write!(f, "embedder dimension mismatch: expected {expected}, actual {actual}")
1279 }
1280 Self::Scheduler => write!(f, "scheduler error"),
1281 Self::OpStore => write!(f, "op-store error"),
1282 Self::WriteValidation => write!(f, "write validation error"),
1283 Self::SchemaValidation => write!(f, "schema validation error"),
1284 Self::Overloaded => write!(f, "engine overloaded"),
1285 Self::Closing => write!(f, "engine is closing"),
1286 }
1287 }
1288}
1289
1290impl EngineError {
1291 fn stable_code(&self) -> &'static str {
1296 match self {
1297 Self::Storage => "StorageError",
1298 Self::Projection => "ProjectionError",
1299 Self::Vector => "VectorError",
1300 Self::Embedder => "EmbedderError",
1301 Self::EmbedderNotConfigured => "EmbedderNotConfiguredError",
1302 Self::KindNotVectorIndexed => "KindNotVectorIndexedError",
1303 Self::EmbedderDimensionMismatch { .. } => "EmbedderDimensionMismatchError",
1304 Self::Scheduler => "SchedulerError",
1305 Self::OpStore => "OpStoreError",
1306 Self::WriteValidation => "WriteValidationError",
1307 Self::SchemaValidation => "SchemaValidationError",
1308 Self::Overloaded => "OverloadedError",
1309 Self::Closing => "ClosingError",
1310 }
1311 }
1312}
1313
1314impl Error for EngineError {}
1315
1316#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
1321pub struct CheckIntegrityOpts {
1322 pub quick: bool,
1323 pub full: bool,
1324 pub round_trip: bool,
1325}
1326
1327#[derive(Clone, Debug, Eq, PartialEq)]
1331pub enum Section {
1332 Clean,
1333 Findings(Vec<Finding>),
1334}
1335
1336#[derive(Clone, Debug, Eq, PartialEq)]
1340pub struct Finding {
1341 pub code: &'static str,
1342 pub stage: &'static str,
1343 pub locator: CorruptionLocator,
1344 pub doc_anchor: &'static str,
1345 pub detail: String,
1346}
1347
1348#[derive(Clone, Debug, Eq, PartialEq)]
1351pub struct IntegrityReport {
1352 pub physical: Section,
1353 pub logical: Section,
1354 pub semantic: Section,
1355}
1356
1357#[derive(Clone, Debug, Eq, PartialEq)]
1362pub struct SafeExportArtifact {
1363 pub export_path: PathBuf,
1364 pub manifest_path: PathBuf,
1365 pub manifest_sha256: String,
1366}
1367
1368#[derive(Clone, Debug, Eq, PartialEq)]
1372pub struct TraceReport {
1373 pub source_ref: String,
1374 pub events: Vec<TraceEvent>,
1375}
1376
1377#[derive(Clone, Debug, Eq, PartialEq)]
1380pub struct TraceEvent {
1381 pub write_cursor: u64,
1382 pub kind: String,
1383 pub table: &'static str,
1384}
1385
1386#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1391pub enum RebuildKind {
1392 Projections,
1393 Vec0,
1394}
1395
1396#[derive(Clone, Debug, Eq, PartialEq)]
1403pub struct RebuildReport {
1404 pub kind: RebuildKind,
1405 pub rows_invalidated: u64,
1406 pub rows_rebuilt: u64,
1407 pub projection_cursor_after: u64,
1408}
1409
1410#[derive(Clone, Debug, Eq, PartialEq)]
1414pub struct ExciseReport {
1415 pub source_ref: String,
1416 pub nodes_excised: u64,
1417 pub edges_excised: u64,
1418 pub projections_invalidated: u64,
1419}
1420
1421#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1425pub enum VerifyEmbedderStatus {
1426 Match,
1427 IdentityMismatch,
1428 DimensionMismatch,
1429 BothMismatch,
1430}
1431
1432#[derive(Clone, Debug, Eq, PartialEq)]
1436pub struct VerifyEmbedderReport {
1437 pub stored_identity: String,
1438 pub stored_dimension: u32,
1439 pub supplied_identity: String,
1440 pub supplied_dimension: u32,
1441 pub status: VerifyEmbedderStatus,
1442}
1443
1444#[derive(Clone, Debug, Eq, PartialEq)]
1446pub struct SchemaObject {
1447 pub name: String,
1448 pub sql: String,
1449}
1450
1451#[derive(Clone, Debug, Eq, PartialEq)]
1456pub struct DumpSchemaReport {
1457 pub user_version: u32,
1458 pub tables: Vec<SchemaObject>,
1459 pub indexes: Vec<SchemaObject>,
1460}
1461
1462#[derive(Clone, Debug, Eq, PartialEq)]
1464pub struct TableRowCount {
1465 pub name: String,
1466 pub rows: u64,
1467}
1468
1469#[derive(Clone, Debug, Eq, PartialEq)]
1473pub struct DumpRowCountsReport {
1474 pub counts: Vec<TableRowCount>,
1475}
1476
1477#[derive(Clone, Debug, Eq, PartialEq)]
1481pub struct DumpProfileReport {
1482 pub embedder_identity: String,
1483 pub embedder_dimension: u32,
1484 pub vectorized_kinds: Vec<String>,
1485}
1486
1487#[derive(Clone, Debug, PartialEq)]
1495pub struct MeanRecomputeReport {
1496 pub dim: u32,
1497 pub old_doc_count: u64,
1498 pub doc_count_requantized: u64,
1499 pub drift_cos_before: f32,
1500 pub mean_was_pinned: bool,
1501 pub elapsed_ms: u64,
1502}
1503
1504#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1508pub enum TruncateWalStatus {
1509 Done,
1510 Busy,
1511}
1512
1513#[derive(Clone, Debug, Eq, PartialEq)]
1517pub struct TruncateWalReport {
1518 pub status: TruncateWalStatus,
1519 pub busy: u32,
1520 pub log_frames: u32,
1521 pub checkpointed_frames: u32,
1522}
1523
1524impl Drop for Engine {
1525 fn drop(&mut self) {
1526 let _ = self.close();
1527 }
1528}
1529
1530impl Engine {
1531 pub fn open(path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
1532 Self::open_with_embedder_and_subscriber(
1533 path,
1534 default_embedder_identity(),
1535 None,
1536 None,
1537 None,
1538 &mut |_| {},
1539 )
1540 }
1541
1542 pub fn open_with_choice(
1551 path: impl Into<PathBuf>,
1552 choice: EmbedderChoice,
1553 ) -> Result<OpenedEngine, EngineOpenError> {
1554 match choice {
1555 EmbedderChoice::Default => Self::open_default_embedder(path),
1556 EmbedderChoice::Caller(embedder) => {
1557 let identity = embedder.identity();
1558 Self::open_with_embedder_and_subscriber(
1559 path,
1560 identity,
1561 Some(embedder),
1562 None,
1563 None,
1564 &mut |_| {},
1565 )
1566 }
1567 EmbedderChoice::None => Self::open_with_embedder_and_subscriber(
1568 path,
1569 default_embedder_identity(),
1570 None,
1571 None,
1572 None,
1573 &mut |_| {},
1574 ),
1575 }
1576 }
1577
1578 #[cfg(feature = "default-embedder")]
1583 fn open_default_embedder(path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
1584 use std::time::Instant as DownloadInstant;
1585 let download_start = DownloadInstant::now();
1586 let weights = fathomdb_embedder::loader::load_pinned_default_embedder().map_err(|err| {
1587 EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
1588 message: format!("default embedder loader: {err}"),
1589 })
1590 })?;
1591 let events = weights.events.clone();
1592 let download_ms = if weights.bytes_downloaded > 0 {
1593 Some(u64::try_from(download_start.elapsed().as_millis()).unwrap_or(u64::MAX))
1594 } else {
1595 None
1596 };
1597 let embedder =
1598 fathomdb_embedder::CandleBgeEmbedder::new_from_weights(weights).map_err(|err| {
1599 EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
1600 message: format!("default embedder construct: {err}"),
1601 })
1602 })?;
1603 let embedder: Arc<dyn Embedder> = Arc::new(embedder);
1604 let identity = embedder.identity();
1605 let loader_info = LoaderInfo { download_ms, events };
1606 Self::open_with_embedder_and_subscriber(
1607 path,
1608 identity,
1609 Some(embedder),
1610 Some(loader_info),
1611 None,
1612 &mut |_| {},
1613 )
1614 }
1615
1616 #[cfg(not(feature = "default-embedder"))]
1617 fn open_default_embedder(_path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
1618 Err(EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
1619 message: "EmbedderChoice::Default requires the `default-embedder` Cargo feature"
1620 .to_string(),
1621 }))
1622 }
1623
1624 pub fn open_with_migration_event_sink(
1625 path: impl Into<PathBuf>,
1626 mut emit_migration_event: impl FnMut(&MigrationStepReport),
1627 ) -> Result<OpenedEngine, EngineOpenError> {
1628 Self::open_with_embedder_and_subscriber(
1629 path,
1630 default_embedder_identity(),
1631 None,
1632 None,
1633 None,
1634 &mut emit_migration_event,
1635 )
1636 }
1637
1638 #[cfg(debug_assertions)]
1639 #[doc(hidden)]
1640 pub fn open_with_migrations_for_test(
1641 path: impl Into<PathBuf>,
1642 migrations: &'static [fathomdb_schema::Migration],
1643 mut emit_migration_event: impl FnMut(&MigrationStepReport),
1644 ) -> Result<OpenedEngine, EngineOpenError> {
1645 Self::open_with_migrations(
1646 path,
1647 migrations,
1648 default_embedder_identity(),
1649 None,
1650 None,
1651 &mut emit_migration_event,
1652 None,
1653 )
1654 }
1655
1656 #[doc(hidden)]
1657 pub fn open_with_subscriber_for_test(
1658 path: impl Into<PathBuf>,
1659 subscriber: Arc<dyn lifecycle::Subscriber>,
1660 ) -> Result<OpenedEngine, EngineOpenError> {
1661 Self::open_with_embedder_and_subscriber(
1662 path,
1663 default_embedder_identity(),
1664 None,
1665 None,
1666 Some(subscriber),
1667 &mut |_| {},
1668 )
1669 }
1670
1671 #[doc(hidden)]
1672 pub fn open_without_embedder_for_test(
1673 path: impl Into<PathBuf>,
1674 ) -> Result<OpenedEngine, EngineOpenError> {
1675 Self::open_with_embedder_and_subscriber(
1676 path,
1677 default_embedder_identity(),
1678 None,
1679 None,
1680 None,
1681 &mut |_| {},
1682 )
1683 }
1684
1685 #[doc(hidden)]
1686 pub fn open_with_embedder_for_test(
1687 path: impl Into<PathBuf>,
1688 embedder: Arc<dyn Embedder>,
1689 ) -> Result<OpenedEngine, EngineOpenError> {
1690 let identity = embedder.identity();
1691 Self::open_with_embedder_and_subscriber(
1692 path,
1693 identity,
1694 Some(embedder),
1695 None,
1696 None,
1697 &mut |_| {},
1698 )
1699 }
1700
1701 fn open_with_embedder_and_subscriber(
1702 path: impl Into<PathBuf>,
1703 embedder_identity: EmbedderIdentity,
1704 runtime_embedder: Option<Arc<dyn Embedder>>,
1705 loader_info: Option<LoaderInfo>,
1706 initial_subscriber: Option<Arc<dyn lifecycle::Subscriber>>,
1707 emit_migration_event: &mut impl FnMut(&MigrationStepReport),
1708 ) -> Result<OpenedEngine, EngineOpenError> {
1709 Self::open_with_migrations(
1710 path,
1711 MIGRATIONS,
1712 embedder_identity,
1713 runtime_embedder,
1714 loader_info,
1715 emit_migration_event,
1716 initial_subscriber,
1717 )
1718 }
1719
1720 fn open_with_migrations(
1721 path: impl Into<PathBuf>,
1722 migrations: &'static [fathomdb_schema::Migration],
1723 embedder_identity: EmbedderIdentity,
1724 runtime_embedder: Option<Arc<dyn Embedder>>,
1725 loader_info: Option<LoaderInfo>,
1726 emit_migration_event: &mut impl FnMut(&MigrationStepReport),
1727 initial_subscriber: Option<Arc<dyn lifecycle::Subscriber>>,
1728 ) -> Result<OpenedEngine, EngineOpenError> {
1729 let canonical_path = canonical_database_path(&path.into())?;
1730 let lock = acquire_lock(&canonical_path)?;
1731 let open_result = Self::open_locked(
1732 canonical_path.clone(),
1733 migrations,
1734 &embedder_identity,
1735 emit_migration_event,
1736 );
1737
1738 match open_result {
1739 Ok((connection, readers, mut report, reader_lookaside_rcs)) => {
1740 if let Some(info) = loader_info {
1746 if info.download_ms.is_some() {
1747 report.embedder_download_ms = info.download_ms;
1748 }
1749 if !info.events.is_empty() {
1750 report.embedder_events = info.events;
1751 }
1752 }
1753 let next_cursor = load_next_cursor(&connection);
1754 let subscribers = Arc::new(lifecycle::SubscriberRegistry::new());
1755 let profiling_enabled = Arc::new(AtomicBool::new(false));
1756 let slow_threshold_ms = Arc::new(AtomicU64::new(DEFAULT_SLOW_THRESHOLD_MS));
1757 let mut profile_contexts: Vec<Box<ProfileContext>> = Vec::new();
1758 let projection_runtime = ProjectionRuntime::new(
1759 canonical_path.clone(),
1760 runtime_embedder.clone(),
1761 embedder_identity.clone(),
1762 report.embedder_mean_vec_pinned,
1763 );
1764
1765 install_profile_callback(
1766 &connection,
1767 &subscribers,
1768 &profiling_enabled,
1769 &slow_threshold_ms,
1770 &mut profile_contexts,
1771 );
1772 for reader in &readers {
1773 install_profile_callback(
1774 reader,
1775 &subscribers,
1776 &profiling_enabled,
1777 &slow_threshold_ms,
1778 &mut profile_contexts,
1779 );
1780 }
1781
1782 let opened = OpenedEngine {
1783 engine: Self {
1784 path: canonical_path.clone(),
1785 next_cursor: AtomicU64::new(next_cursor),
1786 closed: AtomicBool::new(false),
1787 lock: Mutex::new(Some(lock)),
1788 connection: Mutex::new(Some(connection)),
1789 reader_pool: ReaderWorkerPool::new(readers),
1790 counters: lifecycle::Counters::new(),
1791 subscribers,
1792 profiling_enabled,
1793 slow_threshold_ms,
1794 runtime_embedder,
1795 runtime_embedder_identity: embedder_identity,
1796 projection_runtime,
1797 provenance_row_cap: AtomicU64::new(DEFAULT_PROVENANCE_ROW_CAP),
1798 profile_contexts: Mutex::new(profile_contexts),
1799 reader_lookaside_rcs,
1800 #[cfg(debug_assertions)]
1801 force_next_commit_failure: AtomicBool::new(false),
1802 },
1803 report,
1804 };
1805 if let Some(subscriber) = initial_subscriber {
1806 opened.engine.subscribers.attach_persistent(subscriber);
1807 }
1808 if database_has_pending_projection_work(&canonical_path).unwrap_or(false) {
1809 opened.engine.projection_runtime.notify_new_work();
1810 }
1811 Ok(opened)
1812 }
1813 Err(err) => {
1814 if let Some(subscriber) = initial_subscriber {
1815 emit_open_error_event(&subscriber, &err);
1816 }
1817 drop(lock);
1818 Err(err)
1819 }
1820 }
1821 }
1822
1823 fn open_locked(
1824 path: PathBuf,
1825 migrations: &'static [fathomdb_schema::Migration],
1826 embedder_identity: &EmbedderIdentity,
1827 emit_migration_event: &mut impl FnMut(&MigrationStepReport),
1828 ) -> Result<(Connection, Vec<Connection>, OpenReport, Vec<i32>), EngineOpenError> {
1829 init_perf_experiments_runtime();
1830 register_sqlite_vec_extension();
1831 let mut connection = Connection::open(&path)
1832 .map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))?;
1833 probe_database_header(&connection)?;
1841 probe_open_integrity(&connection)?;
1842 probe_wal_sidecar(&path)?;
1843 apply_perf_experiment_writer_pragmas(&connection);
1849 connection
1850 .pragma_update(None, "journal_mode", "WAL")
1851 .map_err(|err| map_open_sqlite_error(err, OpenStage::WalReplay))?;
1852
1853 reject_legacy_shape(&connection)?;
1854 let migration = migrate_with_event_sink(&connection, migrations, emit_migration_event)
1855 .map_err(map_migration_error)?;
1856 if migration.schema_version_after >= SEARCH_INDEX_TOKENIZER_SCHEMA_VERSION
1873 && !search_index_tokenizer_reproject_complete(&connection).map_err(|_| {
1874 EngineOpenError::Io {
1875 message: "could not read search_index tokenizer reproject marker".to_string(),
1876 }
1877 })?
1878 {
1879 reproject_search_index_after_tokenizer_upgrade(&connection).map_err(|_| {
1880 EngineOpenError::Io {
1881 message: "could not re-tokenize search_index after tokenizer upgrade"
1882 .to_string(),
1883 }
1884 })?;
1885 }
1886 let mut embedder_mean_vec_pinned = check_embedder_profile(&connection, embedder_identity)?;
1887 ensure_vector_partition(&mut connection, embedder_identity.dimension).map_err(|_| {
1888 EngineOpenError::Io { message: "could not initialize vector partition".to_string() }
1889 })?;
1890
1891 if identity_requires_mean_centering(embedder_identity) && !embedder_mean_vec_pinned {
1899 let row_count: u64 = connection
1900 .query_row("SELECT COUNT(*) FROM vector_default", [], |row| row.get(0))
1901 .unwrap_or(0);
1902 if row_count >= MEAN_VEC_PIN_THRESHOLD {
1903 recover_mean_vec_pin(&mut connection, embedder_identity).map_err(|_| {
1904 EngineOpenError::Io {
1905 message: "could not recover mean-centering pin".to_string(),
1906 }
1907 })?;
1908 embedder_mean_vec_pinned = true;
1909 }
1910 }
1911
1912 let warmup_started = Instant::now();
1913 let embedder_mean_centering_required = embedder_identity.name == BGE_SMALL_EMBEDDER_NAME;
1918 let report = OpenReport {
1922 schema_version_before: migration.schema_version_before,
1923 schema_version_after: migration.schema_version_after,
1924 migration_steps: migration.migration_steps,
1925 embedder_warmup_ms: u64::try_from(warmup_started.elapsed().as_millis())
1926 .unwrap_or(u64::MAX),
1927 query_backend: "fathomdb-query + sqlite-vec",
1928 default_embedder: embedder_identity.clone(),
1929 embedder_download_ms: None,
1932 embedder_events: Vec::new(),
1934 embedder_mean_centering_required,
1935 embedder_mean_vec_pinned,
1936 };
1937
1938 let mut readers = Vec::with_capacity(READER_POOL_SIZE);
1939 let mut lookaside_rcs: Vec<i32> = Vec::with_capacity(READER_POOL_SIZE);
1940 for _ in 0..READER_POOL_SIZE {
1941 let reader = Connection::open(&path)
1942 .map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))?;
1943 let rc: i32 = configure_reader_lookaside(&reader);
1948 debug_assert_eq!(
1949 rc,
1950 rusqlite::ffi::SQLITE_OK,
1951 "sqlite3_db_config(LOOKASIDE) must return SQLITE_OK on a freshly opened reader",
1952 );
1953 lookaside_rcs.push(rc);
1954 reader
1955 .pragma_update(None, "journal_mode", "WAL")
1956 .map_err(|err| map_open_sqlite_error(err, OpenStage::WalReplay))?;
1957 reader
1958 .pragma_update(None, "query_only", "ON")
1959 .map_err(|err| map_open_sqlite_error(err, OpenStage::SchemaProbe))?;
1960 apply_perf_experiment_reader_pragmas(&reader);
1961 readers.push(reader);
1962 }
1963
1964 Ok((connection, readers, report, lookaside_rcs))
1965 }
1966
1967 #[must_use]
1968 pub fn path(&self) -> &Path {
1969 &self.path
1970 }
1971
1972 pub fn write(&self, batch: &[PreparedWrite]) -> Result<WriteReceipt, EngineError> {
1973 let category = if batch_is_admin(batch) {
1974 lifecycle::EventCategory::Admin
1975 } else {
1976 lifecycle::EventCategory::Writer
1977 };
1978 self.emit_event(lifecycle::Phase::Started, category, None);
1979 let started = Instant::now();
1980 let outcome = self.write_inner(batch);
1981 self.detect_slow(started, category);
1982 match outcome {
1983 Ok(receipt) => {
1984 let rows = u64::try_from(batch.len()).unwrap_or(u64::MAX);
1985 if batch_is_admin(batch) {
1986 self.counters.record_admin();
1987 } else {
1988 self.counters.record_write(rows);
1989 }
1990 self.emit_event(lifecycle::Phase::Finished, category, None);
1991 Ok(receipt)
1992 }
1993 Err(err) => {
1994 let code = err.stable_code();
1995 self.counters.record_error(code);
1996 self.emit_event(lifecycle::Phase::Failed, category, Some(code));
1999 self.emit_event(
2000 lifecycle::Phase::Failed,
2001 lifecycle::EventCategory::Error,
2002 Some(code),
2003 );
2004 Err(err)
2005 }
2006 }
2007 }
2008
2009 fn write_inner(&self, batch: &[PreparedWrite]) -> Result<WriteReceipt, EngineError> {
2010 self.ensure_open()?;
2011
2012 if batch.is_empty() {
2013 return Err(EngineError::WriteValidation);
2014 }
2015
2016 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2017 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
2018 let plans = validate_batch(connection, batch)?;
2019 let projection_jobs = collect_projection_jobs(connection, batch)?;
2020 #[cfg(debug_assertions)]
2021 if self.force_next_commit_failure.swap(false, Ordering::SeqCst) {
2022 return Err(EngineError::Storage);
2023 }
2024 let base_cursor = self.next_cursor.load(Ordering::SeqCst);
2032 let increment = u64::try_from(batch.len()).unwrap_or(u64::MAX);
2033 let last_cursor = base_cursor.saturating_add(increment);
2034 let pending_projection = !projection_jobs.is_empty();
2035
2036 let dangling_edge_endpoints = match commit_batch(
2037 connection,
2038 batch,
2039 &plans,
2040 base_cursor,
2041 self.provenance_row_cap.load(Ordering::Relaxed),
2042 ) {
2043 Ok(count) => count,
2044 Err(err) => {
2045 self.emit_sqlite_internal_error(&err);
2046 return Err(EngineError::Storage);
2047 }
2048 };
2049 self.next_cursor.store(last_cursor, Ordering::SeqCst);
2050 if pending_projection {
2051 self.projection_runtime.notify_new_work();
2052 }
2053
2054 let row_cursors = (0..batch.len())
2057 .map(|i| base_cursor.saturating_add((i as u64).saturating_add(1)))
2058 .collect();
2059 Ok(WriteReceipt { cursor: last_cursor, row_cursors, dangling_edge_endpoints })
2060 }
2061
2062 pub fn search(&self, query: &str) -> Result<SearchResult, EngineError> {
2063 self.search_filtered(query, None)
2064 }
2065
2066 pub fn search_filtered(
2072 &self,
2073 query: &str,
2074 filter: Option<SearchFilter>,
2075 ) -> Result<SearchResult, EngineError> {
2076 self.emit_event(lifecycle::Phase::Started, lifecycle::EventCategory::Search, None);
2077 let started = Instant::now();
2078 let outcome = self.search_inner(query, filter);
2079 self.detect_slow(started, lifecycle::EventCategory::Search);
2080 match outcome {
2081 Ok(result) => {
2082 self.counters.record_query();
2083 self.emit_event(lifecycle::Phase::Finished, lifecycle::EventCategory::Search, None);
2084 Ok(result)
2085 }
2086 Err(err) => {
2087 let code = err.stable_code();
2088 self.counters.record_error(code);
2089 self.emit_event(
2090 lifecycle::Phase::Failed,
2091 lifecycle::EventCategory::Search,
2092 Some(code),
2093 );
2094 self.emit_event(
2095 lifecycle::Phase::Failed,
2096 lifecycle::EventCategory::Error,
2097 Some(code),
2098 );
2099 Err(err)
2100 }
2101 }
2102 }
2103
2104 fn detect_slow(&self, started: Instant, category: lifecycle::EventCategory) {
2105 let elapsed = started.elapsed();
2106 let threshold = self.slow_threshold_ms.load(Ordering::Relaxed);
2107 let threshold_duration = std::time::Duration::from_millis(threshold);
2108 if elapsed > threshold_duration {
2109 self.emit_event(lifecycle::Phase::Slow, category, None);
2116 }
2117 }
2118
2119 fn emit_event(
2120 &self,
2121 phase: lifecycle::Phase,
2122 category: lifecycle::EventCategory,
2123 code: Option<&'static str>,
2124 ) {
2125 let event =
2126 lifecycle::Event { phase, source: lifecycle::EventSource::Engine, category, code };
2127 self.subscribers.dispatch(&event);
2128 }
2129
2130 fn emit_sqlite_internal_error(&self, err: &rusqlite::Error) {
2137 if let Some(code) = sqlite_extended_code_name(err) {
2138 let event = lifecycle::Event {
2139 phase: lifecycle::Phase::Failed,
2140 source: lifecycle::EventSource::SqliteInternal,
2141 category: lifecycle::EventCategory::Error,
2142 code: Some(code),
2143 };
2144 self.subscribers.dispatch(&event);
2145 }
2146 }
2147
2148 fn search_inner(
2149 &self,
2150 query: &str,
2151 filter: Option<SearchFilter>,
2152 ) -> Result<SearchResult, EngineError> {
2153 self.ensure_open()?;
2154 if query.trim().is_empty() {
2155 return Err(EngineError::WriteValidation);
2156 }
2157
2158 let compiled = compile_text_query(query);
2159 let raw_query_vector =
2175 self.runtime_embedder.as_ref().and_then(|embedder| embedder.embed(query).ok());
2176 let query_vector_bin = match raw_query_vector.as_ref() {
2177 Some(vector) if identity_requires_mean_centering(&self.runtime_embedder_identity) => {
2178 let pinned = {
2179 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2180 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2181 read_pinned_mean_vec(connection, self.runtime_embedder_identity.dimension)?
2182 };
2183 match pinned {
2184 Some(mean) => serde_json::to_string(&subtract_mean(vector, &mean)).ok(),
2185 None => serde_json::to_string(vector).ok(),
2186 }
2187 }
2188 Some(vector) => serde_json::to_string(vector).ok(),
2189 None => None,
2190 };
2191 let query_vector = raw_query_vector.and_then(|vector| serde_json::to_string(&vector).ok());
2192 let search_limit = self
2196 .projection_runtime
2197 .shared
2198 .search_limit_override
2199 .load(Ordering::SeqCst)
2200 .max(SEARCH_RERANK_LIMIT);
2201 let recency_enabled =
2202 self.projection_runtime.shared.recency_reweight_enabled.load(Ordering::SeqCst);
2203 let vector_stage_only =
2204 self.projection_runtime.shared.vector_stage_only_for_test.load(Ordering::SeqCst);
2205 let (response_tx, response_rx) = mpsc::sync_channel::<ReaderResponse>(1);
2206 let request = ReaderRequest::Search {
2207 compiled,
2208 query_vector,
2209 query_vector_bin,
2210 search_limit,
2211 filter: filter.map(Box::new),
2212 recency_enabled,
2213 vector_stage_only,
2214 respond: response_tx,
2215 };
2216 if self.reader_pool.dispatch(request).is_err() {
2217 return Err(EngineError::Closing);
2218 }
2219 let search_result = response_rx.recv().map_err(|_| EngineError::Storage)?;
2220 let (cursor, soft_fallback, results) = match search_result {
2221 Ok(result) => result,
2222 Err(err) => {
2223 self.emit_sqlite_internal_error(&err);
2224 return Err(EngineError::Storage);
2225 }
2226 };
2227
2228 Ok(SearchResult { projection_cursor: cursor, soft_fallback, results })
2229 }
2230
2231 pub fn read_get(&self, logical_id: &str) -> Result<Option<NodeRecord>, EngineError> {
2236 let ids = [logical_id.to_string()];
2237 let rows = self.read_get_many(&ids)?;
2238 Ok(rows.into_iter().next().flatten())
2239 }
2240
2241 pub fn read_get_many(
2245 &self,
2246 logical_ids: &[String],
2247 ) -> Result<Vec<Option<NodeRecord>>, EngineError> {
2248 self.ensure_open()?;
2249 if logical_ids.is_empty() {
2250 return Ok(Vec::new());
2251 }
2252 let (response_tx, response_rx) = mpsc::sync_channel(1);
2253 let request =
2254 ReaderRequest::GetById { logical_ids: logical_ids.to_vec(), respond: response_tx };
2255 if self.reader_pool.dispatch(request).is_err() {
2256 return Err(EngineError::Closing);
2257 }
2258 match response_rx.recv().map_err(|_| EngineError::Storage)? {
2259 Ok(rows) => Ok(rows),
2260 Err(err) => {
2261 self.emit_sqlite_internal_error(&err);
2262 Err(EngineError::Storage)
2263 }
2264 }
2265 }
2266
2267 pub fn read_collection(
2272 &self,
2273 collection: &str,
2274 after_id: Option<i64>,
2275 limit: usize,
2276 ) -> Result<Vec<OpStoreRow>, EngineError> {
2277 self.read_collection_dispatch(collection, after_id, limit)
2278 }
2279
2280 pub fn read_mutations(
2283 &self,
2284 collection: &str,
2285 after_id: Option<i64>,
2286 limit: usize,
2287 ) -> Result<Vec<OpStoreRow>, EngineError> {
2288 self.read_collection_dispatch(collection, after_id, limit)
2289 }
2290
2291 fn read_collection_dispatch(
2292 &self,
2293 collection: &str,
2294 after_id: Option<i64>,
2295 limit: usize,
2296 ) -> Result<Vec<OpStoreRow>, EngineError> {
2297 self.ensure_open()?;
2298 let (response_tx, response_rx) = mpsc::sync_channel(1);
2299 let request = ReaderRequest::ReadCollection {
2300 collection: collection.to_string(),
2301 after_id,
2302 limit,
2303 respond: response_tx,
2304 };
2305 if self.reader_pool.dispatch(request).is_err() {
2306 return Err(EngineError::Closing);
2307 }
2308 match response_rx.recv().map_err(|_| EngineError::Storage)? {
2309 Ok(rows) => Ok(rows),
2310 Err(err) => {
2311 self.emit_sqlite_internal_error(&err);
2312 Err(EngineError::Storage)
2313 }
2314 }
2315 }
2316
2317 pub fn close(&self) -> Result<(), EngineError> {
2318 self.closed.store(true, Ordering::SeqCst);
2319 self.projection_runtime.stop();
2320 self.reader_pool.shutdown();
2329 if let Ok(mut connection) = self.connection.lock() {
2330 if let Some(conn) = connection.as_ref() {
2331 uninstall_profile_callback(conn);
2332 }
2333 connection.take();
2334 }
2335 if let Ok(mut contexts) = self.profile_contexts.lock() {
2336 contexts.clear();
2337 }
2338 if let Ok(mut lock) = self.lock.lock() {
2339 lock.take();
2340 }
2341 Ok(())
2342 }
2343
2344 pub fn drain(&self, timeout_ms: u64) -> Result<(), EngineError> {
2349 self.ensure_open()?;
2350 if self.projection_runtime.wait_for_idle(timeout_ms) {
2351 Ok(())
2352 } else {
2353 Err(EngineError::Scheduler)
2354 }
2355 }
2356
2357 #[must_use]
2361 pub fn counters(&self) -> CounterSnapshot {
2362 self.counters.snapshot()
2363 }
2364
2365 pub fn set_profiling(&self, enabled: bool) -> Result<(), EngineError> {
2371 self.profiling_enabled.store(enabled, Ordering::Relaxed);
2372 Ok(())
2373 }
2374
2375 pub fn set_slow_threshold_ms(&self, value: u64) -> Result<(), EngineError> {
2381 self.slow_threshold_ms.store(value, Ordering::Relaxed);
2382 Ok(())
2383 }
2384
2385 #[must_use]
2391 pub fn subscribe(&self, subscriber: Arc<dyn lifecycle::Subscriber>) -> Subscription {
2392 self.subscribers.attach(subscriber)
2393 }
2394
2395 #[cfg(debug_assertions)]
2396 #[doc(hidden)]
2397 pub fn reader_worker_count_for_test(&self) -> usize {
2398 self.reader_pool.worker_count()
2399 }
2400
2401 #[cfg(debug_assertions)]
2402 #[doc(hidden)]
2403 pub fn live_reader_worker_count_for_test(&self) -> usize {
2404 self.reader_pool.live_count()
2405 }
2406
2407 #[cfg(debug_assertions)]
2412 #[doc(hidden)]
2413 pub fn reader_lookaside_config_rcs_for_test(&self) -> Vec<i32> {
2414 self.reader_lookaside_rcs.clone()
2415 }
2416
2417 #[cfg(debug_assertions)]
2423 #[doc(hidden)]
2424 pub fn reader_lookaside_used_per_worker_for_test(&self) -> Vec<i32> {
2425 self.reader_pool.lookaside_used_per_worker()
2426 }
2427
2428 #[cfg(debug_assertions)]
2434 #[doc(hidden)]
2435 pub fn cache_status_per_worker_for_test(&self, label: &str) -> Vec<CacheStatusReply> {
2436 self.reader_pool.cache_status_per_worker(label)
2437 }
2438
2439 #[cfg(debug_assertions)]
2440 #[doc(hidden)]
2441 pub fn force_next_commit_failure_for_test(&self) {
2442 self.force_next_commit_failure.store(true, Ordering::SeqCst);
2443 }
2444
2445 #[cfg(debug_assertions)]
2452 #[doc(hidden)]
2453 pub fn execute_for_test(&self, sql: &str) -> Result<(), EngineError> {
2454 self.ensure_open()?;
2455 let started = Instant::now();
2456 {
2457 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2458 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
2459 connection.execute_batch(sql).map_err(|_| EngineError::Storage)?;
2460 }
2461 self.detect_slow(started, lifecycle::EventCategory::Search);
2462 Ok(())
2463 }
2464
2465 #[doc(hidden)]
2476 #[cfg(debug_assertions)]
2477 pub fn run_one_thread_poison_for_test(&self) -> Result<(), EngineError> {
2478 self.ensure_open()?;
2479
2480 self.write(&[PreparedWrite::Node {
2483 kind: "doc".to_string(),
2484 body: "poison-fixture-seed".to_string(),
2485 source_id: None,
2486 logical_id: None,
2487 }])?;
2488
2489 let poison_outcome: Mutex<Option<EngineError>> = Mutex::new(None);
2490 let poison_thread_id: AtomicU64 = AtomicU64::new(0);
2491
2492 thread::scope(|scope| {
2493 for _ in 0..4 {
2495 scope.spawn(|| {
2496 for _ in 0..4 {
2497 let _ = self.search("poison-fixture-seed");
2498 }
2499 });
2500 }
2501 scope.spawn(|| {
2503 let _ = self.write(&[PreparedWrite::Node {
2504 kind: "doc".to_string(),
2505 body: "writer-progress".to_string(),
2506 source_id: None,
2507 logical_id: None,
2508 }]);
2509 });
2510 scope.spawn(|| {
2513 poison_thread_id.store(1, Ordering::SeqCst);
2516 if let Err(err) = self.write(&[]) {
2517 *poison_outcome.lock().expect("poison_outcome lock") = Some(err);
2518 }
2519 });
2520 });
2521
2522 let err = poison_outcome
2523 .into_inner()
2524 .expect("poison_outcome lock")
2525 .expect("poison thread must produce a deterministic error");
2526
2527 let projection_state = match self.projection_status_for_test("doc") {
2528 Ok(lifecycle::ProjectionStatus::Pending) => "Pending",
2529 Ok(lifecycle::ProjectionStatus::Failed) => "Failed",
2530 Ok(lifecycle::ProjectionStatus::UpToDate) => "UpToDate",
2531 Err(_) => "UpToDate",
2536 };
2537
2538 let context = lifecycle::StressFailureContext {
2539 thread_group_id: poison_thread_id.load(Ordering::SeqCst),
2540 op_kind: "write".to_string(),
2541 last_error_chain: vec![err.stable_code().to_string(), err.to_string()],
2542 projection_state: projection_state.to_string(),
2543 };
2544 self.subscribers.dispatch_stress_failure(&context);
2545 Ok(())
2546 }
2547
2548 #[doc(hidden)]
2549 pub fn set_projection_scheduler_frozen_for_test(&self, frozen: bool) {
2550 self.projection_runtime.set_frozen(frozen);
2551 }
2552
2553 #[doc(hidden)]
2554 pub fn set_projection_retry_delays_for_test(&self, delays_ms: &[u64]) {
2555 self.projection_runtime.set_retry_delays_for_test(delays_ms);
2556 }
2557
2558 #[doc(hidden)]
2561 pub fn set_embed_timeout_ms_for_test(&self, timeout_ms: u64) {
2562 self.projection_runtime.set_embed_timeout_ms_for_test(timeout_ms);
2563 }
2564
2565 #[doc(hidden)]
2568 pub fn set_embed_circuit_threshold_for_test(&self, threshold: u64) {
2569 self.projection_runtime.set_embed_circuit_threshold_for_test(threshold);
2570 }
2571
2572 #[doc(hidden)]
2574 pub fn embed_circuit_open_for_test(&self) -> bool {
2575 self.projection_runtime.embed_circuit_open_for_test()
2576 }
2577
2578 #[doc(hidden)]
2579 pub fn projection_status_for_test(
2580 &self,
2581 kind: &str,
2582 ) -> Result<lifecycle::ProjectionStatus, EngineError> {
2583 self.ensure_open()?;
2584 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2585 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2586 projection_status(connection, kind)
2587 }
2588
2589 #[doc(hidden)]
2590 pub fn has_vector_for_cursor_for_test(&self, cursor: u64) -> Result<bool, EngineError> {
2591 self.ensure_open()?;
2592 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2593 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2594 terminal_state_for_cursor(connection, cursor)
2595 .map(|state| matches!(state.as_deref(), Some("up_to_date")))
2596 .map_err(|_| EngineError::Storage)
2597 }
2598
2599 #[doc(hidden)]
2600 pub fn projection_failure_count_for_test(&self, cursor: u64) -> Result<u64, EngineError> {
2601 self.ensure_open()?;
2602 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2603 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2604 connection
2605 .query_row(
2606 "SELECT COUNT(*) FROM operational_mutations
2607 WHERE collection_name = 'projection_failures'
2608 AND record_key = ?1",
2609 [cursor.to_string()],
2610 |row| row.get::<_, u64>(0),
2611 )
2612 .map_err(|_| EngineError::Storage)
2613 }
2614
2615 #[doc(hidden)]
2616 pub fn set_provenance_row_cap_for_test(&self, cap: Option<u64>) {
2617 self.provenance_row_cap.store(cap.unwrap_or(0), Ordering::Relaxed);
2618 }
2619
2620 #[doc(hidden)]
2621 pub fn provenance_row_count_for_test(&self) -> Result<u64, EngineError> {
2622 self.ensure_open()?;
2623 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2624 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2625 connection
2626 .query_row("SELECT COUNT(*) FROM operational_mutations", [], |row| row.get::<_, u64>(0))
2627 .map_err(|_| EngineError::Storage)
2628 }
2629
2630 #[doc(hidden)]
2631 pub fn oldest_provenance_record_key_for_test(
2632 &self,
2633 collection: &str,
2634 ) -> Result<Option<String>, EngineError> {
2635 self.ensure_open()?;
2636 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2637 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2638 connection
2639 .query_row(
2640 "SELECT record_key FROM operational_mutations
2641 WHERE collection_name = ?1
2642 ORDER BY id
2643 LIMIT 1",
2644 [collection],
2645 |row| row.get::<_, String>(0),
2646 )
2647 .map(Some)
2648 .or_else(|err| match err {
2649 rusqlite::Error::QueryReturnedNoRows => Ok(None),
2650 _ => Err(EngineError::Storage),
2651 })
2652 }
2653
2654 #[doc(hidden)]
2655 pub fn configure_vector_kind_for_test(&self, kind: &str) -> Result<(), EngineError> {
2656 self.ensure_open()?;
2657 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2658 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
2659 connection
2660 .execute(
2661 "INSERT OR REPLACE INTO _fathomdb_vector_kinds(kind, profile, created_at)
2662 VALUES(?1, ?2, 0)",
2663 params![kind, DEFAULT_VECTOR_PROFILE],
2664 )
2665 .map_err(|_| EngineError::Storage)?;
2666 Ok(())
2667 }
2668
2669 #[doc(hidden)]
2670 pub fn write_vector_for_test(
2671 &self,
2672 kind: &str,
2673 text: &str,
2674 ) -> Result<WriteReceipt, EngineError> {
2675 self.ensure_open()?;
2676 let embedder =
2677 self.runtime_embedder.as_ref().cloned().ok_or(EngineError::EmbedderNotConfigured)?;
2678
2679 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2680 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
2681 if !kind_is_vector_indexed(connection, kind)? {
2682 return Err(EngineError::KindNotVectorIndexed);
2683 }
2684
2685 let expected = default_profile_dimension(connection)?;
2686 ensure_vector_partition(connection, expected).map_err(|_| EngineError::Storage)?;
2687 let vector = embedder.embed(text).map_err(map_runtime_embedder_error)?;
2688 let actual = u32::try_from(vector.len()).unwrap_or(u32::MAX);
2689 if actual != expected {
2690 return Err(EngineError::EmbedderDimensionMismatch { expected, actual });
2691 }
2692
2693 let cursor = self.next_cursor.load(Ordering::SeqCst).saturating_add(1);
2694 let blob = encode_vector_blob(&vector);
2700 let bin_blob = if identity_requires_mean_centering(&self.runtime_embedder_identity) {
2701 match read_pinned_mean_vec(connection, self.runtime_embedder_identity.dimension)? {
2702 Some(mean) => encode_vector_blob(&subtract_mean(&vector, &mean)),
2703 None => blob.clone(),
2704 }
2705 } else {
2706 blob.clone()
2707 };
2708 let source_type = resolve_source_type(kind)?;
2709 let now_unix =
2710 SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs() as i64;
2711
2712 let pin_event = {
2717 let runtime = &self.projection_runtime.shared;
2718 let mut accumulator =
2719 runtime.mean_accumulator.lock().map_err(|_| EngineError::Storage)?;
2720 if let Some(acc) = accumulator.as_mut() {
2721 acc.add(&vector);
2722 if acc.count() >= MEAN_VEC_PIN_THRESHOLD {
2723 let mean = acc.materialize();
2724 *accumulator = None;
2725 Some(mean)
2726 } else {
2727 None
2728 }
2729 } else {
2730 None
2731 }
2732 };
2733
2734 let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
2735 tx.execute(
2736 "INSERT INTO _fathomdb_vector_rows(rowid, kind, write_cursor) VALUES(?1, ?2, ?3)",
2737 params![cursor, kind, cursor],
2738 )
2739 .map_err(|_| EngineError::Storage)?;
2740 tx.execute(
2741 "INSERT INTO vector_default(
2747 rowid, embedding, embedding_bin, source_type, kind, created_at, status
2748 ) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
2749 params![cursor, blob, bin_blob, source_type, kind, now_unix],
2750 )
2751 .map_err(|_| EngineError::Storage)?;
2752
2753 let mut emitted_event: Option<EmbedderEvent> = None;
2754 if let Some(mean_vec) = pin_event {
2755 let mean_bytes = encode_vector_blob(&mean_vec);
2756 tx.execute(
2757 "UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
2758 params![mean_bytes],
2759 )
2760 .map_err(|_| EngineError::Storage)?;
2761 let rows: Vec<(i64, Vec<u8>)> = {
2764 let mut statement = tx
2765 .prepare("SELECT rowid, embedding FROM vector_default ORDER BY rowid")
2766 .map_err(|_| EngineError::Storage)?;
2767 let mapped = statement
2768 .query_map([], |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?)))
2769 .map_err(|_| EngineError::Storage)?;
2770 let mut out = Vec::new();
2771 for r in mapped {
2772 out.push(r.map_err(|_| EngineError::Storage)?);
2773 }
2774 out
2775 };
2776 let (doc_count, _) = run_pin_and_requantize_pass(&tx, &rows, &mean_vec)?;
2777 emitted_event = Some(EmbedderEvent::MeanVecPinned {
2778 dim: u32::try_from(mean_vec.len()).unwrap_or(u32::MAX),
2779 doc_count,
2780 });
2781 }
2782
2783 tx.commit().map_err(|_| EngineError::Storage)?;
2784
2785 if let Some(ev) = emitted_event {
2786 if let Ok(mut events) = self.projection_runtime.shared.pending_events.lock() {
2787 events.push(ev);
2788 }
2789 }
2790
2791 self.next_cursor.store(cursor, Ordering::SeqCst);
2792 Ok(WriteReceipt { cursor, row_cursors: vec![cursor], dangling_edge_endpoints: 0 })
2795 }
2796
2797 #[doc(hidden)]
2802 pub fn drain_mean_centering_events_for_test(&self) -> Result<Vec<EmbedderEvent>, EngineError> {
2803 self.ensure_open()?;
2804 let mut events = self
2805 .projection_runtime
2806 .shared
2807 .pending_events
2808 .lock()
2809 .map_err(|_| EngineError::Storage)?;
2810 let out = std::mem::take(&mut *events);
2811 Ok(out)
2812 }
2813
2814 pub fn drain_embedder_events(&self) -> Result<Vec<EmbedderEvent>, EngineError> {
2822 self.ensure_open()?;
2823 let mut events = self
2824 .projection_runtime
2825 .shared
2826 .pending_events
2827 .lock()
2828 .map_err(|_| EngineError::Storage)?;
2829 Ok(std::mem::take(&mut *events))
2830 }
2831
2832 #[cfg(feature = "operator")]
2846 pub fn recompute_mean(&self) -> Result<MeanRecomputeReport, EngineError> {
2847 self.ensure_open()?;
2848 let identity = self.runtime_embedder_identity.clone();
2849 if !identity_requires_mean_centering(&identity) {
2850 return Err(EngineError::EmbedderNotConfigured);
2851 }
2852 let report = {
2853 let _gate = self
2856 .projection_runtime
2857 .shared
2858 .commit_gate
2859 .lock()
2860 .unwrap_or_else(|p| p.into_inner());
2861 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2862 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
2863 let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
2864 #[cfg(debug_assertions)]
2865 let fail = self
2866 .projection_runtime
2867 .shared
2868 .force_recompute_failure
2869 .swap(false, Ordering::SeqCst);
2870 #[cfg(not(debug_assertions))]
2871 let fail = false;
2872 let report = recompute_mean_in_tx_inner(&tx, &identity, fail)?;
2873 tx.commit().map_err(|_| EngineError::Storage)?;
2874 report
2875 };
2876 if let Ok(mut events) = self.projection_runtime.shared.pending_events.lock() {
2878 events.push(EmbedderEvent::MeanVecRecomputed {
2879 dim: report.dim,
2880 doc_count: report.doc_count_requantized,
2881 trigger: MeanRecomputeTrigger::Manual,
2882 });
2883 }
2884 Ok(report)
2885 }
2886
2887 #[doc(hidden)]
2895 pub fn set_search_limit_for_test(&self, limit: usize) {
2896 self.projection_runtime.shared.search_limit_override.store(limit, Ordering::SeqCst);
2897 }
2898
2899 #[doc(hidden)]
2903 pub fn set_recency_reweight_enabled_for_test(&self, enabled: bool) {
2904 self.projection_runtime.shared.recency_reweight_enabled.store(enabled, Ordering::SeqCst);
2905 }
2906
2907 #[doc(hidden)]
2917 pub fn set_vector_stage_only_for_test(&self, enabled: bool) {
2918 self.projection_runtime.shared.vector_stage_only_for_test.store(enabled, Ordering::SeqCst);
2919 }
2920
2921 #[doc(hidden)]
2925 #[cfg(debug_assertions)]
2926 pub fn force_next_recompute_failure_for_test(&self) {
2927 self.projection_runtime.shared.force_recompute_failure.store(true, Ordering::SeqCst);
2928 }
2929
2930 #[doc(hidden)]
2931 pub fn vector_row_count_for_test(&self) -> Result<u64, EngineError> {
2932 self.ensure_open()?;
2933 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2934 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2935 connection
2936 .query_row("SELECT COUNT(*) FROM vector_default", [], |row| row.get::<_, u64>(0))
2937 .map_err(|_| EngineError::Storage)
2938 }
2939
2940 #[doc(hidden)]
2941 pub fn read_vector_blob_for_test(&self, rowid: i64) -> Result<Vec<u8>, EngineError> {
2942 self.ensure_open()?;
2943 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2944 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2945 connection
2946 .query_row("SELECT embedding FROM vector_default WHERE rowid = ?1", [rowid], |row| {
2947 row.get::<_, Vec<u8>>(0)
2948 })
2949 .map_err(|_| EngineError::Storage)
2950 }
2951
2952 #[doc(hidden)]
2953 pub fn default_embedder_profile_for_test(&self) -> Result<EmbedderIdentity, EngineError> {
2954 self.ensure_open()?;
2955 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2956 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2957 load_default_profile(connection).map_err(|_| EngineError::Storage)
2958 }
2959
2960 #[cfg(feature = "operator")]
2964 pub fn check_integrity(
2965 &self,
2966 opts: CheckIntegrityOpts,
2967 ) -> Result<IntegrityReport, EngineError> {
2968 self.ensure_open()?;
2969 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2970 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2971 Ok(IntegrityReport {
2972 physical: physical_section(connection, opts.full),
2973 logical: logical_section(connection),
2974 semantic: semantic_section(connection),
2975 })
2976 }
2977
2978 #[cfg(feature = "operator")]
2983 pub fn safe_export(
2984 &self,
2985 out: &Path,
2986 manifest: &Path,
2987 ) -> Result<SafeExportArtifact, EngineError> {
2988 self.ensure_open()?;
2989 {
2990 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
2991 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
2992 let target = out.to_string_lossy().to_string();
2993 connection
2994 .execute("VACUUM INTO ?1", params![target])
2995 .map_err(|_| EngineError::Storage)?;
2996 }
2997 let bytes = std::fs::read(out).map_err(|_| EngineError::Storage)?;
2998 let digest = sha2::Sha256::digest(&bytes);
2999 let sha256_hex = hex_encode(digest.as_slice());
3000 let export_abs = out.canonicalize().unwrap_or_else(|_| out.to_path_buf());
3001 let manifest_json = serde_json::json!({
3002 "export_path": export_abs.to_string_lossy(),
3003 "sha256": sha256_hex,
3004 "byte_count": bytes.len() as u64,
3005 });
3006 let manifest_bytes =
3007 serde_json::to_vec_pretty(&manifest_json).map_err(|_| EngineError::Storage)?;
3008 std::fs::write(manifest, &manifest_bytes).map_err(|_| EngineError::Storage)?;
3009 Ok(SafeExportArtifact {
3010 export_path: out.to_path_buf(),
3011 manifest_path: manifest.to_path_buf(),
3012 manifest_sha256: sha256_hex,
3013 })
3014 }
3015
3016 #[cfg(feature = "operator")]
3023 pub fn rebuild_projections(&self) -> Result<RebuildReport, EngineError> {
3024 self.ensure_open()?;
3025 self.run_rebuild(true, RebuildKind::Projections)
3026 }
3027
3028 #[cfg(feature = "operator")]
3032 pub fn rebuild_vec0(&self) -> Result<RebuildReport, EngineError> {
3033 self.ensure_open()?;
3034 self.run_rebuild(false, RebuildKind::Vec0)
3035 }
3036
3037 #[cfg(feature = "operator")]
3042 pub fn trace_source_ref(&self, source_id: &str) -> Result<TraceReport, EngineError> {
3043 self.ensure_open()?;
3044 if source_id.is_empty() {
3045 return Err(EngineError::WriteValidation);
3046 }
3047 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3048 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3049
3050 let mut events: Vec<TraceEvent> = Vec::new();
3051 let mut nodes = connection
3052 .prepare(
3053 "SELECT write_cursor, kind FROM canonical_nodes WHERE source_id = ?1
3054 ORDER BY write_cursor",
3055 )
3056 .map_err(|_| EngineError::Storage)?;
3057 let node_rows = nodes
3058 .query_map([source_id], |row| {
3059 Ok(TraceEvent {
3060 write_cursor: row.get::<_, i64>(0)? as u64,
3061 kind: row.get::<_, String>(1)?,
3062 table: "canonical_nodes",
3063 })
3064 })
3065 .map_err(|_| EngineError::Storage)?;
3066 for row in node_rows {
3067 events.push(row.map_err(|_| EngineError::Storage)?);
3068 }
3069
3070 let mut edges = connection
3071 .prepare(
3072 "SELECT write_cursor, kind FROM canonical_edges WHERE source_id = ?1
3073 ORDER BY write_cursor",
3074 )
3075 .map_err(|_| EngineError::Storage)?;
3076 let edge_rows = edges
3077 .query_map([source_id], |row| {
3078 Ok(TraceEvent {
3079 write_cursor: row.get::<_, i64>(0)? as u64,
3080 kind: row.get::<_, String>(1)?,
3081 table: "canonical_edges",
3082 })
3083 })
3084 .map_err(|_| EngineError::Storage)?;
3085 for row in edge_rows {
3086 events.push(row.map_err(|_| EngineError::Storage)?);
3087 }
3088
3089 events.sort_by_key(|e| e.write_cursor);
3090 Ok(TraceReport { source_ref: source_id.to_string(), events })
3091 }
3092
3093 #[cfg(feature = "operator")]
3103 pub fn excise_source(&self, source_id: &str) -> Result<ExciseReport, EngineError> {
3104 self.ensure_open()?;
3105 if source_id.is_empty() {
3106 return Err(EngineError::WriteValidation);
3107 }
3108
3109 self.projection_runtime.set_frozen(true);
3116 let drain_result = self.drain(REBUILD_DRAIN_TIMEOUT_MS);
3117 let outcome = drain_result.and_then(|()| self.excise_source_inner(source_id));
3118 self.projection_runtime.set_frozen(false);
3119 outcome
3120 }
3121
3122 #[cfg(feature = "operator")]
3126 pub fn verify_embedder(
3127 &self,
3128 supplied_identity: &str,
3129 supplied_dimension: u32,
3130 ) -> Result<VerifyEmbedderReport, EngineError> {
3131 self.ensure_open()?;
3132 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3133 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3134 let stored = load_default_profile(connection).map_err(|_| EngineError::Storage)?;
3135 let stored_identity = format!("{}:{}", stored.name, stored.revision);
3136 let identity_match = stored_identity == supplied_identity;
3137 let dimension_match = stored.dimension == supplied_dimension;
3138 let status = match (identity_match, dimension_match) {
3139 (true, true) => VerifyEmbedderStatus::Match,
3140 (false, true) => VerifyEmbedderStatus::IdentityMismatch,
3141 (true, false) => VerifyEmbedderStatus::DimensionMismatch,
3142 (false, false) => VerifyEmbedderStatus::BothMismatch,
3143 };
3144 Ok(VerifyEmbedderReport {
3145 stored_identity,
3146 stored_dimension: stored.dimension,
3147 supplied_identity: supplied_identity.to_string(),
3148 supplied_dimension,
3149 status,
3150 })
3151 }
3152
3153 #[cfg(feature = "operator")]
3158 pub fn dump_schema(&self) -> Result<DumpSchemaReport, EngineError> {
3159 self.ensure_open()?;
3160 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3161 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3162 let user_version: u32 = connection
3163 .query_row("PRAGMA user_version", [], |row| row.get(0))
3164 .map_err(|_| EngineError::Storage)?;
3165 let tables = read_schema_objects(connection, "table")?;
3166 let indexes = read_schema_objects(connection, "index")?;
3167 Ok(DumpSchemaReport { user_version, tables: order_canonical_first(tables), indexes })
3168 }
3169
3170 #[cfg(feature = "operator")]
3173 pub fn dump_row_counts(&self) -> Result<DumpRowCountsReport, EngineError> {
3174 self.ensure_open()?;
3175 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3176 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3177 let mut counts = Vec::with_capacity(CANONICAL_TABLES.len());
3178 for name in CANONICAL_TABLES {
3179 let rows: u64 = connection
3180 .query_row(&format!("SELECT COUNT(*) FROM {name}"), [], |row| row.get(0))
3181 .map_err(|_| EngineError::Storage)?;
3182 counts.push(TableRowCount { name: (*name).to_string(), rows });
3183 }
3184 Ok(DumpRowCountsReport { counts })
3185 }
3186
3187 #[cfg(feature = "operator")]
3191 pub fn dump_profile(&self) -> Result<DumpProfileReport, EngineError> {
3192 self.ensure_open()?;
3193 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3194 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3195 let stored = load_default_profile(connection).map_err(|_| EngineError::Storage)?;
3196 let mut stmt = connection
3197 .prepare("SELECT kind FROM _fathomdb_vector_kinds ORDER BY kind")
3198 .map_err(|_| EngineError::Storage)?;
3199 let rows =
3200 stmt.query_map([], |row| row.get::<_, String>(0)).map_err(|_| EngineError::Storage)?;
3201 let mut vectorized_kinds = Vec::new();
3202 for row in rows {
3203 vectorized_kinds.push(row.map_err(|_| EngineError::Storage)?);
3204 }
3205 Ok(DumpProfileReport {
3206 embedder_identity: format!("{}:{}", stored.name, stored.revision),
3207 embedder_dimension: stored.dimension,
3208 vectorized_kinds,
3209 })
3210 }
3211
3212 #[cfg(feature = "operator")]
3218 pub fn truncate_wal(&self) -> Result<TruncateWalReport, EngineError> {
3219 self.ensure_open()?;
3220 let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3221 let connection = connection.as_ref().ok_or(EngineError::Closing)?;
3222 let (busy, log_frames, checkpointed_frames): (i64, i64, i64) = connection
3223 .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
3224 Ok((row.get(0)?, row.get(1)?, row.get(2)?))
3225 })
3226 .map_err(|_| EngineError::Storage)?;
3227 let status = if busy == 0 { TruncateWalStatus::Done } else { TruncateWalStatus::Busy };
3228 Ok(TruncateWalReport {
3229 status,
3230 busy: busy.max(0) as u32,
3231 log_frames: log_frames.max(0) as u32,
3232 checkpointed_frames: checkpointed_frames.max(0) as u32,
3233 })
3234 }
3235
3236 #[cfg(feature = "operator")]
3237 fn excise_source_inner(&self, source_id: &str) -> Result<ExciseReport, EngineError> {
3238 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3239 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
3240 let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
3241
3242 let node_cursors: Vec<i64> = {
3245 let mut stmt = tx
3246 .prepare("SELECT write_cursor FROM canonical_nodes WHERE source_id = ?1")
3247 .map_err(|_| EngineError::Storage)?;
3248 let rows = stmt
3249 .query_map([source_id], |row| row.get::<_, i64>(0))
3250 .map_err(|_| EngineError::Storage)?;
3251 rows.collect::<rusqlite::Result<Vec<_>>>().map_err(|_| EngineError::Storage)?
3252 };
3253 let edge_cursors: Vec<i64> = {
3254 let mut stmt = tx
3255 .prepare("SELECT write_cursor FROM canonical_edges WHERE source_id = ?1")
3256 .map_err(|_| EngineError::Storage)?;
3257 let rows = stmt
3258 .query_map([source_id], |row| row.get::<_, i64>(0))
3259 .map_err(|_| EngineError::Storage)?;
3260 rows.collect::<rusqlite::Result<Vec<_>>>().map_err(|_| EngineError::Storage)?
3261 };
3262
3263 let mut shadow_invalidated: u64 = 0;
3264 for cursor in node_cursors.iter().chain(edge_cursors.iter()) {
3265 shadow_invalidated = shadow_invalidated.saturating_add(
3266 tx.execute("DELETE FROM search_index WHERE write_cursor = ?1", [cursor])
3267 .map_err(|_| EngineError::Storage)? as u64,
3268 );
3269 shadow_invalidated = shadow_invalidated.saturating_add(
3272 tx.execute("DELETE FROM vector_default WHERE rowid = ?1", [cursor])
3273 .map_err(|_| EngineError::Storage)? as u64,
3274 );
3275 shadow_invalidated = shadow_invalidated.saturating_add(
3276 tx.execute("DELETE FROM _fathomdb_vector_rows WHERE write_cursor = ?1", [cursor])
3277 .map_err(|_| EngineError::Storage)? as u64,
3278 );
3279 shadow_invalidated = shadow_invalidated.saturating_add(
3280 tx.execute(
3281 "DELETE FROM _fathomdb_projection_terminal WHERE write_cursor = ?1",
3282 [cursor],
3283 )
3284 .map_err(|_| EngineError::Storage)? as u64,
3285 );
3286 }
3287
3288 let nodes_excised = tx
3289 .execute("DELETE FROM canonical_nodes WHERE source_id = ?1", [source_id])
3290 .map_err(|_| EngineError::Storage)? as u64;
3291 let edges_excised = tx
3292 .execute("DELETE FROM canonical_edges WHERE source_id = ?1", [source_id])
3293 .map_err(|_| EngineError::Storage)? as u64;
3294
3295 let excised_at = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
3302 let payload = serde_json::json!({
3303 "source_id": source_id,
3304 "excised_at": excised_at,
3305 "nodes_excised": nodes_excised,
3306 "edges_excised": edges_excised,
3307 "projections_invalidated": shadow_invalidated,
3308 })
3309 .to_string();
3310 let audit_cursor = self.next_cursor.load(Ordering::SeqCst).saturating_add(1);
3311 tx.execute(
3312 "INSERT INTO operational_mutations(
3313 collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
3314 ) VALUES('excise_source_audit', ?1, 'append', ?2, NULL, ?3)",
3315 params![source_id, payload, audit_cursor],
3316 )
3317 .map_err(|_| EngineError::Storage)?;
3318
3319 tx.commit().map_err(|_| EngineError::Storage)?;
3320 self.next_cursor.store(audit_cursor, Ordering::SeqCst);
3321 Ok(ExciseReport {
3322 source_ref: source_id.to_string(),
3323 nodes_excised,
3324 edges_excised,
3325 projections_invalidated: shadow_invalidated,
3326 })
3327 }
3328
3329 #[cfg(feature = "operator")]
3330 fn run_rebuild(
3331 &self,
3332 include_fts: bool,
3333 kind: RebuildKind,
3334 ) -> Result<RebuildReport, EngineError> {
3335 self.projection_runtime.set_frozen(true);
3336 let drain_result = self.drain(REBUILD_DRAIN_TIMEOUT_MS);
3343 let result = drain_result.and_then(|()| self.rebuild_shadow_state(include_fts, kind));
3344 self.projection_runtime.set_frozen(false);
3345 result
3346 }
3347
3348 #[cfg(feature = "operator")]
3349 fn rebuild_shadow_state(
3350 &self,
3351 include_fts: bool,
3352 kind: RebuildKind,
3353 ) -> Result<RebuildReport, EngineError> {
3354 let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
3355 let connection = connection.as_mut().ok_or(EngineError::Closing)?;
3356 let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
3357 let mut rows_invalidated: u64 = 0;
3358 if include_fts {
3359 let n = tx.execute("DELETE FROM search_index", []).map_err(|_| EngineError::Storage)?;
3360 rows_invalidated = rows_invalidated.saturating_add(n as u64);
3361 }
3362 let n = tx.execute("DELETE FROM vector_default", []).map_err(|_| EngineError::Storage)?;
3363 rows_invalidated = rows_invalidated.saturating_add(n as u64);
3364 let n = tx
3365 .execute("DELETE FROM _fathomdb_vector_rows", [])
3366 .map_err(|_| EngineError::Storage)?;
3367 rows_invalidated = rows_invalidated.saturating_add(n as u64);
3368 let n = tx
3369 .execute("DELETE FROM _fathomdb_projection_terminal", [])
3370 .map_err(|_| EngineError::Storage)?;
3371 rows_invalidated = rows_invalidated.saturating_add(n as u64);
3372 store_projection_cursor(&tx, 0).map_err(|_| EngineError::Storage)?;
3373 let mut rows_rebuilt: u64 = 0;
3374 if include_fts {
3375 for row in canonical_node_rows(&tx).map_err(|_| EngineError::Storage)? {
3376 tx.execute(
3377 "INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)",
3378 params![row.body, row.kind, row.cursor],
3379 )
3380 .map_err(|_| EngineError::Storage)?;
3381 rows_rebuilt = rows_rebuilt.saturating_add(1);
3382 }
3383 }
3384 let projection_cursor_after =
3385 load_projection_cursor(&tx).map_err(|_| EngineError::Storage)?;
3386 tx.commit().map_err(|_| EngineError::Storage)?;
3387 Ok(RebuildReport { kind, rows_invalidated, rows_rebuilt, projection_cursor_after })
3388 }
3389
3390 fn ensure_open(&self) -> Result<(), EngineError> {
3391 if self.closed.load(Ordering::SeqCst) {
3392 return Err(EngineError::Closing);
3393 }
3394
3395 Ok(())
3396 }
3397}
3398
3399fn batch_is_admin(batch: &[PreparedWrite]) -> bool {
3400 !batch.is_empty() && batch.iter().all(|w| matches!(w, PreparedWrite::AdminSchema { .. }))
3401}
3402
3403pub const TOP_K_BIT_CANDIDATES: usize = 192;
3412
3413pub const MEAN_VEC_PIN_THRESHOLD: u64 = 256;
3419
3420pub const SEARCH_RERANK_LIMIT: usize = 10;
3425
3426#[derive(Clone, Debug)]
3431struct MeanAccumulator {
3432 sum: Vec<f64>,
3433 count: u64,
3434}
3435
3436impl MeanAccumulator {
3437 fn new(dim: usize) -> Self {
3438 Self { sum: vec![0.0; dim], count: 0 }
3439 }
3440
3441 fn add(&mut self, v: &[f32]) {
3442 debug_assert_eq!(v.len(), self.sum.len(), "accumulator dim mismatch");
3443 for (slot, value) in self.sum.iter_mut().zip(v.iter()) {
3444 *slot += f64::from(*value);
3445 }
3446 self.count = self.count.saturating_add(1);
3447 }
3448
3449 fn materialize(&self) -> Vec<f32> {
3450 if self.count == 0 {
3451 return vec![0.0; self.sum.len()];
3452 }
3453 let denom = self.count as f64;
3454 self.sum.iter().map(|s| (s / denom) as f32).collect()
3455 }
3456
3457 fn count(&self) -> u64 {
3458 self.count
3459 }
3460}
3461
3462fn cosine_similarity(a: &[f32], b: &[f32]) -> f32 {
3466 if a.len() != b.len() {
3467 return 1.0;
3468 }
3469 let mut dot = 0.0f64;
3470 let mut na = 0.0f64;
3471 let mut nb = 0.0f64;
3472 for (x, y) in a.iter().zip(b.iter()) {
3473 dot += f64::from(*x) * f64::from(*y);
3474 na += f64::from(*x) * f64::from(*x);
3475 nb += f64::from(*y) * f64::from(*y);
3476 }
3477 if na == 0.0 || nb == 0.0 {
3478 return 1.0;
3479 }
3480 (dot / (na.sqrt() * nb.sqrt())) as f32
3481}
3482
3483fn run_pin_and_requantize_pass(
3491 tx: &rusqlite::Transaction<'_>,
3492 rows: &[(i64, Vec<u8>)],
3493 mean: &[f32],
3494) -> Result<(u64, Vec<EmbedderEvent>), EngineError> {
3495 let mut updated: u64 = 0;
3496 let dim = mean.len();
3497 for (rowid, blob) in rows {
3506 if blob.len() != dim * 4 {
3507 return Err(EngineError::Storage);
3508 }
3509 let un_centered = decode_vector_blob(blob);
3510 let centered = subtract_mean(&un_centered, mean);
3511 let centered_blob = encode_vector_blob(¢ered);
3512
3513 let (source_type, kind, created_at): (String, String, i64) = tx
3514 .query_row(
3515 "SELECT source_type, kind, created_at FROM vector_default WHERE rowid = ?1",
3516 params![rowid],
3517 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
3518 )
3519 .map_err(|_| EngineError::Storage)?;
3520
3521 tx.execute("DELETE FROM vector_default WHERE rowid = ?1", params![rowid])
3522 .map_err(|_| EngineError::Storage)?;
3523
3524 tx.execute(
3525 "INSERT INTO vector_default(
3530 rowid, embedding, embedding_bin, source_type, kind, created_at, status
3531 ) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
3532 params![rowid, blob, centered_blob, source_type, kind, created_at],
3533 )
3534 .map_err(|_| EngineError::Storage)?;
3535
3536 updated = updated.saturating_add(1);
3537 }
3538 let events = vec![EmbedderEvent::MeanVecPinned {
3539 dim: u32::try_from(dim).unwrap_or(u32::MAX),
3540 doc_count: updated,
3541 }];
3542 Ok((updated, events))
3543}
3544
3545fn run_requantize_pass(rows: &[(i64, Vec<u8>)], mean: &[f32]) -> (u64, Vec<EmbedderEvent>) {
3549 let mut updated: u64 = 0;
3550 let dim = mean.len();
3551 for (_rowid, blob) in rows {
3552 if blob.len() != dim * 4 {
3553 continue;
3554 }
3555 updated = updated.saturating_add(1);
3556 }
3557 let events = vec![EmbedderEvent::MeanVecPinned {
3558 dim: u32::try_from(dim).unwrap_or(u32::MAX),
3559 doc_count: updated,
3560 }];
3561 (updated, events)
3562}
3563
3564#[doc(hidden)]
3568pub mod mean_centering_internals_for_test {
3569 use super::{EmbedderEvent, MeanAccumulator};
3570
3571 pub struct AccumulatorHandle(MeanAccumulator);
3572
3573 #[must_use]
3574 pub fn new_mean_accumulator(dim: usize) -> AccumulatorHandle {
3575 AccumulatorHandle(MeanAccumulator::new(dim))
3576 }
3577
3578 pub fn accumulator_add(handle: &mut AccumulatorHandle, v: &[f32]) {
3579 handle.0.add(v);
3580 }
3581
3582 #[must_use]
3583 pub fn accumulator_materialize(handle: &AccumulatorHandle) -> Vec<f32> {
3584 handle.0.materialize()
3585 }
3586
3587 #[must_use]
3588 pub fn accumulator_count(handle: &AccumulatorHandle) -> u64 {
3589 handle.0.count()
3590 }
3591
3592 #[must_use]
3593 pub fn run_requantize_pass(rows: &[(i64, Vec<u8>)], mean: &[f32]) -> (u64, Vec<EmbedderEvent>) {
3594 super::run_requantize_pass(rows, mean)
3595 }
3596}
3597
3598pub const RRF_K: f64 = 60.0;
3601
3602pub const RECENCY_WEIGHT: f64 = 0.5 / RRF_K;
3606
3607#[doc(hidden)]
3619#[must_use]
3620pub fn fuse_rrf(vector_hits: Vec<SearchHit>, text_hits: Vec<SearchHit>) -> Vec<SearchHit> {
3621 struct Entry {
3622 hit: SearchHit,
3623 score: f64,
3624 in_vector: bool,
3625 order: usize,
3626 }
3627 let mut entries: Vec<Entry> = Vec::new();
3628 let mut accumulate = |hit: SearchHit, rank0: usize, in_vector: bool| {
3629 let contrib = 1.0 / (RRF_K + (rank0 as f64 + 1.0));
3630 if let Some(existing) = entries.iter_mut().find(|e| e.hit.body == hit.body) {
3631 existing.score += contrib;
3633 } else {
3634 let order = entries.len();
3635 entries.push(Entry { hit, score: contrib, in_vector, order });
3636 }
3637 };
3638 for (rank0, hit) in vector_hits.into_iter().enumerate() {
3639 accumulate(hit, rank0, true);
3640 }
3641 for (rank0, hit) in text_hits.into_iter().enumerate() {
3642 accumulate(hit, rank0, false);
3643 }
3644 entries.sort_by(|a, b| {
3645 b.score
3646 .partial_cmp(&a.score)
3647 .unwrap_or(std::cmp::Ordering::Equal)
3648 .then_with(|| b.in_vector.cmp(&a.in_vector))
3650 .then_with(|| a.order.cmp(&b.order))
3651 });
3652 entries
3653 .into_iter()
3654 .map(|mut e| {
3655 e.hit.score = e.score;
3656 e.hit
3657 })
3658 .collect()
3659}
3660
3661#[doc(hidden)]
3665#[must_use]
3666pub fn apply_recency_reweight(hits: Vec<SearchHit>, enabled: bool) -> Vec<SearchHit> {
3667 if !enabled || hits.len() < 2 {
3668 return hits;
3669 }
3670 let min_id = hits.iter().map(|h| h.id).min().unwrap_or(0);
3671 let max_id = hits.iter().map(|h| h.id).max().unwrap_or(0);
3672 if max_id == min_id {
3673 return hits;
3674 }
3675 let span = (max_id - min_id) as f64;
3676 let mut reweighted: Vec<SearchHit> = hits
3677 .into_iter()
3678 .map(|mut h| {
3679 let norm = (h.id - min_id) as f64 / span;
3680 h.score += RECENCY_WEIGHT * norm;
3681 h
3682 })
3683 .collect();
3684 reweighted.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal));
3686 reweighted
3687}
3688
3689#[doc(hidden)]
3693#[must_use]
3694pub fn rerank_fused(hits: Vec<SearchHit>) -> Vec<SearchHit> {
3695 hits
3696}
3697
3698fn vector_filter_clause(filter: Option<&SearchFilter>) -> String {
3704 let Some(filter) = filter else {
3705 return String::new();
3706 };
3707 if filter.is_unfiltered() {
3708 return String::new();
3709 }
3710 let mut cols: Vec<(&str, &str)> = Vec::new();
3711 if filter.source_type.is_some() {
3712 cols.push(("source_type", "="));
3713 }
3714 if filter.kind.is_some() {
3715 cols.push(("kind", "="));
3716 }
3717 if filter.created_after.is_some() {
3718 cols.push(("created_at", ">="));
3719 }
3720 if filter.status.is_some() {
3721 cols.push(("status", "="));
3722 }
3723 let mut clause = String::new();
3724 for (i, (col, op)) in cols.iter().enumerate() {
3725 clause.push_str(&format!(" AND {col}{op}?{}", i + 3));
3726 }
3727 clause
3728}
3729
3730fn vector_filter_values(filter: Option<&SearchFilter>) -> Vec<rusqlite::types::Value> {
3734 use rusqlite::types::Value;
3735 let mut out = Vec::new();
3736 let Some(filter) = filter else {
3737 return out;
3738 };
3739 if filter.is_unfiltered() {
3740 return out;
3741 }
3742 if let Some(s) = &filter.source_type {
3743 out.push(Value::Text(s.clone()));
3744 }
3745 if let Some(s) = &filter.kind {
3746 out.push(Value::Text(s.clone()));
3747 }
3748 if let Some(c) = filter.created_after {
3749 out.push(Value::Integer(c));
3750 }
3751 if let Some(s) = &filter.status {
3752 out.push(Value::Text(s.clone()));
3753 }
3754 out
3755}
3756
3757fn build_vector_phase1_sql(filter: Option<&SearchFilter>, final_limit: usize) -> String {
3763 let filter_clause = vector_filter_clause(filter);
3764 format!(
3765 "WITH candidates AS (
3766 SELECT rowid
3767 FROM vector_default
3768 WHERE embedding_bin MATCH vec_quantize_binary(vec_f32(?1)){filter_clause}
3769 ORDER BY distance
3770 LIMIT {top_k}
3771 )
3772 SELECT c.rowid, vec_distance_l2(v.embedding, vec_f32(?2)) AS l2
3773 FROM candidates c
3774 JOIN vector_default v ON v.rowid = c.rowid
3775 ORDER BY l2
3776 LIMIT {final_limit}",
3777 top_k = TOP_K_BIT_CANDIDATES,
3778 )
3779}
3780
3781#[doc(hidden)]
3785#[must_use]
3786pub fn vector_phase1_sql_for_test(filter: Option<&SearchFilter>) -> String {
3787 build_vector_phase1_sql(filter, SEARCH_RERANK_LIMIT)
3788}
3789
3790fn text_hit_passes_filter(
3798 tx: &rusqlite::Transaction<'_>,
3799 id: u64,
3800 kind: &str,
3801 filter: Option<&SearchFilter>,
3802) -> rusqlite::Result<bool> {
3803 let Some(filter) = filter else {
3804 return Ok(true);
3805 };
3806 if filter.is_unfiltered() {
3807 return Ok(true);
3808 }
3809 if let Some(k) = &filter.kind {
3810 if kind != k {
3811 return Ok(false);
3812 }
3813 }
3814 if let Some(st) = &filter.source_type {
3815 match resolve_source_type(kind) {
3816 Ok(resolved) if resolved == st.as_str() => {}
3817 _ => return Ok(false),
3818 }
3819 }
3820 if filter.created_after.is_some() || filter.status.is_some() {
3821 let meta: Option<(i64, Option<String>)> = tx
3822 .query_row(
3823 "SELECT created_at, status FROM vector_default WHERE rowid = ?1 LIMIT 1",
3824 [id as i64],
3825 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?)),
3826 )
3827 .optional()?;
3828 let Some((created_at, status)) = meta else {
3829 return Ok(false);
3831 };
3832 if let Some(bound) = filter.created_after {
3833 if created_at < bound {
3834 return Ok(false);
3835 }
3836 }
3837 if let Some(want) = &filter.status {
3838 if status.as_deref() != Some(want.as_str()) {
3839 return Ok(false);
3840 }
3841 }
3842 }
3843 Ok(true)
3844}
3845
3846#[allow(clippy::too_many_arguments)]
3852fn read_search_in_tx(
3853 reader: &mut Connection,
3854 compiled: &fathomdb_query::CompiledQuery,
3855 query_vector: Option<&str>,
3856 query_vector_bin: Option<&str>,
3857 final_limit: usize,
3858 filter: Option<&SearchFilter>,
3859 recency_enabled: bool,
3860 vector_stage_only: bool,
3861) -> rusqlite::Result<(u64, Option<SoftFallback>, Vec<SearchHit>)> {
3862 let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
3863 let cursor = load_projection_cursor(&tx)?;
3864 let vector_results = if let Some(query_vector) = query_vector {
3865 let mut rowids = Vec::new();
3866 let bin_vector = query_vector_bin.unwrap_or(query_vector);
3867 {
3868 let sql = build_vector_phase1_sql(filter, final_limit);
3888 let mut params: Vec<rusqlite::types::Value> = vec![
3889 rusqlite::types::Value::Text(bin_vector.to_string()),
3890 rusqlite::types::Value::Text(query_vector.to_string()),
3891 ];
3892 params.extend(vector_filter_values(filter));
3893 let mut statement = tx.prepare(&sql)?;
3894 let rows = statement.query_map(rusqlite::params_from_iter(params.iter()), |row| {
3895 Ok((row.get::<_, i64>(0)?, row.get::<_, f64>(1)?))
3896 })?;
3897 for row in rows.flatten() {
3898 rowids.push(row);
3899 }
3900 }
3901 let mut results = Vec::new();
3906 let mut statement =
3907 tx.prepare("SELECT kind, body FROM canonical_nodes WHERE write_cursor = ?1 LIMIT 1")?;
3908 for (rowid, score) in rowids {
3909 if let Ok((kind, body)) = statement
3910 .query_row([rowid], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)))
3911 {
3912 results.push(SearchHit {
3913 id: rowid as u64,
3914 kind,
3915 body,
3916 score,
3917 branch: SoftFallbackBranch::Vector,
3918 });
3919 }
3920 }
3921 results
3922 } else {
3923 Vec::new()
3924 };
3925 let vector_rows_visible = !vector_results.is_empty();
3926 let soft_fallback = if query_vector.is_some() && !vector_rows_visible {
3927 tx.query_row(
3928 "SELECT 1
3929 FROM search_index
3930 JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = search_index.kind
3931 LEFT JOIN _fathomdb_projection_terminal
3932 ON _fathomdb_projection_terminal.write_cursor = search_index.write_cursor
3933 WHERE search_index MATCH ?1
3934 AND _fathomdb_projection_terminal.write_cursor IS NULL
3935 LIMIT 1",
3936 [compiled.match_expression.as_str()],
3937 |_row| Ok(SoftFallback { branch: SoftFallbackBranch::Vector }),
3938 )
3939 .ok()
3940 } else {
3941 None
3942 };
3943 let text_candidates: Vec<SearchHit> = {
3948 let perf_limit: Option<usize> = if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_some() {
3955 std::env::var("FATHOMDB_PERF_SEARCH_LIMIT").ok().and_then(|s| s.parse().ok())
3956 } else {
3957 None
3958 };
3959 let sql = match perf_limit {
3964 Some(k) => format!(
3965 "SELECT body, kind, write_cursor, bm25(search_index) FROM search_index \
3966 WHERE search_index MATCH ?1 ORDER BY write_cursor LIMIT {k}"
3967 ),
3968 None => "SELECT body, kind, write_cursor, bm25(search_index) FROM search_index \
3969 WHERE search_index MATCH ?1 ORDER BY write_cursor"
3970 .to_string(),
3971 };
3972 let mut statement = tx.prepare(&sql)?;
3973 let rows = statement.query_map([compiled.match_expression.as_str()], |row| {
3974 Ok(SearchHit {
3975 body: row.get::<_, String>(0)?,
3976 kind: row.get::<_, String>(1)?,
3977 id: row.get::<_, i64>(2)? as u64,
3978 score: row.get::<_, f64>(3)?,
3979 branch: SoftFallbackBranch::Text,
3980 })
3981 })?;
3982 rows.flatten().collect()
3983 };
3984 let mut text_results: Vec<SearchHit> = Vec::with_capacity(text_candidates.len());
3985 for hit in text_candidates {
3986 if text_hit_passes_filter(&tx, hit.id, &hit.kind, filter)? {
3987 text_results.push(hit);
3988 }
3989 }
3990 tx.commit()?;
3991
3992 let results = if vector_stage_only {
4002 vector_results
4003 } else {
4004 rerank_fused(apply_recency_reweight(
4010 fuse_rrf(vector_results, text_results),
4011 recency_enabled,
4012 ))
4013 };
4014 Ok((cursor, soft_fallback, results))
4015}
4016
4017const READ_COLLECTION_MAX_LIMIT: usize = 1_000_000;
4022
4023fn read_get_by_id_in_tx(
4029 reader: &mut Connection,
4030 logical_ids: &[String],
4031) -> rusqlite::Result<Vec<Option<NodeRecord>>> {
4032 if logical_ids.is_empty() {
4033 return Ok(Vec::new());
4034 }
4035 let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
4036 let mut found: HashMap<String, NodeRecord> = HashMap::new();
4039 {
4040 let unique: Vec<&String> = {
4041 let mut seen = std::collections::HashSet::new();
4042 logical_ids.iter().filter(|id| seen.insert((*id).clone())).collect()
4043 };
4044 let placeholders = std::iter::repeat_n("?", unique.len()).collect::<Vec<_>>().join(", ");
4045 let sql = format!(
4046 "SELECT logical_id, kind, body, write_cursor
4047 FROM canonical_nodes
4048 WHERE logical_id IN ({placeholders}) AND superseded_at IS NULL"
4049 );
4050 let mut statement = tx.prepare(&sql)?;
4051 let params = rusqlite::params_from_iter(unique.iter().map(|s| s.as_str()));
4052 let rows = statement.query_map(params, |row| {
4053 let logical_id: String = row.get(0)?;
4054 Ok(NodeRecord {
4055 logical_id,
4056 kind: row.get(1)?,
4057 body: row.get(2)?,
4058 write_cursor: row.get::<_, i64>(3)? as u64,
4059 })
4060 })?;
4061 for row in rows {
4062 let record = row?;
4063 found.insert(record.logical_id.clone(), record);
4064 }
4065 }
4066 let out = logical_ids.iter().map(|id| found.get(id).cloned()).collect();
4068 Ok(out)
4069}
4070
4071fn read_collection_in_tx(
4088 reader: &mut Connection,
4089 collection: &str,
4090 after_id: Option<i64>,
4091 limit: usize,
4092) -> rusqlite::Result<Vec<OpStoreRow>> {
4093 if limit == 0 {
4094 return Ok(Vec::new());
4095 }
4096 let clamped = limit.min(READ_COLLECTION_MAX_LIMIT) as i64;
4097 let after = after_id.unwrap_or(0).max(0);
4102 let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
4103 let mut statement = tx.prepare(
4104 "SELECT id, collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
4105 FROM operational_mutations
4106 WHERE collection_name = ?1 AND id > ?2
4107 ORDER BY id
4108 LIMIT ?3",
4109 )?;
4110 let rows = statement.query_map(params![collection, after, clamped], |row| {
4111 Ok(OpStoreRow {
4112 id: row.get(0)?,
4113 collection: row.get(1)?,
4114 record_key: row.get(2)?,
4115 op_kind: row.get(3)?,
4116 payload: row.get(4)?,
4117 schema_id: row.get(5)?,
4118 write_cursor: row.get::<_, i64>(6)? as u64,
4119 })
4120 })?;
4121 let mut out = Vec::new();
4122 for row in rows {
4123 out.push(row?);
4124 }
4125 Ok(out)
4126}
4127
4128fn projection_dispatcher_loop(shared: Arc<ProjectionRuntimeShared>) {
4129 let connection = match open_runtime_connection(&shared.path) {
4130 Ok(connection) => connection,
4131 Err(_) => return,
4132 };
4133 loop {
4134 let in_flight = {
4135 let mut state = match shared.state.lock() {
4136 Ok(state) => state,
4137 Err(_) => return,
4138 };
4139 while !state.stopping
4140 && (!state.pending_scan
4141 || state.frozen
4142 || state.active_jobs + state.queued_jobs >= PROJECTION_INFLIGHT_LIMIT)
4143 {
4144 state = match shared.state_cvar.wait(state) {
4145 Ok(state) => state,
4146 Err(_) => return,
4147 };
4148 }
4149 if state.stopping {
4150 return;
4151 }
4152 state.pending_scan = false;
4153 state.in_flight.clone()
4154 };
4155
4156 let budget = {
4162 let state = match shared.state.lock() {
4163 Ok(state) => state,
4164 Err(_) => return,
4165 };
4166 PROJECTION_INFLIGHT_LIMIT.saturating_sub(state.active_jobs + state.queued_jobs)
4167 };
4168 let fetch_cap = budget.clamp(1, PROJECTION_SCAN_FETCH);
4169 match next_pending_projection_jobs(&connection, &in_flight, fetch_cap) {
4170 Ok(jobs) if !jobs.is_empty() => {
4171 if let Ok(mut state) = shared.state.lock() {
4172 state.queued_jobs = state.queued_jobs.saturating_add(jobs.len());
4173 for job in &jobs {
4174 state.in_flight.insert(job.cursor);
4175 }
4176 state.pending_scan = true;
4177 shared.state_cvar.notify_all();
4178 }
4179 if let Ok(mut queue) = shared.queue.lock() {
4180 for job in jobs {
4181 queue.push_back(job);
4182 }
4183 shared.queue_cvar.notify_all();
4184 }
4185 }
4186 Ok(_) => {}
4187 Err(_) => {
4188 if let Ok(mut state) = shared.state.lock() {
4189 state.pending_scan = false;
4190 shared.state_cvar.notify_all();
4191 }
4192 }
4193 }
4194 }
4195}
4196
4197fn projection_worker_loop(shared: Arc<ProjectionRuntimeShared>) {
4198 let mut connection = match open_runtime_connection(&shared.path) {
4199 Ok(connection) => connection,
4200 Err(_) => return,
4201 };
4202 if ensure_vector_partition(&mut connection, shared.embedder_identity.dimension).is_err() {
4203 return;
4204 }
4205 loop {
4206 let jobs = {
4207 let mut queue = match shared.queue.lock() {
4208 Ok(queue) => queue,
4209 Err(_) => return,
4210 };
4211 loop {
4212 let stopping = shared.state.lock().map(|state| state.stopping).unwrap_or(true);
4213 if stopping && queue.is_empty() {
4214 return;
4215 }
4216 if let Some(job) = queue.pop_front() {
4217 let mut jobs = vec![job];
4218 while jobs.len() < PROJECTION_COMMIT_BATCH {
4219 let Some(job) = queue.pop_front() else {
4220 break;
4221 };
4222 jobs.push(job);
4223 }
4224 if let Ok(mut state) = shared.state.lock() {
4225 state.queued_jobs = state.queued_jobs.saturating_sub(jobs.len());
4226 state.active_jobs = state.active_jobs.saturating_add(jobs.len());
4227 shared.state_cvar.notify_all();
4228 }
4229 break jobs;
4230 }
4231 queue = match shared.queue_cvar.wait(queue) {
4232 Ok(queue) => queue,
4233 Err(_) => return,
4234 };
4235 }
4236 };
4237
4238 let panicked = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4245 run_projection_jobs(&shared, &mut connection, &jobs);
4246 }))
4247 .is_err();
4248 if panicked {
4249 commit_projection_panic_failures(&shared, &mut connection, &jobs);
4250 }
4251
4252 if let Ok(mut state) = shared.state.lock() {
4253 state.active_jobs = state.active_jobs.saturating_sub(jobs.len());
4254 for job in &jobs {
4255 state.in_flight.remove(&job.cursor);
4256 }
4257 if !state.stopping {
4258 state.pending_scan = true;
4259 }
4260 shared.state_cvar.notify_all();
4261 }
4262 }
4263}
4264
4265enum ProjectionOutcome {
4266 Success {
4272 cursor: u64,
4273 kind: String,
4274 blob: Vec<u8>,
4275 bin_blob: Vec<u8>,
4276 },
4277 Failure {
4278 cursor: u64,
4279 failure_code: &'static str,
4280 },
4281}
4282
4283fn run_projection_jobs(
4284 shared: &ProjectionRuntimeShared,
4285 connection: &mut Connection,
4286 jobs: &[ProjectionJob],
4287) {
4288 let mut outcomes = Vec::with_capacity(jobs.len());
4289 for job in jobs {
4290 outcomes.push(run_projection_job(shared, job));
4291 }
4292 let _ = commit_projection_outcomes(connection, &outcomes, shared);
4293}
4294
4295fn commit_projection_panic_failures(
4299 shared: &ProjectionRuntimeShared,
4300 connection: &mut Connection,
4301 jobs: &[ProjectionJob],
4302) {
4303 let outcomes: Vec<ProjectionOutcome> = jobs
4304 .iter()
4305 .map(|job| ProjectionOutcome::Failure {
4306 cursor: job.cursor,
4307 failure_code: "ProjectionPanic",
4308 })
4309 .collect();
4310 let _ = commit_projection_outcomes(connection, &outcomes, shared);
4311}
4312
4313fn embed_with_watchdog(
4337 embedder: &Arc<dyn Embedder>,
4338 body: &str,
4339 timeout: Duration,
4340 live: &Arc<AtomicU64>,
4341) -> Result<Vec<f32>, RuntimeEmbedderError> {
4342 let (tx, rx) = mpsc::channel();
4343 let embedder = Arc::clone(embedder);
4344 let body = body.to_string();
4345 live.fetch_add(1, Ordering::Relaxed);
4348 let live_thread = Arc::clone(live);
4349 thread::spawn(move || {
4350 let outcome =
4351 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| embedder.embed(&body)));
4352 let _ = tx.send(outcome);
4356 live_thread.fetch_sub(1, Ordering::Relaxed);
4357 });
4358 match rx.recv_timeout(timeout) {
4359 Ok(Ok(result)) => result,
4360 Ok(Err(panic_payload)) => std::panic::resume_unwind(panic_payload),
4361 Err(mpsc::RecvTimeoutError::Timeout) => Err(RuntimeEmbedderError::Timeout),
4362 Err(mpsc::RecvTimeoutError::Disconnected) => Err(RuntimeEmbedderError::Failed {
4366 message: "embed watchdog thread dropped its result channel".to_string(),
4367 }),
4368 }
4369}
4370
4371fn run_projection_job(shared: &ProjectionRuntimeShared, job: &ProjectionJob) -> ProjectionOutcome {
4372 if shared.embed_circuit_open.load(Ordering::Relaxed) {
4379 return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: "EmbedderError" };
4380 }
4381 let delays = shared.retry_delays_ms.lock().map(|delays| delays.clone()).unwrap_or_default();
4382 let mut last_code = "EmbedderError";
4383 for (attempt, delay_ms) in std::iter::once(0_u64).chain(delays.iter().copied()).enumerate() {
4384 if attempt > 0 {
4385 if shared.state.lock().map(|state| state.stopping).unwrap_or(true) {
4386 return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code };
4387 }
4388 thread::sleep(Duration::from_millis(delay_ms));
4389 }
4390 if shared.embed_circuit_open.load(Ordering::Relaxed) {
4396 return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code };
4397 }
4398 let embed_timeout = Duration::from_millis(shared.embed_timeout_ms.load(Ordering::Relaxed));
4402 let vector = match shared.embedder.as_ref() {
4403 Some(embedder) => {
4404 let _embed_permit =
4414 shared.embed_serialize.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
4415 let threshold = shared.embed_circuit_threshold.load(Ordering::Relaxed);
4424 if shared.embed_circuit_open.load(Ordering::Relaxed)
4425 || (threshold != 0
4426 && shared.live_embed_threads.load(Ordering::Relaxed) >= threshold)
4427 {
4428 shared.embed_circuit_open.store(true, Ordering::Relaxed);
4429 return ProjectionOutcome::Failure {
4430 cursor: job.cursor,
4431 failure_code: last_code,
4432 };
4433 }
4434 match embed_with_watchdog(
4435 embedder,
4436 &job.body,
4437 embed_timeout,
4438 &shared.live_embed_threads,
4439 ) {
4440 Ok(vector) => vector,
4441 Err(RuntimeEmbedderError::Timeout) => {
4442 last_code = "EmbedderError";
4446 continue;
4447 }
4448 Err(RuntimeEmbedderError::Failed { .. }) => {
4449 last_code = "EmbedderError";
4450 continue;
4451 }
4452 }
4453 }
4454 None => {
4455 last_code = "EmbedderNotConfiguredError";
4456 continue;
4457 }
4458 };
4459
4460 if u32::try_from(vector.len()).unwrap_or(u32::MAX) != shared.embedder_identity.dimension {
4461 last_code = "EmbedderDimensionMismatchError";
4462 continue;
4463 }
4464
4465 let blob = encode_vector_blob(&vector);
4466 let bin_blob = blob.clone();
4475 return ProjectionOutcome::Success {
4476 cursor: job.cursor,
4477 kind: job.kind.clone(),
4478 blob,
4479 bin_blob,
4480 };
4481 }
4482
4483 ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code }
4484}
4485
4486fn next_pending_projection_jobs(
4487 connection: &Connection,
4488 in_flight: &BTreeSet<u64>,
4489 max_jobs: usize,
4490) -> rusqlite::Result<Vec<ProjectionJob>> {
4491 if max_jobs == 0 {
4492 return Ok(Vec::new());
4493 }
4494 let cursor = load_projection_cursor(connection)?;
4495 let sql_limit = max_jobs.saturating_add(in_flight.len()).min(256);
4498 let sql = format!(
4499 "SELECT canonical_nodes.write_cursor, canonical_nodes.kind, canonical_nodes.body
4500 FROM canonical_nodes
4501 JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = canonical_nodes.kind
4502 LEFT JOIN _fathomdb_projection_terminal
4503 ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
4504 WHERE canonical_nodes.write_cursor > ?1
4505 AND _fathomdb_projection_terminal.write_cursor IS NULL
4506 ORDER BY canonical_nodes.write_cursor
4507 LIMIT {sql_limit}"
4508 );
4509 let mut statement = connection.prepare_cached(&sql)?;
4510 let rows = statement.query_map([cursor], |row| {
4511 Ok(ProjectionJob { cursor: row.get(0)?, kind: row.get(1)?, body: row.get(2)? })
4512 })?;
4513 let mut jobs = Vec::with_capacity(max_jobs);
4514 for row in rows {
4515 let job = row?;
4516 if in_flight.contains(&job.cursor) {
4517 continue;
4518 }
4519 jobs.push(job);
4520 if jobs.len() >= max_jobs {
4521 break;
4522 }
4523 }
4524 Ok(jobs)
4525}
4526
4527fn database_has_pending_projection_work(path: &Path) -> rusqlite::Result<bool> {
4528 let connection = open_runtime_connection(path)?;
4529 let cursor = load_projection_cursor(&connection)?;
4530 connection
4531 .query_row(
4532 "SELECT 1
4533 FROM canonical_nodes
4534 JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = canonical_nodes.kind
4535 LEFT JOIN _fathomdb_projection_terminal
4536 ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
4537 WHERE canonical_nodes.write_cursor > ?1
4538 AND _fathomdb_projection_terminal.write_cursor IS NULL
4539 LIMIT 1",
4540 [cursor],
4541 |_row| Ok(true),
4542 )
4543 .or_else(|err| match err {
4544 rusqlite::Error::QueryReturnedNoRows => Ok(false),
4545 _ => Err(err),
4546 })
4547}
4548
4549struct CanonicalNodeRow {
4550 cursor: u64,
4551 kind: String,
4552 body: String,
4553}
4554
4555fn reproject_search_index_after_tokenizer_upgrade(connection: &Connection) -> rusqlite::Result<()> {
4571 let rows = canonical_node_rows(connection)?;
4572 connection.execute_batch("BEGIN IMMEDIATE")?;
4573 let result = (|| {
4574 connection.execute("DELETE FROM search_index", [])?;
4575 {
4576 let mut statement = connection
4577 .prepare("INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)")?;
4578 for row in &rows {
4579 statement.execute(params![row.body, row.kind, row.cursor])?;
4580 }
4581 }
4582 connection.execute(
4583 "INSERT INTO _fathomdb_open_state(key, value) VALUES(?1, ?2)
4584 ON CONFLICT(key) DO UPDATE SET value = excluded.value",
4585 params![SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY, "1"],
4586 )?;
4587 Ok(())
4588 })();
4589 match result {
4590 Ok(()) => connection.execute_batch("COMMIT"),
4591 Err(err) => {
4592 let _ = connection.execute_batch("ROLLBACK");
4593 Err(err)
4594 }
4595 }
4596}
4597
4598fn search_index_tokenizer_reproject_complete(connection: &Connection) -> rusqlite::Result<bool> {
4612 match connection.query_row(
4613 "SELECT value FROM _fathomdb_open_state WHERE key = ?1",
4614 [SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY],
4615 |row| row.get::<_, String>(0),
4616 ) {
4617 Ok(value) => Ok(value == "1"),
4618 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(false),
4619 Err(rusqlite::Error::SqliteFailure(_, Some(ref message)))
4620 if message.contains("no such table") =>
4621 {
4622 Ok(true)
4623 }
4624 Err(err) => Err(err),
4625 }
4626}
4627
4628fn canonical_node_rows(connection: &Connection) -> rusqlite::Result<Vec<CanonicalNodeRow>> {
4629 let mut statement = connection
4630 .prepare("SELECT write_cursor, kind, body FROM canonical_nodes ORDER BY write_cursor")?;
4631 let rows = statement.query_map([], |row| {
4632 Ok(CanonicalNodeRow {
4633 cursor: row.get::<_, u64>(0)?,
4634 kind: row.get::<_, String>(1)?,
4635 body: row.get::<_, String>(2)?,
4636 })
4637 })?;
4638 rows.collect()
4639}
4640
4641#[cfg(feature = "operator")]
4642fn hex_encode(bytes: &[u8]) -> String {
4643 let mut out = String::with_capacity(bytes.len() * 2);
4644 for byte in bytes {
4645 out.push(hex_nibble(byte >> 4));
4646 out.push(hex_nibble(byte & 0x0f));
4647 }
4648 out
4649}
4650
4651#[cfg(feature = "operator")]
4652fn hex_nibble(value: u8) -> char {
4653 match value {
4654 0..=9 => (b'0' + value) as char,
4655 10..=15 => (b'a' + value - 10) as char,
4656 _ => unreachable!(),
4657 }
4658}
4659
4660#[cfg(feature = "operator")]
4661fn physical_section(connection: &Connection, full: bool) -> Section {
4662 let mut findings = Vec::new();
4663 if let Err(err) = connection.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0)) {
4664 findings.push(Finding {
4665 code: "E_CORRUPT_HEADER",
4666 stage: "PhysicalProbe",
4667 locator: locator_from_rusqlite_error(&err),
4668 doc_anchor: "design/recovery.md#header-malformed",
4669 detail: format!("page_count probe failed: {err}"),
4670 });
4671 }
4672 if full {
4673 match collect_integrity_check_findings(connection) {
4674 Ok(rows) => findings.extend(rows),
4675 Err(err) => findings.push(Finding {
4676 code: "E_CORRUPT_INTEGRITY_CHECK",
4677 stage: "IntegrityCheck",
4678 locator: locator_from_rusqlite_error(&err),
4679 doc_anchor: "design/recovery.md#integrity-check-full-findings",
4680 detail: format!("PRAGMA integrity_check failed: {err}"),
4681 }),
4682 }
4683 }
4684 if findings.is_empty() {
4685 Section::Clean
4686 } else {
4687 Section::Findings(findings)
4688 }
4689}
4690
4691#[cfg(feature = "operator")]
4692fn logical_section(connection: &Connection) -> Section {
4693 let mut findings = Vec::new();
4694 if let Err(err) = connection.query_row("PRAGMA schema_version", [], |row| row.get::<_, i64>(0))
4695 {
4696 findings.push(Finding {
4697 code: "E_CORRUPT_SCHEMA",
4698 stage: "SchemaProbe",
4699 locator: locator_from_rusqlite_error(&err),
4700 doc_anchor: "design/recovery.md#schema-inconsistent",
4701 detail: format!("schema_version probe failed: {err}"),
4702 });
4703 }
4704 match connection.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)) {
4705 Ok(0) => findings.push(Finding {
4706 code: "E_CORRUPT_SCHEMA",
4707 stage: "SchemaProbe",
4708 locator: CorruptionLocator::MigrationStep { from: 0, to: 0 },
4709 doc_anchor: "design/recovery.md#schema-inconsistent",
4710 detail: "user_version is zero".to_string(),
4711 }),
4712 Ok(_) => {}
4713 Err(err) => findings.push(Finding {
4714 code: "E_CORRUPT_SCHEMA",
4715 stage: "SchemaProbe",
4716 locator: locator_from_rusqlite_error(&err),
4717 doc_anchor: "design/recovery.md#schema-inconsistent",
4718 detail: format!("user_version probe failed: {err}"),
4719 }),
4720 }
4721 if findings.is_empty() {
4722 Section::Clean
4723 } else {
4724 Section::Findings(findings)
4725 }
4726}
4727
4728#[cfg(feature = "operator")]
4729fn semantic_section(connection: &Connection) -> Section {
4730 match load_default_profile(connection) {
4731 Ok(_) => Section::Clean,
4732 Err(rusqlite::Error::QueryReturnedNoRows) => Section::Findings(vec![Finding {
4733 code: "E_CORRUPT_EMBEDDER_IDENTITY",
4734 stage: "EmbedderIdentity",
4735 locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
4736 doc_anchor: "design/recovery.md#embedder-identity-drift",
4737 detail: "default embedder profile row is missing".to_string(),
4738 }]),
4739 Err(err) => Section::Findings(vec![Finding {
4740 code: "E_CORRUPT_EMBEDDER_IDENTITY",
4741 stage: "EmbedderIdentity",
4742 locator: locator_from_rusqlite_error(&err),
4743 doc_anchor: "design/recovery.md#embedder-identity-drift",
4744 detail: format!("default embedder profile probe failed: {err}"),
4745 }]),
4746 }
4747}
4748
4749#[cfg(feature = "operator")]
4750fn collect_integrity_check_findings(connection: &Connection) -> rusqlite::Result<Vec<Finding>> {
4751 let mut statement = connection.prepare("PRAGMA integrity_check")?;
4752 let rows = statement.query_map([], |row| row.get::<_, String>(0))?;
4753 let mut findings = Vec::new();
4754 for row in rows {
4755 let message = row?;
4756 if message == "ok" {
4757 continue;
4758 }
4759 findings.push(Finding {
4760 code: "E_CORRUPT_INTEGRITY_CHECK",
4761 stage: "IntegrityCheck",
4762 locator: CorruptionLocator::OpaqueSqliteError {
4763 sqlite_extended_code: rusqlite::ffi::SQLITE_CORRUPT,
4764 },
4765 doc_anchor: "design/recovery.md#integrity-check-full-findings",
4766 detail: message,
4767 });
4768 }
4769 Ok(findings)
4770}
4771
4772#[cfg(feature = "operator")]
4773fn locator_from_rusqlite_error(err: &rusqlite::Error) -> CorruptionLocator {
4774 let extended = err.sqlite_error().map(|inner| inner.extended_code).unwrap_or(0);
4775 CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: extended }
4776}
4777
4778fn open_runtime_connection(path: &Path) -> rusqlite::Result<Connection> {
4779 let connection = Connection::open(path)?;
4780 connection.pragma_update(None, "journal_mode", "WAL")?;
4781 Ok(connection)
4782}
4783
4784fn load_projection_cursor(connection: &Connection) -> rusqlite::Result<u64> {
4785 connection
4786 .query_row(
4787 "SELECT value FROM _fathomdb_open_state WHERE key = ?1",
4788 [PROJECTION_CURSOR_KEY],
4789 |row| row.get::<_, String>(0),
4790 )
4791 .map(|value| value.parse::<u64>().unwrap_or(0))
4792 .or_else(|err| match err {
4793 rusqlite::Error::QueryReturnedNoRows => Ok(0),
4794 _ => Err(err),
4795 })
4796}
4797
4798fn store_projection_cursor(connection: &Connection, cursor: u64) -> rusqlite::Result<()> {
4799 connection.execute(
4800 "INSERT INTO _fathomdb_open_state(key, value) VALUES(?1, ?2)
4801 ON CONFLICT(key) DO UPDATE SET value = excluded.value",
4802 params![PROJECTION_CURSOR_KEY, cursor.to_string()],
4803 )?;
4804 Ok(())
4805}
4806
4807fn record_projection_terminal(
4808 connection: &Connection,
4809 cursor: u64,
4810 state: &str,
4811) -> rusqlite::Result<()> {
4812 connection.execute(
4813 "INSERT OR IGNORE INTO _fathomdb_projection_terminal(write_cursor, state) VALUES(?1, ?2)",
4814 params![cursor, state],
4815 )?;
4816 Ok(())
4817}
4818
4819fn terminal_state_for_cursor(
4820 connection: &Connection,
4821 cursor: u64,
4822) -> rusqlite::Result<Option<String>> {
4823 connection
4824 .query_row(
4825 "SELECT state FROM _fathomdb_projection_terminal WHERE write_cursor = ?1",
4826 [cursor],
4827 |row| row.get::<_, String>(0),
4828 )
4829 .map(Some)
4830 .or_else(|err| match err {
4831 rusqlite::Error::QueryReturnedNoRows => Ok(None),
4832 _ => Err(err),
4833 })
4834}
4835
4836fn advance_projection_cursor(connection: &Connection) -> rusqlite::Result<u64> {
4837 let mut cursor = load_projection_cursor(connection)?;
4838 loop {
4839 let next = cursor.saturating_add(1);
4840 if terminal_state_for_cursor(connection, next)?.is_some() {
4841 cursor = next;
4842 } else {
4843 break;
4844 }
4845 }
4846 store_projection_cursor(connection, cursor)?;
4847 Ok(cursor)
4848}
4849
4850fn commit_projection_outcomes(
4851 connection: &mut Connection,
4852 outcomes: &[ProjectionOutcome],
4853 shared: &ProjectionRuntimeShared,
4854) -> rusqlite::Result<()> {
4855 let embedder_identity = &shared.embedder_identity;
4856 let mc = identity_requires_mean_centering(embedder_identity);
4857 let _gate = shared.commit_gate.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
4860 let tx = connection.transaction()?;
4861 let mut current_mean: Option<Vec<f32>> = if mc {
4864 tx.query_row(
4865 "SELECT mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
4866 [],
4867 |row| row.get::<_, Option<Vec<u8>>>(0),
4868 )
4869 .ok()
4870 .flatten()
4871 .map(|bytes| decode_vector_blob(&bytes))
4872 } else {
4873 None
4874 };
4875 let mut staged_events: Vec<EmbedderEvent> = Vec::new();
4876 for outcome in outcomes {
4877 match outcome {
4878 ProjectionOutcome::Success { cursor, kind, blob, bin_blob } => {
4879 if terminal_state_for_cursor(&tx, *cursor)?.is_some() {
4880 continue;
4881 }
4882 let pin_mean: Option<Vec<f32>> = if mc && current_mean.is_none() {
4887 let mut acc = shared.mean_accumulator.lock().unwrap_or_else(|p| p.into_inner());
4888 match acc.as_mut() {
4889 Some(a) => {
4890 a.add(&decode_vector_blob(bin_blob));
4891 if a.count() >= MEAN_VEC_PIN_THRESHOLD {
4892 let mean = a.materialize();
4893 *acc = None;
4894 Some(mean)
4895 } else {
4896 None
4897 }
4898 }
4899 None => None,
4900 }
4901 } else {
4902 None
4903 };
4904
4905 let source_type = resolve_source_type(kind).map_err(|_| {
4906 rusqlite::Error::SqliteFailure(
4907 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
4908 Some(format!("unknown kind for source_type mapping: {kind}")),
4909 )
4910 })?;
4911 let now_unix =
4912 SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs()
4913 as i64;
4914 tx.execute(
4915 "INSERT OR IGNORE INTO _fathomdb_vector_rows(rowid, kind, write_cursor) VALUES(?1, ?2, ?3)",
4916 params![cursor, kind, cursor],
4917 )?;
4918 let centered_blob: Vec<u8> = match ¤t_mean {
4924 Some(mean) if mean.len() * 4 == bin_blob.len() => {
4925 encode_vector_blob(&subtract_mean(&decode_vector_blob(bin_blob), mean))
4926 }
4927 _ => bin_blob.clone(),
4928 };
4929 tx.execute(
4930 "INSERT OR IGNORE INTO vector_default(
4934 rowid, embedding, embedding_bin, source_type, kind, created_at, status
4935 ) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
4936 params![cursor, blob, centered_blob, source_type, kind, now_unix],
4937 )?;
4938 record_projection_terminal(&tx, *cursor, "up_to_date")?;
4939
4940 if let Some(mean) = pin_mean {
4945 tx.execute(
4946 "UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
4947 params![encode_vector_blob(&mean)],
4948 )?;
4949 let rows: Vec<(i64, Vec<u8>)> = {
4950 let mut statement = tx.prepare(
4951 "SELECT rowid, embedding FROM vector_default ORDER BY rowid",
4952 )?;
4953 let mapped = statement.query_map([], |row| {
4954 Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?))
4955 })?;
4956 let mut out = Vec::new();
4957 for r in mapped {
4958 out.push(r?);
4959 }
4960 out
4961 };
4962 let (doc_count, _) =
4963 run_pin_and_requantize_pass(&tx, &rows, &mean).map_err(|_| {
4964 rusqlite::Error::SqliteFailure(
4965 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_ERROR),
4966 Some("mean-centering re-quantize pass failed".to_string()),
4967 )
4968 })?;
4969 staged_events.push(EmbedderEvent::MeanVecPinned {
4970 dim: u32::try_from(mean.len()).unwrap_or(u32::MAX),
4971 doc_count,
4972 });
4973 current_mean = Some(mean);
4974 }
4975 }
4976 ProjectionOutcome::Failure { cursor, failure_code } => {
4977 if terminal_state_for_cursor(&tx, *cursor)?.is_some() {
4978 continue;
4979 }
4980 let existing: u64 = tx.query_row(
4981 "SELECT COUNT(*) FROM operational_mutations
4982 WHERE collection_name = 'projection_failures'
4983 AND json_extract(payload_json, '$.write_cursor') = ?1",
4984 [cursor],
4985 |row| row.get(0),
4986 )?;
4987 if existing == 0 {
4988 let payload = format!(
4989 r#"{{"write_cursor":{cursor},"failure_code":"{failure_code}","recorded_at":0}}"#
4990 );
4991 tx.execute(
4992 "INSERT INTO operational_mutations(
4993 collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
4994 ) VALUES('projection_failures', ?1, 'append', ?2, NULL, ?3)",
4995 params![cursor.to_string(), payload, cursor],
4996 )?;
4997 }
4998 record_projection_terminal(&tx, *cursor, "failed")?;
4999 }
5000 }
5001 }
5002 advance_projection_cursor(&tx)?;
5012 tx.commit()?;
5013 if !staged_events.is_empty() {
5016 if let Ok(mut events) = shared.pending_events.lock() {
5017 events.extend(staged_events);
5018 }
5019 }
5020 Ok(())
5021}
5022
5023fn recover_mean_vec_pin(
5030 connection: &mut Connection,
5031 identity: &EmbedderIdentity,
5032) -> Result<(), EngineError> {
5033 let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
5034 recompute_mean_in_tx(&tx, identity)?;
5035 tx.commit().map_err(|_| EngineError::Storage)?;
5036 Ok(())
5037}
5038
5039fn recompute_mean_in_tx(
5053 tx: &rusqlite::Transaction<'_>,
5054 identity: &EmbedderIdentity,
5055) -> Result<MeanRecomputeReport, EngineError> {
5056 recompute_mean_in_tx_inner(tx, identity, false)
5057}
5058
5059fn recompute_mean_in_tx_inner(
5064 tx: &rusqlite::Transaction<'_>,
5065 identity: &EmbedderIdentity,
5066 fail_after_mean_update: bool,
5067) -> Result<MeanRecomputeReport, EngineError> {
5068 let started = Instant::now();
5069 let dim = identity.dimension as usize;
5070 let old_mean = read_pinned_mean_vec(tx, identity.dimension)?;
5073 let rows: Vec<(i64, Vec<u8>)> = {
5074 let mut statement = tx
5075 .prepare("SELECT rowid, embedding FROM vector_default ORDER BY rowid")
5076 .map_err(|_| EngineError::Storage)?;
5077 let mapped = statement
5078 .query_map([], |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?)))
5079 .map_err(|_| EngineError::Storage)?;
5080 let mut out = Vec::new();
5081 for r in mapped {
5082 out.push(r.map_err(|_| EngineError::Storage)?);
5083 }
5084 out
5085 };
5086 let mut accumulator = MeanAccumulator::new(dim);
5087 for (_rowid, blob) in &rows {
5088 if blob.len() != dim * 4 {
5089 return Err(EngineError::Storage);
5090 }
5091 accumulator.add(&decode_vector_blob(blob));
5092 }
5093 let old_doc_count = accumulator.count();
5094 let mean = accumulator.materialize();
5095 let drift_cos_before = match &old_mean {
5096 Some(old) => cosine_similarity(&mean, old),
5097 None => 1.0,
5098 };
5099 tx.execute(
5100 "UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
5101 params![encode_vector_blob(&mean)],
5102 )
5103 .map_err(|_| EngineError::Storage)?;
5104 if fail_after_mean_update {
5105 return Err(EngineError::Storage);
5108 }
5109 let (doc_count, _) = run_pin_and_requantize_pass(tx, &rows, &mean)?;
5110 Ok(MeanRecomputeReport {
5111 dim: u32::try_from(dim).unwrap_or(u32::MAX),
5112 old_doc_count,
5113 doc_count_requantized: doc_count,
5114 drift_cos_before,
5115 mean_was_pinned: old_mean.is_some(),
5116 elapsed_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
5117 })
5118}
5119
5120fn enforce_provenance_retention(connection: &Connection, cap: u64) -> rusqlite::Result<()> {
5121 if cap == 0 {
5122 return Ok(());
5123 }
5124 let slack = cap.max(20) / 20;
5125 let upper = cap.saturating_add(slack.max(1));
5126 let count: u64 =
5127 connection.query_row("SELECT COUNT(*) FROM operational_mutations", [], |row| row.get(0))?;
5128 if count <= upper {
5129 return Ok(());
5130 }
5131 let to_delete = count.saturating_sub(cap);
5132 connection.execute(
5133 "DELETE FROM operational_mutations
5134 WHERE id IN (
5135 SELECT id FROM operational_mutations
5136 ORDER BY id
5137 LIMIT ?1
5138 )",
5139 [to_delete],
5140 )?;
5141 Ok(())
5142}
5143
5144fn projection_status(
5145 connection: &Connection,
5146 kind: &str,
5147) -> Result<lifecycle::ProjectionStatus, EngineError> {
5148 let latest = connection
5149 .query_row(
5150 "SELECT COALESCE(MAX(write_cursor), 0) FROM canonical_nodes WHERE kind = ?1",
5151 [kind],
5152 |row| row.get::<_, u64>(0),
5153 )
5154 .map_err(|_| EngineError::Storage)?;
5155 if latest == 0 {
5156 return Ok(lifecycle::ProjectionStatus::UpToDate);
5157 }
5158 let pending: u64 = connection
5159 .query_row(
5160 "SELECT COUNT(*)
5161 FROM canonical_nodes
5162 LEFT JOIN _fathomdb_projection_terminal
5163 ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
5164 WHERE canonical_nodes.kind = ?1
5165 AND _fathomdb_projection_terminal.write_cursor IS NULL",
5166 [kind],
5167 |row| row.get(0),
5168 )
5169 .map_err(|_| EngineError::Storage)?;
5170 if pending > 0 {
5171 return Ok(lifecycle::ProjectionStatus::Pending);
5172 }
5173 match terminal_state_for_cursor(connection, latest).map_err(|_| EngineError::Storage)? {
5174 Some(state) if state == "failed" => Ok(lifecycle::ProjectionStatus::Failed),
5175 _ => Ok(lifecycle::ProjectionStatus::UpToDate),
5176 }
5177}
5178
5179fn canonical_database_path(path: &Path) -> Result<PathBuf, EngineOpenError> {
5180 let parent = path
5181 .parent()
5182 .filter(|parent| !parent.as_os_str().is_empty())
5183 .unwrap_or_else(|| Path::new("."));
5184 let canonical_parent = parent.canonicalize().map_err(|_| EngineOpenError::Io {
5185 message: "database parent directory is not accessible".to_string(),
5186 })?;
5187 let file_name = path.file_name().ok_or_else(|| EngineOpenError::Io {
5188 message: "database path has no file name".to_string(),
5189 })?;
5190
5191 Ok(canonical_parent.join(file_name))
5192}
5193
5194fn acquire_lock(path: &Path) -> Result<File, EngineOpenError> {
5195 let lock_path = lock_path(path);
5196 let mut options = OpenOptions::new();
5197 options.read(true).write(true).create(true);
5198 #[cfg(unix)]
5199 options.mode(0o600);
5200
5201 let mut file = options.open(&lock_path).map_err(|_| EngineOpenError::Io {
5202 message: "could not open database lock file".to_string(),
5203 })?;
5204
5205 match file.try_lock() {
5206 Ok(()) => {
5207 let pid = std::process::id().to_string();
5208 let _ = file.set_len(0);
5209 let _ = file.seek(SeekFrom::Start(0));
5210 let _ = file.write_all(pid.as_bytes());
5211 Ok(file)
5212 }
5213 Err(std::fs::TryLockError::WouldBlock) => {
5214 Err(EngineOpenError::DatabaseLocked { holder_pid: read_holder_pid(&lock_path) })
5215 }
5216 Err(_) => {
5217 Err(EngineOpenError::Io { message: "could not acquire database lock".to_string() })
5218 }
5219 }
5220}
5221
5222fn lock_path(path: &Path) -> PathBuf {
5223 let mut lock_path = path.as_os_str().to_os_string();
5224 lock_path.push(LOCK_SUFFIX);
5225 PathBuf::from(lock_path)
5226}
5227
5228fn read_holder_pid(path: &Path) -> Option<u32> {
5229 std::fs::read_to_string(path).ok()?.trim().parse().ok()
5230}
5231
5232fn map_migration_error(err: SchemaMigrationError) -> EngineOpenError {
5233 match err {
5234 SchemaMigrationError::IncompatibleSchemaVersion { seen, supported } => {
5235 EngineOpenError::IncompatibleSchemaVersion { seen, supported }
5236 }
5237 SchemaMigrationError::MigrationError(report) => EngineOpenError::MigrationError {
5238 schema_version_before: report.schema_version_before,
5239 schema_version_current: report.schema_version_current,
5240 step_id: report.migration_steps.last().map_or(0, |step| step.step_id),
5241 },
5242 SchemaMigrationError::Storage { message } => {
5243 EngineOpenError::Io { message: message.to_string() }
5244 }
5245 }
5246}
5247
5248fn init_perf_experiments_runtime() {
5264 static INIT: Once = Once::new();
5265 INIT.call_once(|| {
5266 if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
5267 return;
5268 }
5269 let memstatus_off =
5270 std::env::var_os("FATHOMDB_PERF_SQLITE_MEMSTATUS_OFF").is_some_and(|v| v == "1");
5271 let pagecache = std::env::var("FATHOMDB_PERF_SQLITE_PAGECACHE").ok();
5277 let pcache2_on =
5281 std::env::var_os("FATHOMDB_PERF_SQLITE_PCACHE2").is_some_and(|v| v == "1");
5282 if !memstatus_off && pagecache.is_none() && !pcache2_on {
5283 return;
5284 }
5285 unsafe {
5292 let rc_shutdown = rusqlite::ffi::sqlite3_shutdown();
5293 let rc_memstatus = if memstatus_off {
5294 rusqlite::ffi::sqlite3_config(rusqlite::ffi::SQLITE_CONFIG_MEMSTATUS, 0_i32)
5295 } else {
5296 -1
5297 };
5298 let rc_pagecache = if let Some(spec) = pagecache.as_ref() {
5302 let mut parts = spec.split(':');
5303 let sz = parts.next().and_then(|s| s.parse::<i32>().ok()).unwrap_or(0);
5304 let n = parts.next().and_then(|s| s.parse::<i32>().ok()).unwrap_or(0);
5305 if sz > 0 && n > 0 {
5306 rusqlite::ffi::sqlite3_config(
5307 7, std::ptr::null_mut::<std::ffi::c_void>(),
5309 sz,
5310 n,
5311 )
5312 } else {
5313 eprintln!(
5314 "perf-experiment: bad FATHOMDB_PERF_SQLITE_PAGECACHE spec '{spec}' (expect '<bytes>:<count>')"
5315 );
5316 -1
5317 }
5318 } else {
5319 -1
5320 };
5321 let rc_pcache2 = if pcache2_on {
5322 rusqlite::ffi::sqlite3_config(
5326 rusqlite::ffi::SQLITE_CONFIG_PCACHE2,
5327 &raw const pcache2::PCACHE2_METHODS.0,
5328 )
5329 } else {
5330 -1
5331 };
5332 let rc_init = rusqlite::ffi::sqlite3_initialize();
5333 eprintln!(
5334 "perf-experiment: runtime-config rcs shutdown={rc_shutdown} \
5335 memstatus={rc_memstatus} pagecache={rc_pagecache} pcache2={rc_pcache2} \
5336 initialize={rc_init} (0=SQLITE_OK; 21=SQLITE_MISUSE; -1=not configured)"
5337 );
5338 }
5339 });
5340}
5341
5342fn register_sqlite_vec_extension() {
5343 static REGISTER: Once = Once::new();
5344 REGISTER.call_once(|| unsafe {
5345 let entrypoint: unsafe extern "C" fn(
5346 *mut rusqlite::ffi::sqlite3,
5347 *mut *const std::os::raw::c_char,
5348 *const rusqlite::ffi::sqlite3_api_routines,
5349 ) -> std::os::raw::c_int = std::mem::transmute(sqlite3_vec_init as *const ());
5350 rusqlite::ffi::sqlite3_auto_extension(Some(entrypoint));
5351 });
5352}
5353
5354fn probe_open_integrity(connection: &Connection) -> Result<(), EngineOpenError> {
5355 connection
5360 .query_row("SELECT COUNT(*) FROM sqlite_schema", [], |row| row.get::<_, i64>(0))
5361 .map(|_| ())
5362 .map_err(|err| map_open_sqlite_error(err, OpenStage::SchemaProbe))
5363}
5364
5365fn probe_database_header(connection: &Connection) -> Result<(), EngineOpenError> {
5366 connection
5367 .query_row("PRAGMA application_id", [], |row| row.get::<_, i64>(0))
5368 .map(|_| ())
5369 .map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))
5370}
5371
5372fn probe_wal_sidecar(db_path: &Path) -> Result<(), EngineOpenError> {
5379 let mut wal_path = db_path.as_os_str().to_owned();
5380 wal_path.push("-wal");
5381 let wal_path = PathBuf::from(wal_path);
5382 use std::io::Read;
5390 let mut file = match std::fs::File::open(&wal_path) {
5391 Ok(file) => file,
5392 Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(()),
5393 Err(_) => return Ok(()),
5394 };
5395 let mut bytes = [0u8; 32];
5396 if file.read_exact(&mut bytes).is_err() {
5397 return Ok(());
5400 }
5401 let magic = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
5402 let page_size = u32::from_be_bytes([bytes[8], bytes[9], bytes[10], bytes[11]]);
5403 const WAL_MAGIC_MASK: u32 = 0xFFFF_FFFE;
5407 const WAL_MAGIC: u32 = 0x377F_0682;
5408 const SQLITE_MAX_PAGE_SIZE: u32 = 65536;
5409 let magic_ok = (magic & WAL_MAGIC_MASK) == WAL_MAGIC;
5410 let page_size_ok =
5411 page_size.is_power_of_two() && (512..=SQLITE_MAX_PAGE_SIZE).contains(&page_size);
5412 if magic_ok && page_size_ok {
5413 return Ok(());
5414 }
5415 Err(EngineOpenError::Corruption(CorruptionDetail {
5416 kind: CorruptionKind::WalReplayFailure,
5417 stage: OpenStage::WalReplay,
5418 locator: CorruptionLocator::FileOffset { offset: if !magic_ok { 0 } else { 8 } },
5419 recovery_hint: RecoveryHint {
5420 code: "E_CORRUPT_WAL_REPLAY",
5421 doc_anchor: "design/recovery.md#wal-replay-failures",
5422 },
5423 }))
5424}
5425
5426fn reject_legacy_shape(connection: &Connection) -> Result<(), EngineOpenError> {
5427 let has_legacy_table = table_exists(connection, "fathom_nodes")
5428 || table_exists(connection, "fathom_edges")
5429 || table_exists(connection, "fathom_chunks");
5430 if !has_legacy_table {
5431 return Ok(());
5432 }
5433
5434 let seen =
5435 connection.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)).unwrap_or(0);
5436 Err(EngineOpenError::IncompatibleSchemaVersion { seen, supported: SCHEMA_VERSION })
5437}
5438
5439fn table_exists(connection: &Connection, table: &str) -> bool {
5440 connection
5441 .query_row(
5442 "SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = ?1",
5443 [table],
5444 |_row| Ok(()),
5445 )
5446 .is_ok()
5447}
5448
5449#[cfg(feature = "operator")]
5450fn read_schema_objects(
5451 connection: &Connection,
5452 obj_type: &str,
5453) -> Result<Vec<SchemaObject>, EngineError> {
5454 let mut stmt = connection
5455 .prepare(
5456 "SELECT name, sql FROM sqlite_schema
5457 WHERE type = ?1 AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL
5458 ORDER BY name",
5459 )
5460 .map_err(|_| EngineError::Storage)?;
5461 let rows = stmt
5462 .query_map([obj_type], |row| {
5463 Ok(SchemaObject { name: row.get::<_, String>(0)?, sql: row.get::<_, String>(1)? })
5464 })
5465 .map_err(|_| EngineError::Storage)?;
5466 let mut out = Vec::new();
5467 for row in rows {
5468 out.push(row.map_err(|_| EngineError::Storage)?);
5469 }
5470 Ok(out)
5471}
5472
5473#[cfg(feature = "operator")]
5474fn order_canonical_first(mut objects: Vec<SchemaObject>) -> Vec<SchemaObject> {
5475 let mut canonical: Vec<SchemaObject> = Vec::new();
5476 for name in CANONICAL_TABLES {
5477 if let Some(pos) = objects.iter().position(|o| o.name == *name) {
5478 canonical.push(objects.remove(pos));
5479 }
5480 }
5481 canonical.extend(objects);
5482 canonical
5483}
5484
5485fn load_default_profile(connection: &Connection) -> rusqlite::Result<EmbedderIdentity> {
5486 connection.query_row(
5487 "SELECT name, revision, dimension FROM _fathomdb_embedder_profiles WHERE profile = ?1",
5488 [DEFAULT_VECTOR_PROFILE],
5489 |row| {
5490 Ok(EmbedderIdentity::new(
5491 row.get::<_, String>(0)?,
5492 row.get::<_, String>(1)?,
5493 row.get::<_, u32>(2)?,
5494 ))
5495 },
5496 )
5497}
5498
5499fn default_profile_dimension(connection: &Connection) -> Result<u32, EngineError> {
5500 load_default_profile(connection)
5501 .map(|identity| identity.dimension)
5502 .map_err(|_| EngineError::Storage)
5503}
5504
5505fn kind_is_vector_indexed(connection: &Connection, kind: &str) -> Result<bool, EngineError> {
5506 connection
5507 .query_row("SELECT 1 FROM _fathomdb_vector_kinds WHERE kind = ?1", [kind], |_row| Ok(()))
5508 .map(|_| true)
5509 .or_else(|err| match err {
5510 rusqlite::Error::QueryReturnedNoRows => Ok(false),
5511 _ => Err(EngineError::Storage),
5512 })
5513}
5514
5515fn ensure_vector_partition(connection: &mut Connection, dimension: u32) -> rusqlite::Result<()> {
5516 let existing_sql: Option<String> = connection
5530 .query_row(
5531 "SELECT sql FROM sqlite_master WHERE type='table' AND name=?1",
5532 [DEFAULT_VECTOR_PARTITION],
5533 |row| row.get::<_, String>(0),
5534 )
5535 .optional()?;
5536
5537 match existing_sql {
5544 None => create_vector_partition(connection, dimension),
5545 Some(sql) if sql.contains("status") => Ok(()),
5546 Some(sql) if sql.contains("embedding_bin") => {
5547 migrate_vector_partition_pack1_to_pack2(connection, dimension)
5548 }
5549 Some(_) => migrate_vector_partition_to_pack1(connection, dimension),
5550 }
5551}
5552
5553fn vector_partition_create_sql(dimension: u32, if_not_exists: bool) -> String {
5559 let guard = if if_not_exists { "IF NOT EXISTS " } else { "" };
5560 format!(
5561 "CREATE VIRTUAL TABLE {guard}{DEFAULT_VECTOR_PARTITION} USING vec0(\
5562 embedding float[{dimension}],\
5563 embedding_bin bit[{dimension}],\
5564 source_type TEXT partition key,\
5565 kind TEXT,\
5566 created_at INTEGER,\
5567 status TEXT\
5568 )"
5569 )
5570}
5571
5572fn create_vector_partition(connection: &Connection, dimension: u32) -> rusqlite::Result<()> {
5573 connection.execute_batch(&vector_partition_create_sql(dimension, true))
5574}
5575
5576fn migrate_vector_partition_pack1_to_pack2(
5586 connection: &mut Connection,
5587 dimension: u32,
5588) -> rusqlite::Result<()> {
5589 let tx = connection.transaction()?;
5590 tx.execute_batch(
5591 "CREATE TABLE _fathomdb_vector_pack2_stage (
5592 rowid INTEGER PRIMARY KEY,
5593 embedding BLOB NOT NULL,
5594 embedding_bin BLOB NOT NULL,
5595 source_type TEXT,
5596 kind TEXT,
5597 created_at INTEGER
5598 );
5599 INSERT INTO _fathomdb_vector_pack2_stage(
5600 rowid, embedding, embedding_bin, source_type, kind, created_at
5601 )
5602 SELECT rowid, embedding, embedding_bin, source_type, kind, created_at
5603 FROM vector_default;
5604 DROP TABLE vector_default;",
5605 )?;
5606 tx.execute_batch(&vector_partition_create_sql(dimension, false))?;
5607 tx.execute_batch(
5614 "INSERT INTO vector_default(
5615 rowid, embedding, embedding_bin, source_type, kind, created_at, status
5616 )
5617 SELECT rowid, embedding, vec_bit(embedding_bin), source_type, kind, created_at, ''
5618 FROM _fathomdb_vector_pack2_stage;
5619 DROP TABLE _fathomdb_vector_pack2_stage;",
5620 )?;
5621 tx.commit()
5622}
5623
5624const KIND_TO_SOURCE_TYPE_CASE_SQL: &str = "CASE s.kind
5628 WHEN 'email' THEN 'email'
5629 WHEN 'article' THEN 'article'
5630 WHEN 'paper' THEN 'paper'
5631 WHEN 'meeting' THEN 'meeting'
5632 WHEN 'note' THEN 'note'
5633 WHEN 'todo' THEN 'todo'
5634 WHEN 'doc' THEN 'article'
5635 ELSE 'article'
5636END";
5637
5638fn migrate_vector_partition_to_pack1(
5654 connection: &mut Connection,
5655 dimension: u32,
5656) -> rusqlite::Result<()> {
5657 let tx = connection.transaction()?;
5658 tx.execute_batch(
5659 "CREATE TABLE _fathomdb_vector_migration_v0_7_0 (
5660 rowid INTEGER PRIMARY KEY,
5661 embedding BLOB NOT NULL,
5662 kind TEXT NOT NULL
5663 );
5664 INSERT INTO _fathomdb_vector_migration_v0_7_0(rowid, embedding, kind)
5665 SELECT v.rowid, v.embedding, r.kind
5666 FROM vector_default v
5667 JOIN _fathomdb_vector_rows r ON r.rowid = v.rowid;
5668 DROP TABLE vector_default;",
5669 )?;
5670 tx.execute_batch(&vector_partition_create_sql(dimension, false))?;
5673 let repopulate_sql = format!(
5678 "INSERT INTO vector_default(
5679 rowid, embedding, embedding_bin, source_type, kind, created_at, status
5680 )
5681 SELECT
5682 s.rowid,
5683 s.embedding,
5684 vec_quantize_binary(s.embedding),
5685 {KIND_TO_SOURCE_TYPE_CASE_SQL},
5686 s.kind,
5687 strftime('%s', 'now'),
5688 ''
5689 FROM _fathomdb_vector_migration_v0_7_0 s;
5690 DROP TABLE _fathomdb_vector_migration_v0_7_0;"
5691 );
5692 tx.execute_batch(&repopulate_sql)?;
5693 tx.commit()
5694}
5695
5696fn encode_vector_blob(vector: &[f32]) -> Vec<u8> {
5697 vector.iter().flat_map(|value| value.to_le_bytes()).collect()
5698}
5699
5700fn decode_vector_blob(bytes: &[u8]) -> Vec<f32> {
5701 debug_assert_eq!(bytes.len() % 4, 0, "f32 BLOB length must be multiple of 4");
5702 bytes.chunks_exact(4).map(|c| f32::from_le_bytes([c[0], c[1], c[2], c[3]])).collect()
5703}
5704
5705fn identity_requires_mean_centering(identity: &EmbedderIdentity) -> bool {
5709 identity.name == BGE_SMALL_EMBEDDER_NAME
5710}
5711
5712fn read_pinned_mean_vec(
5719 connection: &Connection,
5720 dimension: u32,
5721) -> Result<Option<Vec<f32>>, EngineError> {
5722 let bytes: Option<Vec<u8>> = connection
5723 .query_row(
5724 "SELECT mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
5725 [],
5726 |row| row.get::<_, Option<Vec<u8>>>(0),
5727 )
5728 .or_else(|err| match err {
5729 rusqlite::Error::QueryReturnedNoRows => Ok(None),
5730 other => Err(other),
5731 })
5732 .map_err(|_| EngineError::Storage)?;
5733 let Some(bytes) = bytes else { return Ok(None) };
5734 let expected_len = (dimension as usize).saturating_mul(4);
5735 if bytes.len() != expected_len {
5736 return Err(EngineError::Storage);
5737 }
5738 let mut out = Vec::with_capacity(dimension as usize);
5739 for chunk in bytes.chunks_exact(4) {
5740 let arr = [chunk[0], chunk[1], chunk[2], chunk[3]];
5741 out.push(f32::from_le_bytes(arr));
5742 }
5743 Ok(Some(out))
5744}
5745
5746fn subtract_mean(v: &[f32], mean: &[f32]) -> Vec<f32> {
5749 debug_assert_eq!(v.len(), mean.len(), "subtract_mean dim mismatch");
5750 v.iter().zip(mean.iter()).map(|(a, b)| *a - *b).collect()
5751}
5752
5753fn resolve_source_type(kind: &str) -> Result<&'static str, EngineError> {
5760 Ok(match kind {
5761 "email" => "email",
5762 "article" => "article",
5763 "paper" => "paper",
5764 "meeting" => "meeting",
5765 "note" => "note",
5766 "todo" => "todo",
5767 "doc" => "article",
5769 _ => return Err(EngineError::Storage),
5770 })
5771}
5772
5773fn map_runtime_embedder_error(err: RuntimeEmbedderError) -> EngineError {
5774 match err {
5775 RuntimeEmbedderError::Failed { .. } | RuntimeEmbedderError::Timeout => {
5776 EngineError::Embedder
5777 }
5778 }
5779}
5780
5781fn default_embedder_identity() -> EmbedderIdentity {
5782 EmbedderIdentity::new(
5783 DEFAULT_EMBEDDER_NAME,
5784 DEFAULT_EMBEDDER_REVISION,
5785 DEFAULT_EMBEDDER_DIMENSION,
5786 )
5787}
5788
5789fn check_embedder_profile(
5790 connection: &Connection,
5791 supplied: &EmbedderIdentity,
5792) -> Result<bool, EngineOpenError> {
5793 let mut statement = match connection.prepare(
5797 "SELECT name, revision, dimension, mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
5798 ) {
5799 Ok(statement) => statement,
5800 Err(_) => return Ok(false),
5801 };
5802 let mut rows = statement.query([]).map_err(|_| {
5803 EngineOpenError::Corruption(CorruptionDetail {
5804 kind: CorruptionKind::EmbedderIdentityDrift,
5805 stage: OpenStage::EmbedderIdentity,
5806 locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
5807 recovery_hint: RecoveryHint {
5808 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5809 doc_anchor: "design/recovery.md#embedder-identity-drift",
5810 },
5811 })
5812 })?;
5813
5814 let Some(row) = rows.next().map_err(|_| {
5815 EngineOpenError::Corruption(CorruptionDetail {
5816 kind: CorruptionKind::EmbedderIdentityDrift,
5817 stage: OpenStage::EmbedderIdentity,
5818 locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
5819 recovery_hint: RecoveryHint {
5820 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5821 doc_anchor: "design/recovery.md#embedder-identity-drift",
5822 },
5823 })
5824 })?
5825 else {
5826 connection
5827 .execute(
5828 "INSERT INTO _fathomdb_embedder_profiles(profile, name, revision, dimension)
5829 VALUES(?1, ?2, ?3, ?4)",
5830 params![
5831 DEFAULT_VECTOR_PROFILE,
5832 supplied.name,
5833 supplied.revision,
5834 supplied.dimension
5835 ],
5836 )
5837 .map_err(|_| EngineOpenError::Io {
5838 message: "could not persist embedder profile".to_string(),
5839 })?;
5840 return Ok(false);
5841 };
5842
5843 let stored_name = row.get::<_, String>(0).map_err(|_| {
5844 EngineOpenError::Corruption(CorruptionDetail {
5845 kind: CorruptionKind::EmbedderIdentityDrift,
5846 stage: OpenStage::EmbedderIdentity,
5847 locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
5848 recovery_hint: RecoveryHint {
5849 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5850 doc_anchor: "design/recovery.md#embedder-identity-drift",
5851 },
5852 })
5853 })?;
5854 let stored_revision = row.get::<_, String>(1).map_err(|_| {
5855 EngineOpenError::Corruption(CorruptionDetail {
5856 kind: CorruptionKind::EmbedderIdentityDrift,
5857 stage: OpenStage::EmbedderIdentity,
5858 locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
5859 recovery_hint: RecoveryHint {
5860 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5861 doc_anchor: "design/recovery.md#embedder-identity-drift",
5862 },
5863 })
5864 })?;
5865 let dimension = row.get::<_, u32>(2).map_err(|_| {
5866 EngineOpenError::Corruption(CorruptionDetail {
5867 kind: CorruptionKind::EmbedderIdentityDrift,
5868 stage: OpenStage::EmbedderIdentity,
5869 locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
5870 recovery_hint: RecoveryHint {
5871 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5872 doc_anchor: "design/recovery.md#embedder-identity-drift",
5873 },
5874 })
5875 })?;
5876
5877 let stored = EmbedderIdentity::new(stored_name, stored_revision, dimension);
5878
5879 if stored.name != supplied.name || stored.revision != supplied.revision {
5880 return Err(EngineOpenError::EmbedderIdentityMismatch {
5881 stored,
5882 supplied: supplied.clone(),
5883 });
5884 }
5885 if dimension != supplied.dimension {
5886 return Err(EngineOpenError::EmbedderDimensionMismatch {
5887 stored: dimension,
5888 supplied: supplied.dimension,
5889 });
5890 }
5891
5892 let mean_vec: Option<Vec<u8>> = row.get::<_, Option<Vec<u8>>>(3).map_err(|_| {
5897 EngineOpenError::Corruption(CorruptionDetail {
5898 kind: CorruptionKind::EmbedderIdentityDrift,
5899 stage: OpenStage::EmbedderIdentity,
5900 locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
5901 recovery_hint: RecoveryHint {
5902 code: "E_CORRUPT_EMBEDDER_IDENTITY",
5903 doc_anchor: "design/recovery.md#embedder-identity-drift",
5904 },
5905 })
5906 })?;
5907 let pinned = match mean_vec {
5908 Some(bytes) => {
5909 let expected_len = (dimension as usize).saturating_mul(4);
5910 if bytes.len() != expected_len {
5916 return Err(EngineOpenError::EmbedderIdentityMismatch {
5917 stored,
5918 supplied: supplied.clone(),
5919 });
5920 }
5921 true
5922 }
5923 None => false,
5924 };
5925
5926 Ok(pinned)
5927}
5928
5929#[derive(Clone, Debug, Eq, PartialEq)]
5930enum WritePlan {
5931 Node,
5932 Edge,
5933 AppendOnlyLog,
5934 LatestState,
5935 AdminSchema,
5936}
5937
5938fn validate_batch(
5939 connection: &Connection,
5940 batch: &[PreparedWrite],
5941) -> Result<Vec<WritePlan>, EngineError> {
5942 batch.iter().map(|write| validate_write(connection, write)).collect()
5943}
5944
5945fn collect_projection_jobs(
5946 connection: &Connection,
5947 batch: &[PreparedWrite],
5948) -> Result<Vec<ProjectionJob>, EngineError> {
5949 let mut jobs = Vec::new();
5950 for write in batch {
5951 if let PreparedWrite::Node { kind, body, .. } = write {
5952 if kind_is_vector_indexed(connection, kind)? {
5953 jobs.push(ProjectionJob { cursor: 0, kind: kind.clone(), body: body.clone() });
5954 }
5955 }
5956 }
5957 Ok(jobs)
5958}
5959
5960fn validate_write(
5961 connection: &Connection,
5962 write: &PreparedWrite,
5963) -> Result<WritePlan, EngineError> {
5964 match write {
5965 PreparedWrite::Node { kind, body, source_id, logical_id } => {
5966 if kind.trim().is_empty() || body.trim().is_empty() {
5967 return Err(EngineError::WriteValidation);
5968 }
5969 if let Some(source_id) = source_id {
5970 if source_id.is_empty() {
5971 return Err(EngineError::WriteValidation);
5972 }
5973 }
5974 if let Some(logical_id) = logical_id {
5977 if logical_id.is_empty() {
5978 return Err(EngineError::WriteValidation);
5979 }
5980 }
5981 Ok(WritePlan::Node)
5982 }
5983 PreparedWrite::Edge { kind, from, to, source_id, logical_id } => {
5984 if kind.trim().is_empty() || from.trim().is_empty() || to.trim().is_empty() {
5985 return Err(EngineError::WriteValidation);
5986 }
5987 if let Some(source_id) = source_id {
5988 if source_id.is_empty() {
5989 return Err(EngineError::WriteValidation);
5990 }
5991 }
5992 if let Some(logical_id) = logical_id {
5993 if logical_id.is_empty() {
5994 return Err(EngineError::WriteValidation);
5995 }
5996 }
5997 Ok(WritePlan::Edge)
5998 }
5999 PreparedWrite::AdminSchema { name, kind, schema_json, retention_json } => {
6000 if name.trim().is_empty()
6001 || !matches!(kind.as_str(), "append_only_log" | "latest_state")
6002 || serde_json::from_str::<Value>(schema_json).is_err()
6003 || serde_json::from_str::<Value>(retention_json).is_err()
6004 || contains_external_ref(schema_json)
6005 {
6006 return Err(EngineError::SchemaValidation);
6007 }
6008 Ok(WritePlan::AdminSchema)
6009 }
6010 PreparedWrite::OpStore { collection, record_key, schema_id, body } => {
6011 if collection.trim().is_empty() || record_key.trim().is_empty() {
6012 return Err(EngineError::WriteValidation);
6013 }
6014 let (kind, schema_json) = collection_metadata(connection, collection)?;
6015 if let Some(schema_id) = schema_id {
6016 if schema_id != collection {
6017 return Err(EngineError::SchemaValidation);
6018 }
6019 validate_payload(&schema_json, body)?;
6020 } else if serde_json::from_str::<Value>(body).is_err() {
6021 return Err(EngineError::SchemaValidation);
6022 }
6023
6024 match kind.as_str() {
6025 "append_only_log" => Ok(WritePlan::AppendOnlyLog),
6026 "latest_state" => Ok(WritePlan::LatestState),
6027 _ => Err(EngineError::OpStore),
6028 }
6029 }
6030 }
6031}
6032
6033fn collection_metadata(
6034 connection: &Connection,
6035 collection: &str,
6036) -> Result<(String, String), EngineError> {
6037 connection
6038 .query_row(
6039 "SELECT kind, schema_json FROM operational_collections WHERE name = ?1",
6040 [collection],
6041 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
6042 )
6043 .map_err(|_| EngineError::OpStore)
6044}
6045
6046fn validate_payload(schema_json: &str, body: &str) -> Result<(), EngineError> {
6047 let schema =
6048 serde_json::from_str::<Value>(schema_json).map_err(|_| EngineError::SchemaValidation)?;
6049 let payload = serde_json::from_str::<Value>(body).map_err(|_| EngineError::SchemaValidation)?;
6050
6051 let compiled = JSONSchema::compile(&schema).map_err(|_| EngineError::SchemaValidation)?;
6052 compiled.validate(&payload).map_err(|_| EngineError::SchemaValidation)?;
6053
6054 Ok(())
6055}
6056
6057fn contains_external_ref(schema_json: &str) -> bool {
6058 let Ok(value) = serde_json::from_str::<Value>(schema_json) else {
6059 return false;
6060 };
6061 value_contains_external_ref(&value)
6062}
6063
6064fn value_contains_external_ref(value: &Value) -> bool {
6065 match value {
6066 Value::Object(object) => object.iter().any(|(key, value)| {
6067 if key == "$ref" {
6068 return value.as_str().is_some_and(|uri| !uri.starts_with('#'));
6069 }
6070 value_contains_external_ref(value)
6071 }),
6072 Value::Array(values) => values.iter().any(value_contains_external_ref),
6073 _ => false,
6074 }
6075}
6076
6077fn commit_batch(
6078 connection: &mut Connection,
6079 batch: &[PreparedWrite],
6080 plans: &[WritePlan],
6081 base_cursor: u64,
6082 provenance_row_cap: u64,
6083) -> rusqlite::Result<u64> {
6084 let tx = connection.transaction()?;
6085
6086 for (i, (write, plan)) in batch.iter().zip(plans).enumerate() {
6087 let cursor = base_cursor.saturating_add((i as u64).saturating_add(1));
6090 match (write, plan) {
6091 (PreparedWrite::Node { kind, body, source_id, logical_id }, WritePlan::Node) => {
6092 if let Some(logical_id) = logical_id {
6100 tx.execute(
6101 "UPDATE canonical_nodes SET superseded_at = ?1
6102 WHERE logical_id = ?2 AND superseded_at IS NULL",
6103 params![cursor, logical_id],
6104 )?;
6105 }
6106 tx.execute(
6107 "INSERT INTO canonical_nodes(write_cursor, kind, body, source_id, logical_id)
6108 VALUES(?1, ?2, ?3, ?4, ?5)",
6109 params![cursor, kind, body, source_id, logical_id],
6110 )?;
6111 tx.execute(
6112 "INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)",
6113 params![body, kind, cursor],
6114 )?;
6115 if kind_is_vector_indexed(&tx, kind).unwrap_or(false) {
6116 tx.execute(
6117 "INSERT INTO _fathomdb_projection_state(kind, last_enqueued_cursor, updated_at)
6118 VALUES(?1, ?2, 0)
6119 ON CONFLICT(kind) DO UPDATE SET last_enqueued_cursor = excluded.last_enqueued_cursor",
6120 params![kind, cursor],
6121 )?;
6122 } else {
6123 record_projection_terminal(&tx, cursor, "up_to_date")?;
6127 }
6128 }
6129 (PreparedWrite::Edge { kind, from, to, source_id, logical_id }, WritePlan::Edge) => {
6130 if let Some(logical_id) = logical_id {
6136 tx.execute(
6137 "UPDATE canonical_edges SET superseded_at = ?1
6138 WHERE logical_id = ?2 AND superseded_at IS NULL",
6139 params![cursor, logical_id],
6140 )?;
6141 }
6142 tx.execute(
6143 "INSERT INTO canonical_edges(write_cursor, kind, from_id, to_id, source_id, logical_id)
6144 VALUES(?1, ?2, ?3, ?4, ?5, ?6)",
6145 params![cursor, kind, from, to, source_id, logical_id],
6146 )?;
6147 record_projection_terminal(&tx, cursor, "up_to_date")?;
6148 }
6149 (
6150 PreparedWrite::AdminSchema { name, kind, schema_json, retention_json },
6151 WritePlan::AdminSchema,
6152 ) => {
6153 tx.execute(
6154 "INSERT INTO operational_collections(
6155 name, kind, schema_json, retention_json, format_version, created_at
6156 ) VALUES(?1, ?2, ?3, ?4, 1, 0)
6157 ON CONFLICT(name) DO UPDATE SET
6158 schema_json = excluded.schema_json,
6159 retention_json = excluded.retention_json",
6160 params![name, kind, schema_json, retention_json],
6161 )?;
6162 record_projection_terminal(&tx, cursor, "up_to_date")?;
6163 }
6164 (
6165 PreparedWrite::OpStore { collection, record_key, schema_id, body },
6166 WritePlan::AppendOnlyLog,
6167 ) => {
6168 tx.execute(
6169 "INSERT INTO operational_mutations(
6170 collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
6171 ) VALUES(?1, ?2, 'append', ?3, ?4, ?5)",
6172 params![collection, record_key, body, schema_id, cursor],
6173 )?;
6174 record_projection_terminal(&tx, cursor, "up_to_date")?;
6175 }
6176 (
6177 PreparedWrite::OpStore { collection, record_key, schema_id, body },
6178 WritePlan::LatestState,
6179 ) => {
6180 tx.execute(
6181 "INSERT INTO operational_state(
6182 collection_name, record_key, payload_json, schema_id, write_cursor
6183 ) VALUES(?1, ?2, ?3, ?4, ?5)
6184 ON CONFLICT(collection_name, record_key) DO UPDATE SET
6185 payload_json = excluded.payload_json,
6186 schema_id = excluded.schema_id,
6187 write_cursor = excluded.write_cursor",
6188 params![collection, record_key, body, schema_id, cursor],
6189 )?;
6190 record_projection_terminal(&tx, cursor, "up_to_date")?;
6191 }
6192 _ => return Err(rusqlite::Error::InvalidQuery),
6193 }
6194 }
6195
6196 let dangling_edge_endpoints = {
6211 let mut last_index: HashMap<&str, usize> = HashMap::new();
6222 for (i, write) in batch.iter().enumerate() {
6223 if let PreparedWrite::Edge { logical_id: Some(lid), .. } = write {
6224 last_index.insert(lid.as_str(), i);
6225 }
6226 }
6227
6228 let mut probe = tx.prepare(
6229 "SELECT 1 FROM canonical_nodes WHERE logical_id = ?1 AND superseded_at IS NULL LIMIT 1",
6230 )?;
6231 let mut count: u64 = 0;
6232 for (i, write) in batch.iter().enumerate() {
6233 if let PreparedWrite::Edge { from, to, logical_id, .. } = write {
6234 if let Some(lid) = logical_id {
6240 let superseded_in_batch =
6241 last_index.get(lid.as_str()).is_some_and(|&last| last > i);
6242 if superseded_in_batch {
6243 continue;
6244 }
6245 }
6246 for endpoint in [from, to] {
6248 if !probe.exists(params![endpoint])? {
6249 count = count.saturating_add(1);
6250 }
6251 }
6252 }
6253 }
6254 count
6255 };
6256
6257 enforce_provenance_retention(&tx, provenance_row_cap)?;
6258 advance_projection_cursor(&tx)?;
6259
6260 tx.commit()?;
6261 Ok(dangling_edge_endpoints)
6262}
6263
6264fn load_next_cursor(connection: &Connection) -> u64 {
6265 let nodes = max_cursor(connection, "canonical_nodes").unwrap_or(0);
6266 let edges = max_cursor(connection, "canonical_edges").unwrap_or(0);
6267 let mutations = max_cursor(connection, "operational_mutations").unwrap_or(0);
6268 let state = max_cursor(connection, "operational_state").unwrap_or(0);
6269 nodes.max(edges).max(mutations).max(state)
6270}
6271
6272fn max_cursor(connection: &Connection, table: &str) -> rusqlite::Result<u64> {
6273 let sql = format!("SELECT COALESCE(MAX(write_cursor), 0) FROM {table}");
6274 connection.query_row(&sql, [], |row| row.get::<_, u64>(0))
6275}
6276
6277fn sqlite_extended_code_name(err: &rusqlite::Error) -> Option<&'static str> {
6298 let sqlite_error = err.sqlite_error()?;
6299 let extended = sqlite_error.extended_code;
6300 Some(match extended {
6301 rusqlite::ffi::SQLITE_SCHEMA => "SQLITE_SCHEMA",
6302 rusqlite::ffi::SQLITE_BUSY => "SQLITE_BUSY",
6303 rusqlite::ffi::SQLITE_LOCKED => "SQLITE_LOCKED",
6304 rusqlite::ffi::SQLITE_CORRUPT => "SQLITE_CORRUPT",
6305 rusqlite::ffi::SQLITE_NOTADB => "SQLITE_NOTADB",
6306 rusqlite::ffi::SQLITE_IOERR => "SQLITE_IOERR",
6307 rusqlite::ffi::SQLITE_FULL => "SQLITE_FULL",
6308 rusqlite::ffi::SQLITE_READONLY => "SQLITE_READONLY",
6309 rusqlite::ffi::SQLITE_CONSTRAINT => "SQLITE_CONSTRAINT",
6310 rusqlite::ffi::SQLITE_MISUSE => "SQLITE_MISUSE",
6311 rusqlite::ffi::SQLITE_INTERRUPT => "SQLITE_INTERRUPT",
6312 rusqlite::ffi::SQLITE_NOMEM => "SQLITE_NOMEM",
6313 rusqlite::ffi::SQLITE_PERM => "SQLITE_PERM",
6314 rusqlite::ffi::SQLITE_ABORT => "SQLITE_ABORT",
6315 rusqlite::ffi::SQLITE_PROTOCOL => "SQLITE_PROTOCOL",
6316 rusqlite::ffi::SQLITE_RANGE => "SQLITE_RANGE",
6317 rusqlite::ffi::SQLITE_TOOBIG => "SQLITE_TOOBIG",
6318 rusqlite::ffi::SQLITE_MISMATCH => "SQLITE_MISMATCH",
6319 rusqlite::ffi::SQLITE_AUTH => "SQLITE_AUTH",
6320 rusqlite::ffi::SQLITE_NOTFOUND => "SQLITE_NOTFOUND",
6321 rusqlite::ffi::SQLITE_CANTOPEN => "SQLITE_CANTOPEN",
6322 _ => "SQLITE_UNKNOWN",
6323 })
6324}
6325
6326fn sqlite_extended_code_name_from_int(extended: i32) -> &'static str {
6327 match extended {
6328 rusqlite::ffi::SQLITE_SCHEMA => "SQLITE_SCHEMA",
6329 rusqlite::ffi::SQLITE_BUSY => "SQLITE_BUSY",
6330 rusqlite::ffi::SQLITE_LOCKED => "SQLITE_LOCKED",
6331 rusqlite::ffi::SQLITE_CORRUPT => "SQLITE_CORRUPT",
6332 rusqlite::ffi::SQLITE_NOTADB => "SQLITE_NOTADB",
6333 rusqlite::ffi::SQLITE_IOERR => "SQLITE_IOERR",
6334 rusqlite::ffi::SQLITE_FULL => "SQLITE_FULL",
6335 rusqlite::ffi::SQLITE_READONLY => "SQLITE_READONLY",
6336 rusqlite::ffi::SQLITE_CONSTRAINT => "SQLITE_CONSTRAINT",
6337 rusqlite::ffi::SQLITE_MISUSE => "SQLITE_MISUSE",
6338 rusqlite::ffi::SQLITE_INTERRUPT => "SQLITE_INTERRUPT",
6339 rusqlite::ffi::SQLITE_NOMEM => "SQLITE_NOMEM",
6340 rusqlite::ffi::SQLITE_PERM => "SQLITE_PERM",
6341 rusqlite::ffi::SQLITE_ABORT => "SQLITE_ABORT",
6342 rusqlite::ffi::SQLITE_PROTOCOL => "SQLITE_PROTOCOL",
6343 rusqlite::ffi::SQLITE_RANGE => "SQLITE_RANGE",
6344 rusqlite::ffi::SQLITE_TOOBIG => "SQLITE_TOOBIG",
6345 rusqlite::ffi::SQLITE_MISMATCH => "SQLITE_MISMATCH",
6346 rusqlite::ffi::SQLITE_AUTH => "SQLITE_AUTH",
6347 rusqlite::ffi::SQLITE_NOTFOUND => "SQLITE_NOTFOUND",
6348 rusqlite::ffi::SQLITE_CANTOPEN => "SQLITE_CANTOPEN",
6349 _ => "SQLITE_UNKNOWN",
6350 }
6351}
6352
6353fn map_open_sqlite_error(err: rusqlite::Error, stage: OpenStage) -> EngineOpenError {
6354 let Some(sqlite_error) = err.sqlite_error() else {
6355 return EngineOpenError::Io { message: "could not open database".to_string() };
6356 };
6357 match sqlite_error.extended_code {
6358 rusqlite::ffi::SQLITE_CORRUPT | rusqlite::ffi::SQLITE_NOTADB => {
6359 EngineOpenError::Corruption(CorruptionDetail {
6360 kind: match stage {
6361 OpenStage::WalReplay => CorruptionKind::WalReplayFailure,
6362 OpenStage::HeaderProbe => CorruptionKind::HeaderMalformed,
6363 OpenStage::SchemaProbe => CorruptionKind::SchemaInconsistent,
6364 OpenStage::EmbedderIdentity => CorruptionKind::EmbedderIdentityDrift,
6365 },
6366 stage,
6367 locator: CorruptionLocator::OpaqueSqliteError {
6368 sqlite_extended_code: sqlite_error.extended_code,
6369 },
6370 recovery_hint: RecoveryHint {
6371 code: match stage {
6372 OpenStage::WalReplay => "E_CORRUPT_WAL_REPLAY",
6373 OpenStage::HeaderProbe => "E_CORRUPT_HEADER",
6374 OpenStage::SchemaProbe => "E_CORRUPT_SCHEMA",
6375 OpenStage::EmbedderIdentity => "E_CORRUPT_EMBEDDER_IDENTITY",
6376 },
6377 doc_anchor: match stage {
6378 OpenStage::WalReplay => "design/recovery.md#wal-replay-failures",
6379 OpenStage::HeaderProbe => "design/recovery.md#header-malformed",
6380 OpenStage::SchemaProbe => "design/recovery.md#schema-inconsistent",
6381 OpenStage::EmbedderIdentity => "design/recovery.md#embedder-identity-drift",
6382 },
6383 },
6384 })
6385 }
6386 _ => EngineOpenError::Io { message: "could not open database".to_string() },
6387 }
6388}
6389
6390fn emit_open_error_event(subscriber: &Arc<dyn lifecycle::Subscriber>, err: &EngineOpenError) {
6391 if let EngineOpenError::Corruption(detail) = err {
6392 let code = match detail.locator {
6393 CorruptionLocator::OpaqueSqliteError { sqlite_extended_code } => {
6394 Some(sqlite_extended_code_name_from_int(sqlite_extended_code))
6395 }
6396 _ => None,
6397 };
6398 let event = lifecycle::Event {
6399 phase: lifecycle::Phase::Failed,
6400 source: lifecycle::EventSource::SqliteInternal,
6401 category: lifecycle::EventCategory::Corruption,
6402 code,
6403 };
6404 subscriber.on_event(&event);
6405 }
6406}
6407
6408#[allow(clippy::vec_box)]
6423fn install_profile_callback(
6424 connection: &Connection,
6425 subscribers: &Arc<lifecycle::SubscriberRegistry>,
6426 profiling_enabled: &Arc<AtomicBool>,
6427 slow_threshold_ms: &Arc<AtomicU64>,
6428 contexts: &mut Vec<Box<ProfileContext>>,
6429) {
6430 let mut ctx = Box::new(ProfileContext {
6431 subscribers: Arc::clone(subscribers),
6432 profiling_enabled: Arc::clone(profiling_enabled),
6433 slow_threshold_ms: Arc::clone(slow_threshold_ms),
6434 });
6435 let ctx_ptr: *mut ProfileContext = &mut *ctx;
6436
6437 unsafe {
6449 rusqlite::ffi::sqlite3_profile(
6450 connection.handle(),
6451 Some(profile_callback_trampoline),
6452 ctx_ptr.cast::<std::ffi::c_void>(),
6453 );
6454 }
6455 contexts.push(ctx);
6456}
6457
6458fn uninstall_profile_callback(connection: &Connection) {
6462 unsafe {
6465 rusqlite::ffi::sqlite3_profile(connection.handle(), None, std::ptr::null_mut());
6466 }
6467}
6468
6469fn apply_perf_experiment_writer_pragmas(connection: &Connection) {
6507 if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
6508 return;
6509 }
6510 let raw = match std::env::var("FATHOMDB_PERF_WRITER_PRAGMAS") {
6511 Ok(s) if !s.is_empty() => s,
6512 _ => return,
6513 };
6514 for entry in raw.split(',') {
6515 let entry = entry.trim();
6516 if entry.is_empty() {
6517 continue;
6518 }
6519 let (name, value) = match entry.split_once('=') {
6520 Some((n, v)) => (n.trim(), v.trim()),
6521 None => {
6522 eprintln!("perf-experiment: bad writer pragma entry (expect name=value): {entry}");
6523 continue;
6524 }
6525 };
6526 if name.is_empty() {
6527 eprintln!("perf-experiment: empty pragma name in writer entry: {entry}");
6528 continue;
6529 }
6530 match connection.pragma_update(None, name, value) {
6531 Ok(()) => {
6532 eprintln!(
6533 "perf-experiment: applied PRAGMA {name}={value} on writer (pre-migration)"
6534 );
6535 }
6536 Err(err) => {
6537 eprintln!("perf-experiment: writer PRAGMA {name}={value} failed: {err}");
6538 }
6539 }
6540 }
6541}
6542
6543fn apply_perf_experiment_reader_pragmas(connection: &Connection) {
6544 if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
6545 return;
6546 }
6547 let raw = match std::env::var("FATHOMDB_PERF_READER_PRAGMAS") {
6548 Ok(s) if !s.is_empty() => s,
6549 _ => return,
6550 };
6551 for entry in raw.split(',') {
6552 let entry = entry.trim();
6553 if entry.is_empty() {
6554 continue;
6555 }
6556 let (name, value) = match entry.split_once('=') {
6557 Some((n, v)) => (n.trim(), v.trim()),
6558 None => {
6559 eprintln!("perf-experiment: bad pragma entry (expect name=value): {entry}");
6560 continue;
6561 }
6562 };
6563 if name.is_empty() {
6564 eprintln!("perf-experiment: empty pragma name in entry: {entry}");
6565 continue;
6566 }
6567 match connection.pragma_update(None, name, value) {
6568 Ok(()) => {
6569 eprintln!("perf-experiment: applied PRAGMA {name}={value} on reader");
6570 }
6571 Err(err) => {
6572 eprintln!("perf-experiment: PRAGMA {name}={value} failed: {err}");
6573 }
6574 }
6575 }
6576}
6577
6578fn configure_reader_lookaside(connection: &Connection) -> std::os::raw::c_int {
6579 unsafe {
6589 rusqlite::ffi::sqlite3_db_config(
6590 connection.handle(),
6591 rusqlite::ffi::SQLITE_DBCONFIG_LOOKASIDE,
6592 std::ptr::null_mut::<std::ffi::c_void>(),
6593 READER_LOOKASIDE_SLOT_SIZE,
6594 READER_LOOKASIDE_SLOT_COUNT,
6595 )
6596 }
6597}
6598
6599#[cfg(debug_assertions)]
6607fn read_lookaside_used_hiwtr(connection: &Connection) -> std::os::raw::c_int {
6608 let mut current: std::os::raw::c_int = 0;
6609 let mut hiwtr: std::os::raw::c_int = 0;
6610 unsafe {
6613 rusqlite::ffi::sqlite3_db_status(
6614 connection.handle(),
6615 rusqlite::ffi::SQLITE_DBSTATUS_LOOKASIDE_USED,
6616 &mut current,
6617 &mut hiwtr,
6618 0,
6619 );
6620 }
6621 hiwtr
6622}
6623
6624#[cfg(debug_assertions)]
6631fn read_cache_status(
6632 connection: &Connection,
6633) -> (std::os::raw::c_int, std::os::raw::c_int, std::os::raw::c_int) {
6634 let mut hit_current: std::os::raw::c_int = 0;
6635 let mut hit_hiwtr: std::os::raw::c_int = 0;
6636 let mut miss_current: std::os::raw::c_int = 0;
6637 let mut miss_hiwtr: std::os::raw::c_int = 0;
6638 let mut used_current: std::os::raw::c_int = 0;
6639 let mut used_hiwtr: std::os::raw::c_int = 0;
6640 unsafe {
6644 rusqlite::ffi::sqlite3_db_status(
6645 connection.handle(),
6646 rusqlite::ffi::SQLITE_DBSTATUS_CACHE_HIT,
6647 &mut hit_current,
6648 &mut hit_hiwtr,
6649 0,
6650 );
6651 rusqlite::ffi::sqlite3_db_status(
6652 connection.handle(),
6653 rusqlite::ffi::SQLITE_DBSTATUS_CACHE_MISS,
6654 &mut miss_current,
6655 &mut miss_hiwtr,
6656 0,
6657 );
6658 rusqlite::ffi::sqlite3_db_status(
6659 connection.handle(),
6660 rusqlite::ffi::SQLITE_DBSTATUS_CACHE_USED,
6661 &mut used_current,
6662 &mut used_hiwtr,
6663 0,
6664 );
6665 }
6666 (hit_current, miss_current, used_current)
6670}
6671
6672unsafe extern "C" fn profile_callback_trampoline(
6686 user_data: *mut std::ffi::c_void,
6687 sql: *const std::os::raw::c_char,
6688 nanoseconds: u64,
6689) {
6690 if user_data.is_null() || sql.is_null() {
6691 return;
6692 }
6693 let ctx = unsafe { &*(user_data.cast::<ProfileContext>()) };
6694 let sql_text = match unsafe { std::ffi::CStr::from_ptr(sql) }.to_str() {
6695 Ok(s) => s,
6696 Err(_) => return,
6697 };
6698
6699 let wall_clock_ms = nanoseconds / 1_000_000;
6700
6701 if ctx.profiling_enabled.load(Ordering::Relaxed) {
6702 let record = lifecycle::ProfileRecord {
6703 wall_clock_ms,
6704 step_count: 0,
6710 cache_delta: 0,
6711 };
6712 ctx.subscribers.dispatch_profile(&record);
6713 }
6714
6715 let threshold = ctx.slow_threshold_ms.load(Ordering::Relaxed);
6716 if wall_clock_ms > threshold {
6717 let signal = lifecycle::SlowStatement { statement: sql_text.to_string(), wall_clock_ms };
6718 ctx.subscribers.dispatch_slow_statement(&signal);
6719 }
6720}
6721
6722#[cfg(test)]
6723mod tests {
6724 use super::{resolve_source_type, Engine, PreparedWrite, KIND_TO_SOURCE_TYPE_CASE_SQL};
6725 use rusqlite::Connection;
6726 use tempfile::TempDir;
6727
6728 #[test]
6738 fn resolve_source_type_drift_check() {
6739 let kinds = ["email", "article", "paper", "meeting", "note", "todo", "doc"];
6740
6741 let want: &[(&str, &str)] = &[
6744 ("email", "email"),
6745 ("article", "article"),
6746 ("paper", "paper"),
6747 ("meeting", "meeting"),
6748 ("note", "note"),
6749 ("todo", "todo"),
6750 ("doc", "article"),
6751 ];
6752 for (kind, expected) in want {
6753 let got = resolve_source_type(kind).unwrap_or_else(|_| {
6754 panic!("resolve_source_type({kind}) returned Err; want Ok({expected})")
6755 });
6756 assert_eq!(got, *expected, "Rust helper drift for kind={kind}");
6757 }
6758 assert!(
6759 resolve_source_type("banana").is_err(),
6760 "unknown kind must surface as writer error"
6761 );
6762
6763 let conn = Connection::open_in_memory().expect("in-memory sqlite");
6768 conn.execute_batch("CREATE TABLE s(kind TEXT NOT NULL)").expect("create s");
6769 for kind in &kinds {
6770 conn.execute("INSERT INTO s(kind) VALUES (?1)", [kind]).expect("insert kind");
6771 }
6772 let sql = format!("SELECT s.kind, {KIND_TO_SOURCE_TYPE_CASE_SQL} FROM s");
6773 let mut stmt = conn.prepare(&sql).expect("prepare CASE");
6774 let rows: Vec<(String, String)> = stmt
6775 .query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)))
6776 .expect("query")
6777 .map(|r| r.expect("row"))
6778 .collect();
6779 assert_eq!(rows.len(), kinds.len(), "row count drift");
6780 for (kind, sql_result) in &rows {
6781 let rust_result = resolve_source_type(kind).expect("known kind");
6782 assert_eq!(
6783 sql_result, rust_result,
6784 "SQL CASE vs Rust helper drift for kind={kind}: SQL={sql_result}, Rust={rust_result}"
6785 );
6786 }
6787 }
6788
6789 #[test]
6790 fn write_advances_cursor() {
6791 let dir = TempDir::new().unwrap();
6792 let opened = Engine::open(dir.path().join("rewrite.sqlite")).expect("engine should open");
6793 let receipt = opened
6794 .engine
6795 .write(&[PreparedWrite::Node {
6796 kind: "doc".to_string(),
6797 body: "hello".to_string(),
6798 source_id: None,
6799 logical_id: None,
6800 }])
6801 .expect("write should succeed");
6802
6803 assert_eq!(receipt.cursor, 1);
6804 }
6805}