pub mod lifecycle;
mod pcache2;
use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
use std::error::Error;
use std::fmt::{Display, Formatter};
use std::fs::{File, OpenOptions};
use std::io::{Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::mpsc::{self, Receiver, SyncSender};
use std::sync::Once;
use std::sync::{Arc, Condvar, Mutex};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use fathomdb_embedder::EmbedderEvent;
#[cfg(feature = "operator")]
use fathomdb_embedder::MeanRecomputeTrigger;
use fathomdb_embedder_api::{Embedder, EmbedderError as RuntimeEmbedderError, EmbedderIdentity};
use fathomdb_query::compile_text_query;
use fathomdb_schema::{
migrate_with_event_sink, MigrationError as SchemaMigrationError, MigrationStepReport,
LOCK_SUFFIX, MIGRATIONS, SCHEMA_VERSION,
};
#[cfg(feature = "operator")]
use fathomdb_schema::CANONICAL_TABLES;
use jsonschema::JSONSchema;
use rusqlite::{params, Connection, OptionalExtension};
use serde_json::Value;
#[cfg(feature = "operator")]
use sha2::Digest;
use sqlite_vec::sqlite3_vec_init;
#[cfg(unix)]
use std::os::unix::fs::OpenOptionsExt;
const DEFAULT_EMBEDDER_NAME: &str = "fathomdb-bge-small-en-v1.5";
const DEFAULT_EMBEDDER_REVISION: &str = "5c38ec7c405ec4b44b94cc5a9bb96e735b38267a";
const DEFAULT_EMBEDDER_DIMENSION: u32 = 384;
const BGE_SMALL_EMBEDDER_NAME: &str = "fathomdb-bge-small-en-v1.5";
const DEFAULT_SLOW_THRESHOLD_MS: u64 = 100;
const DEFAULT_VECTOR_PROFILE: &str = "default";
const DEFAULT_VECTOR_PARTITION: &str = "vector_default";
#[cfg(feature = "operator")]
const REBUILD_DRAIN_TIMEOUT_MS: u64 = 30_000;
const SEARCH_INDEX_TOKENIZER_SCHEMA_VERSION: u32 = 11;
const SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY: &str =
"search_index_tokenizer_reproject_complete";
const DEFAULT_PROVENANCE_ROW_CAP: u64 = 1_000_000;
const PROJECTION_CURSOR_KEY: &str = "projection_cursor";
const PROJECTION_WORKERS: usize = 2;
const DEFAULT_EMBED_TIMEOUT_MS: u64 = 30_000;
const DEFAULT_EMBED_CIRCUIT_THRESHOLD: u64 = 8;
const PROJECTION_COMMIT_BATCH: usize = 16;
const PROJECTION_INFLIGHT_LIMIT: usize = PROJECTION_WORKERS * PROJECTION_COMMIT_BATCH;
const PROJECTION_SCAN_FETCH: usize = PROJECTION_INFLIGHT_LIMIT;
const DEFAULT_PROJECTION_RETRY_DELAYS_MS: [u64; 3] = [1_000, 4_000, 16_000];
const READER_POOL_SIZE: usize = 8;
const READER_LOOKASIDE_SLOT_SIZE: std::os::raw::c_int = 1200;
const READER_LOOKASIDE_SLOT_COUNT: std::os::raw::c_int = 500;
pub struct Engine {
path: PathBuf,
next_cursor: AtomicU64,
closed: AtomicBool,
lock: Mutex<Option<File>>,
connection: Mutex<Option<Connection>>,
reader_pool: ReaderWorkerPool,
counters: lifecycle::Counters,
subscribers: Arc<lifecycle::SubscriberRegistry>,
profiling_enabled: Arc<AtomicBool>,
slow_threshold_ms: Arc<AtomicU64>,
runtime_embedder: Option<Arc<dyn Embedder>>,
runtime_embedder_identity: EmbedderIdentity,
projection_runtime: ProjectionRuntime,
provenance_row_cap: AtomicU64,
#[allow(clippy::vec_box)]
profile_contexts: Mutex<Vec<Box<ProfileContext>>>,
#[allow(dead_code)]
reader_lookaside_rcs: Vec<i32>,
#[cfg(debug_assertions)]
force_next_commit_failure: AtomicBool,
}
#[derive(Clone, Debug)]
struct ProjectionJob {
cursor: u64,
kind: String,
body: String,
}
#[derive(Debug, Default)]
struct ProjectionRuntimeState {
active_jobs: usize,
queued_jobs: usize,
frozen: bool,
pending_scan: bool,
stopping: bool,
in_flight: BTreeSet<u64>,
}
struct ProjectionRuntimeShared {
path: PathBuf,
embedder: Option<Arc<dyn Embedder>>,
embedder_identity: EmbedderIdentity,
state: Mutex<ProjectionRuntimeState>,
state_cvar: Condvar,
queue: Mutex<VecDeque<ProjectionJob>>,
queue_cvar: Condvar,
retry_delays_ms: Mutex<Vec<u64>>,
embed_timeout_ms: AtomicU64,
embed_serialize: Mutex<()>,
live_embed_threads: Arc<AtomicU64>,
embed_circuit_open: AtomicBool,
embed_circuit_threshold: AtomicU64,
mean_accumulator: Mutex<Option<MeanAccumulator>>,
pending_events: Mutex<Vec<EmbedderEvent>>,
commit_gate: Mutex<()>,
search_limit_override: AtomicUsize,
recency_reweight_enabled: AtomicBool,
vector_stage_only_for_test: AtomicBool,
#[cfg(debug_assertions)]
force_recompute_failure: AtomicBool,
}
impl std::fmt::Debug for ProjectionRuntimeShared {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ProjectionRuntimeShared")
.field("path", &self.path)
.field("embedder_identity", &self.embedder_identity)
.finish_non_exhaustive()
}
}
#[derive(Debug)]
struct ProjectionRuntime {
shared: Arc<ProjectionRuntimeShared>,
dispatcher: Mutex<Option<JoinHandle<()>>>,
workers: Mutex<Vec<JoinHandle<()>>>,
}
impl std::fmt::Debug for Engine {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Engine")
.field("path", &self.path)
.field("closed", &self.closed.load(Ordering::SeqCst))
.field("runtime_embedder_identity", &self.runtime_embedder_identity)
.finish_non_exhaustive()
}
}
#[derive(Debug)]
struct ProfileContext {
subscribers: Arc<lifecycle::SubscriberRegistry>,
profiling_enabled: Arc<AtomicBool>,
slow_threshold_ms: Arc<AtomicU64>,
}
struct ReaderWorkerPool {
senders: Vec<SyncSender<ReaderRequest>>,
handles: Mutex<Option<Vec<JoinHandle<()>>>>,
next: AtomicUsize,
shutdown: AtomicBool,
live_workers: Arc<AtomicUsize>,
}
impl std::fmt::Debug for ReaderWorkerPool {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ReaderWorkerPool")
.field("worker_count", &self.senders.len())
.field("live_workers", &self.live_workers.load(Ordering::Relaxed))
.field("shutdown", &self.shutdown.load(Ordering::Relaxed))
.finish()
}
}
enum ReaderRequest {
Search {
compiled: fathomdb_query::CompiledQuery,
query_vector: Option<String>,
query_vector_bin: Option<String>,
search_limit: usize,
filter: Option<Box<SearchFilter>>,
recency_enabled: bool,
vector_stage_only: bool,
respond: SyncSender<ReaderResponse>,
},
GetById {
logical_ids: Vec<String>,
respond: SyncSender<rusqlite::Result<Vec<Option<NodeRecord>>>>,
},
ReadCollection {
collection: String,
after_id: Option<i64>,
limit: usize,
respond: SyncSender<rusqlite::Result<Vec<OpStoreRow>>>,
},
Shutdown,
#[cfg(debug_assertions)]
LookasideStatus {
respond: SyncSender<i32>,
},
#[cfg(debug_assertions)]
CacheStatus {
snapshot_label: String,
respond: SyncSender<(String, i32, i32, i32)>,
},
}
type ReaderResponse = rusqlite::Result<(u64, Option<SoftFallback>, Vec<SearchHit>)>;
#[cfg(debug_assertions)]
#[doc(hidden)]
#[derive(Clone, Debug)]
pub struct CacheStatusReply {
pub worker_idx: usize,
pub snapshot_label: String,
pub cache_hit: i32,
pub cache_miss: i32,
pub cache_used_bytes: i32,
}
const READER_WORKER_CHANNEL_CAPACITY: usize = 4;
impl ReaderWorkerPool {
fn new(connections: Vec<Connection>) -> Self {
let live_workers = Arc::new(AtomicUsize::new(0));
let mut senders = Vec::with_capacity(connections.len());
let mut handles = Vec::with_capacity(connections.len());
for (idx, connection) in connections.into_iter().enumerate() {
let (tx, rx) = mpsc::sync_channel::<ReaderRequest>(READER_WORKER_CHANNEL_CAPACITY);
let live = Arc::clone(&live_workers);
let handle = thread::Builder::new()
.name(format!("fathomdb-reader-{idx}"))
.spawn(move || reader_worker_loop(connection, rx, live))
.expect("spawn reader worker");
senders.push(tx);
handles.push(handle);
}
Self {
senders,
handles: Mutex::new(Some(handles)),
next: AtomicUsize::new(0),
shutdown: AtomicBool::new(false),
live_workers,
}
}
fn worker_count(&self) -> usize {
self.senders.len()
}
fn live_count(&self) -> usize {
self.live_workers.load(Ordering::SeqCst)
}
#[cfg(debug_assertions)]
fn lookaside_used_per_worker(&self) -> Vec<i32> {
let mut results = Vec::with_capacity(self.senders.len());
for sender in &self.senders {
let (tx, rx) = mpsc::sync_channel::<i32>(1);
if sender.send(ReaderRequest::LookasideStatus { respond: tx }).is_ok() {
results.push(rx.recv().unwrap_or(-1));
} else {
results.push(-1);
}
}
results
}
#[cfg(debug_assertions)]
fn cache_status_per_worker(&self, snapshot_label: &str) -> Vec<CacheStatusReply> {
let mut results = Vec::with_capacity(self.senders.len());
for (idx, sender) in self.senders.iter().enumerate() {
let (tx, rx) = mpsc::sync_channel::<(String, i32, i32, i32)>(1);
let request = ReaderRequest::CacheStatus {
snapshot_label: snapshot_label.to_string(),
respond: tx,
};
if sender.send(request).is_ok() {
if let Ok((label, hit, miss, used)) = rx.recv() {
results.push(CacheStatusReply {
worker_idx: idx,
snapshot_label: label,
cache_hit: hit,
cache_miss: miss,
cache_used_bytes: used,
});
continue;
}
}
results.push(CacheStatusReply {
worker_idx: idx,
snapshot_label: snapshot_label.to_string(),
cache_hit: -1,
cache_miss: -1,
cache_used_bytes: -1,
});
}
results
}
fn dispatch(&self, request: ReaderRequest) -> Result<(), ReaderRequest> {
if self.shutdown.load(Ordering::Relaxed) {
return Err(request);
}
let n = self.senders.len();
if n == 0 {
return Err(request);
}
let idx = self.next.fetch_add(1, Ordering::Relaxed) % n;
self.senders[idx].send(request).map_err(|err| err.0)
}
fn shutdown(&self) {
if self.shutdown.swap(true, Ordering::SeqCst) {
return;
}
for sender in &self.senders {
let _ = sender.send(ReaderRequest::Shutdown);
}
if let Ok(mut slot) = self.handles.lock() {
if let Some(handles) = slot.take() {
for handle in handles {
let _ = handle.join();
}
}
}
}
}
impl Drop for ReaderWorkerPool {
fn drop(&mut self) {
self.shutdown();
}
}
fn reader_worker_loop(
mut connection: Connection,
rx: Receiver<ReaderRequest>,
live_workers: Arc<AtomicUsize>,
) {
live_workers.fetch_add(1, Ordering::SeqCst);
struct LiveGuard(Arc<AtomicUsize>);
impl Drop for LiveGuard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
let _guard = LiveGuard(live_workers);
while let Ok(request) = rx.recv() {
match request {
ReaderRequest::Shutdown => break,
ReaderRequest::Search {
compiled,
query_vector,
query_vector_bin,
search_limit,
filter,
recency_enabled,
vector_stage_only,
respond,
} => {
let result = read_search_in_tx(
&mut connection,
&compiled,
query_vector.as_deref(),
query_vector_bin.as_deref(),
search_limit,
filter.as_deref(),
recency_enabled,
vector_stage_only,
);
let _ = respond.send(result);
}
ReaderRequest::GetById { logical_ids, respond } => {
let result = read_get_by_id_in_tx(&mut connection, &logical_ids);
let _ = respond.send(result);
}
ReaderRequest::ReadCollection { collection, after_id, limit, respond } => {
let result = read_collection_in_tx(&mut connection, &collection, after_id, limit);
let _ = respond.send(result);
}
#[cfg(debug_assertions)]
ReaderRequest::LookasideStatus { respond } => {
let _ = respond.send(read_lookaside_used_hiwtr(&connection));
}
#[cfg(debug_assertions)]
ReaderRequest::CacheStatus { snapshot_label, respond } => {
let (hit, miss, used) = read_cache_status(&connection);
let _ = respond.send((snapshot_label, hit, miss, used));
}
}
}
uninstall_profile_callback(&connection);
drop(connection);
}
impl ProjectionRuntime {
fn new(
path: PathBuf,
embedder: Option<Arc<dyn Embedder>>,
embedder_identity: EmbedderIdentity,
mean_already_pinned: bool,
) -> Self {
let mc_required = identity_requires_mean_centering(&embedder_identity);
let mean_accumulator = if mc_required && !mean_already_pinned {
Some(MeanAccumulator::new(embedder_identity.dimension as usize))
} else {
None
};
let shared = Arc::new(ProjectionRuntimeShared {
path,
embedder,
embedder_identity,
state: Mutex::new(ProjectionRuntimeState::default()),
state_cvar: Condvar::new(),
queue: Mutex::new(VecDeque::new()),
queue_cvar: Condvar::new(),
retry_delays_ms: Mutex::new(DEFAULT_PROJECTION_RETRY_DELAYS_MS.to_vec()),
embed_timeout_ms: AtomicU64::new(DEFAULT_EMBED_TIMEOUT_MS),
embed_serialize: Mutex::new(()),
live_embed_threads: Arc::new(AtomicU64::new(0)),
embed_circuit_open: AtomicBool::new(false),
embed_circuit_threshold: AtomicU64::new(DEFAULT_EMBED_CIRCUIT_THRESHOLD),
mean_accumulator: Mutex::new(mean_accumulator),
pending_events: Mutex::new(Vec::new()),
commit_gate: Mutex::new(()),
search_limit_override: AtomicUsize::new(SEARCH_RERANK_LIMIT),
recency_reweight_enabled: AtomicBool::new(false),
vector_stage_only_for_test: AtomicBool::new(false),
#[cfg(debug_assertions)]
force_recompute_failure: AtomicBool::new(false),
});
let dispatcher_shared = Arc::clone(&shared);
let dispatcher = thread::spawn(move || projection_dispatcher_loop(dispatcher_shared));
let mut workers = Vec::with_capacity(PROJECTION_WORKERS);
for _ in 0..PROJECTION_WORKERS {
let worker_shared = Arc::clone(&shared);
workers.push(thread::spawn(move || projection_worker_loop(worker_shared)));
}
Self { shared, dispatcher: Mutex::new(Some(dispatcher)), workers: Mutex::new(workers) }
}
fn notify_new_work(&self) {
if let Ok(mut state) = self.shared.state.lock() {
state.pending_scan = true;
self.shared.state_cvar.notify_all();
}
}
fn set_frozen(&self, frozen: bool) {
if let Ok(mut state) = self.shared.state.lock() {
state.frozen = frozen;
if !frozen {
state.pending_scan = true;
}
self.shared.state_cvar.notify_all();
}
}
fn wait_for_idle(&self, timeout_ms: u64) -> bool {
let deadline = Instant::now() + Duration::from_millis(timeout_ms);
let mut state = match self.shared.state.lock() {
Ok(state) => state,
Err(_) => return false,
};
loop {
if state.active_jobs == 0 && state.queued_jobs == 0 {
drop(state);
if !database_has_pending_projection_work(&self.shared.path).unwrap_or(true) {
return true;
}
state = match self.shared.state.lock() {
Ok(state) => state,
Err(_) => return false,
};
}
let now = Instant::now();
if now >= deadline {
return false;
}
let wait = deadline.saturating_duration_since(now);
let Ok((next_state, _)) = self.shared.state_cvar.wait_timeout(state, wait) else {
return false;
};
state = next_state;
}
}
fn set_retry_delays_for_test(&self, delays_ms: &[u64]) {
if let Ok(mut delays) = self.shared.retry_delays_ms.lock() {
*delays = delays_ms.to_vec();
}
}
fn set_embed_timeout_ms_for_test(&self, timeout_ms: u64) {
self.shared.embed_timeout_ms.store(timeout_ms, Ordering::Relaxed);
}
fn set_embed_circuit_threshold_for_test(&self, threshold: u64) {
self.shared.embed_circuit_threshold.store(threshold, Ordering::Relaxed);
}
fn embed_circuit_open_for_test(&self) -> bool {
self.shared.embed_circuit_open.load(Ordering::Relaxed)
}
fn stop(&self) {
if let Ok(mut state) = self.shared.state.lock() {
if state.stopping {
return;
}
state.stopping = true;
state.pending_scan = false;
self.shared.state_cvar.notify_all();
}
if let Ok(mut queue) = self.shared.queue.lock() {
queue.clear();
self.shared.queue_cvar.notify_all();
}
if let Ok(mut dispatcher) = self.dispatcher.lock() {
if let Some(handle) = dispatcher.take() {
let _ = handle.join();
}
}
if let Ok(mut workers) = self.workers.lock() {
for handle in workers.drain(..) {
let _ = handle.join();
}
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OpenReport {
pub schema_version_before: u32,
pub schema_version_after: u32,
pub migration_steps: Vec<MigrationStepReport>,
pub embedder_warmup_ms: u64,
pub query_backend: &'static str,
pub default_embedder: EmbedderIdentity,
pub embedder_download_ms: Option<u64>,
pub embedder_events: Vec<EmbedderEvent>,
pub embedder_mean_centering_required: bool,
pub embedder_mean_vec_pinned: bool,
}
#[derive(Debug)]
pub struct OpenedEngine {
pub engine: Engine,
pub report: OpenReport,
}
#[derive(Clone, Debug)]
struct LoaderInfo {
download_ms: Option<u64>,
events: Vec<EmbedderEvent>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WriteReceipt {
pub cursor: u64,
pub row_cursors: Vec<u64>,
pub dangling_edge_endpoints: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SoftFallback {
pub branch: SoftFallbackBranch,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SoftFallbackBranch {
Vector,
Text,
}
#[derive(Clone, Debug, PartialEq)]
pub struct SearchHit {
pub id: u64,
pub kind: String,
pub body: String,
pub score: f64,
pub branch: SoftFallbackBranch,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct NodeRecord {
pub logical_id: String,
pub kind: String,
pub body: String,
pub write_cursor: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OpStoreRow {
pub id: i64,
pub collection: String,
pub record_key: String,
pub op_kind: String,
pub payload: String,
pub schema_id: Option<String>,
pub write_cursor: u64,
}
#[derive(Clone, Debug, PartialEq)]
pub struct SearchResult {
pub projection_cursor: u64,
pub soft_fallback: Option<SoftFallback>,
pub results: Vec<SearchHit>,
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct SearchFilter {
pub source_type: Option<String>,
pub kind: Option<String>,
pub created_after: Option<i64>,
pub status: Option<String>,
}
impl SearchFilter {
fn is_unfiltered(&self) -> bool {
self.source_type.is_none()
&& self.kind.is_none()
&& self.created_after.is_none()
&& self.status.is_none()
}
}
#[non_exhaustive]
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum PreparedWrite {
Node {
kind: String,
body: String,
source_id: Option<String>,
logical_id: Option<String>,
},
Edge {
kind: String,
from: String,
to: String,
source_id: Option<String>,
logical_id: Option<String>,
},
OpStore {
collection: String,
record_key: String,
schema_id: Option<String>,
body: String,
},
AdminSchema {
name: String,
kind: String,
schema_json: String,
retention_json: String,
},
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct CounterSnapshot {
pub queries: u64,
pub writes: u64,
pub write_rows: u64,
pub errors_by_code: BTreeMap<String, u64>,
pub admin_ops: u64,
pub cache_hit: u64,
pub cache_miss: u64,
}
pub use lifecycle::Subscription;
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct CorruptionDetail {
pub kind: CorruptionKind,
pub stage: OpenStage,
pub locator: CorruptionLocator,
pub recovery_hint: RecoveryHint,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CorruptionKind {
WalReplayFailure,
HeaderMalformed,
SchemaInconsistent,
EmbedderIdentityDrift,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum OpenStage {
WalReplay,
HeaderProbe,
SchemaProbe,
EmbedderIdentity,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CorruptionLocator {
FileOffset { offset: u64 },
PageId { page: u32 },
TableRow { table: &'static str, rowid: i64 },
Vec0ShadowRow { partition: &'static str, rowid: i64 },
MigrationStep { from: u32, to: u32 },
OpaqueSqliteError { sqlite_extended_code: i32 },
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct RecoveryHint {
pub code: &'static str,
pub doc_anchor: &'static str,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum EngineOpenError {
DatabaseLocked {
holder_pid: Option<u32>,
},
Corruption(CorruptionDetail),
IncompatibleSchemaVersion {
seen: u32,
supported: u32,
},
MigrationError {
schema_version_before: u32,
schema_version_current: u32,
step_id: u32,
},
EmbedderIdentityMismatch {
stored: EmbedderIdentity,
supplied: EmbedderIdentity,
},
EmbedderDimensionMismatch {
stored: u32,
supplied: u32,
},
Embedder(RuntimeEmbedderError),
Io {
message: String,
},
}
#[derive(Clone)]
pub enum EmbedderChoice {
Default,
Caller(Arc<dyn Embedder>),
None,
}
impl Display for EngineOpenError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::DatabaseLocked { holder_pid } => match holder_pid {
Some(pid) => write!(f, "database is locked by process {pid}"),
None => write!(f, "database is locked by another engine instance"),
},
Self::Corruption(detail) => {
write!(
f,
"engine corruption at {:?} stage: {}",
detail.stage, detail.recovery_hint.code
)
}
Self::IncompatibleSchemaVersion { seen, supported } => write!(
f,
"database schema version {seen} is incompatible with supported version {supported}"
),
Self::MigrationError {
schema_version_before,
schema_version_current,
step_id,
} => write!(
f,
"schema migration failed at step {step_id}; schema version remained between {schema_version_before} and {schema_version_current}"
),
Self::EmbedderIdentityMismatch { stored, supplied } => write!(
f,
"embedder identity mismatch: stored {}@{}, supplied {}@{}",
stored.name, stored.revision, supplied.name, supplied.revision,
),
Self::EmbedderDimensionMismatch { stored, supplied } => write!(
f,
"embedder vector dimension mismatch: stored {stored}, supplied {supplied}",
),
Self::Embedder(err) => match err {
RuntimeEmbedderError::Timeout => write!(f, "embedder timeout during open"),
RuntimeEmbedderError::Failed { message } => {
write!(f, "embedder failure during open: {message}")
}
},
Self::Io { message } => write!(f, "database I/O error: {message}"),
}
}
}
impl Error for EngineOpenError {}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum EngineError {
Storage,
Projection,
Vector,
Embedder,
EmbedderNotConfigured,
KindNotVectorIndexed,
EmbedderDimensionMismatch { expected: u32, actual: u32 },
Scheduler,
OpStore,
WriteValidation,
SchemaValidation,
Overloaded,
Closing,
}
impl Display for EngineError {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
match self {
Self::Storage => write!(f, "storage error"),
Self::Projection => write!(f, "projection error"),
Self::Vector => write!(f, "vector error"),
Self::Embedder => write!(f, "embedder error"),
Self::EmbedderNotConfigured => write!(f, "embedder is not configured"),
Self::KindNotVectorIndexed => write!(f, "kind is not configured for vector indexing"),
Self::EmbedderDimensionMismatch { expected, actual } => {
write!(f, "embedder dimension mismatch: expected {expected}, actual {actual}")
}
Self::Scheduler => write!(f, "scheduler error"),
Self::OpStore => write!(f, "op-store error"),
Self::WriteValidation => write!(f, "write validation error"),
Self::SchemaValidation => write!(f, "schema validation error"),
Self::Overloaded => write!(f, "engine overloaded"),
Self::Closing => write!(f, "engine is closing"),
}
}
}
impl EngineError {
fn stable_code(&self) -> &'static str {
match self {
Self::Storage => "StorageError",
Self::Projection => "ProjectionError",
Self::Vector => "VectorError",
Self::Embedder => "EmbedderError",
Self::EmbedderNotConfigured => "EmbedderNotConfiguredError",
Self::KindNotVectorIndexed => "KindNotVectorIndexedError",
Self::EmbedderDimensionMismatch { .. } => "EmbedderDimensionMismatchError",
Self::Scheduler => "SchedulerError",
Self::OpStore => "OpStoreError",
Self::WriteValidation => "WriteValidationError",
Self::SchemaValidation => "SchemaValidationError",
Self::Overloaded => "OverloadedError",
Self::Closing => "ClosingError",
}
}
}
impl Error for EngineError {}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct CheckIntegrityOpts {
pub quick: bool,
pub full: bool,
pub round_trip: bool,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum Section {
Clean,
Findings(Vec<Finding>),
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Finding {
pub code: &'static str,
pub stage: &'static str,
pub locator: CorruptionLocator,
pub doc_anchor: &'static str,
pub detail: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct IntegrityReport {
pub physical: Section,
pub logical: Section,
pub semantic: Section,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SafeExportArtifact {
pub export_path: PathBuf,
pub manifest_path: PathBuf,
pub manifest_sha256: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TraceReport {
pub source_ref: String,
pub events: Vec<TraceEvent>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TraceEvent {
pub write_cursor: u64,
pub kind: String,
pub table: &'static str,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum RebuildKind {
Projections,
Vec0,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RebuildReport {
pub kind: RebuildKind,
pub rows_invalidated: u64,
pub rows_rebuilt: u64,
pub projection_cursor_after: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ExciseReport {
pub source_ref: String,
pub nodes_excised: u64,
pub edges_excised: u64,
pub projections_invalidated: u64,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum VerifyEmbedderStatus {
Match,
IdentityMismatch,
DimensionMismatch,
BothMismatch,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct VerifyEmbedderReport {
pub stored_identity: String,
pub stored_dimension: u32,
pub supplied_identity: String,
pub supplied_dimension: u32,
pub status: VerifyEmbedderStatus,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SchemaObject {
pub name: String,
pub sql: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DumpSchemaReport {
pub user_version: u32,
pub tables: Vec<SchemaObject>,
pub indexes: Vec<SchemaObject>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TableRowCount {
pub name: String,
pub rows: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DumpRowCountsReport {
pub counts: Vec<TableRowCount>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DumpProfileReport {
pub embedder_identity: String,
pub embedder_dimension: u32,
pub vectorized_kinds: Vec<String>,
}
#[derive(Clone, Debug, PartialEq)]
pub struct MeanRecomputeReport {
pub dim: u32,
pub old_doc_count: u64,
pub doc_count_requantized: u64,
pub drift_cos_before: f32,
pub mean_was_pinned: bool,
pub elapsed_ms: u64,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum TruncateWalStatus {
Done,
Busy,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TruncateWalReport {
pub status: TruncateWalStatus,
pub busy: u32,
pub log_frames: u32,
pub checkpointed_frames: u32,
}
impl Drop for Engine {
fn drop(&mut self) {
let _ = self.close();
}
}
impl Engine {
pub fn open(path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_embedder_and_subscriber(
path,
default_embedder_identity(),
None,
None,
None,
&mut |_| {},
)
}
pub fn open_with_choice(
path: impl Into<PathBuf>,
choice: EmbedderChoice,
) -> Result<OpenedEngine, EngineOpenError> {
match choice {
EmbedderChoice::Default => Self::open_default_embedder(path),
EmbedderChoice::Caller(embedder) => {
let identity = embedder.identity();
Self::open_with_embedder_and_subscriber(
path,
identity,
Some(embedder),
None,
None,
&mut |_| {},
)
}
EmbedderChoice::None => Self::open_with_embedder_and_subscriber(
path,
default_embedder_identity(),
None,
None,
None,
&mut |_| {},
),
}
}
#[cfg(feature = "default-embedder")]
fn open_default_embedder(path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
use std::time::Instant as DownloadInstant;
let download_start = DownloadInstant::now();
let weights = fathomdb_embedder::loader::load_pinned_default_embedder().map_err(|err| {
EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
message: format!("default embedder loader: {err}"),
})
})?;
let events = weights.events.clone();
let download_ms = if weights.bytes_downloaded > 0 {
Some(u64::try_from(download_start.elapsed().as_millis()).unwrap_or(u64::MAX))
} else {
None
};
let embedder =
fathomdb_embedder::CandleBgeEmbedder::new_from_weights(weights).map_err(|err| {
EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
message: format!("default embedder construct: {err}"),
})
})?;
let embedder: Arc<dyn Embedder> = Arc::new(embedder);
let identity = embedder.identity();
let loader_info = LoaderInfo { download_ms, events };
Self::open_with_embedder_and_subscriber(
path,
identity,
Some(embedder),
Some(loader_info),
None,
&mut |_| {},
)
}
#[cfg(not(feature = "default-embedder"))]
fn open_default_embedder(_path: impl Into<PathBuf>) -> Result<OpenedEngine, EngineOpenError> {
Err(EngineOpenError::Embedder(RuntimeEmbedderError::Failed {
message: "EmbedderChoice::Default requires the `default-embedder` Cargo feature"
.to_string(),
}))
}
pub fn open_with_migration_event_sink(
path: impl Into<PathBuf>,
mut emit_migration_event: impl FnMut(&MigrationStepReport),
) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_embedder_and_subscriber(
path,
default_embedder_identity(),
None,
None,
None,
&mut emit_migration_event,
)
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn open_with_migrations_for_test(
path: impl Into<PathBuf>,
migrations: &'static [fathomdb_schema::Migration],
mut emit_migration_event: impl FnMut(&MigrationStepReport),
) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_migrations(
path,
migrations,
default_embedder_identity(),
None,
None,
&mut emit_migration_event,
None,
)
}
#[doc(hidden)]
pub fn open_with_subscriber_for_test(
path: impl Into<PathBuf>,
subscriber: Arc<dyn lifecycle::Subscriber>,
) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_embedder_and_subscriber(
path,
default_embedder_identity(),
None,
None,
Some(subscriber),
&mut |_| {},
)
}
#[doc(hidden)]
pub fn open_without_embedder_for_test(
path: impl Into<PathBuf>,
) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_embedder_and_subscriber(
path,
default_embedder_identity(),
None,
None,
None,
&mut |_| {},
)
}
#[doc(hidden)]
pub fn open_with_embedder_for_test(
path: impl Into<PathBuf>,
embedder: Arc<dyn Embedder>,
) -> Result<OpenedEngine, EngineOpenError> {
let identity = embedder.identity();
Self::open_with_embedder_and_subscriber(
path,
identity,
Some(embedder),
None,
None,
&mut |_| {},
)
}
fn open_with_embedder_and_subscriber(
path: impl Into<PathBuf>,
embedder_identity: EmbedderIdentity,
runtime_embedder: Option<Arc<dyn Embedder>>,
loader_info: Option<LoaderInfo>,
initial_subscriber: Option<Arc<dyn lifecycle::Subscriber>>,
emit_migration_event: &mut impl FnMut(&MigrationStepReport),
) -> Result<OpenedEngine, EngineOpenError> {
Self::open_with_migrations(
path,
MIGRATIONS,
embedder_identity,
runtime_embedder,
loader_info,
emit_migration_event,
initial_subscriber,
)
}
fn open_with_migrations(
path: impl Into<PathBuf>,
migrations: &'static [fathomdb_schema::Migration],
embedder_identity: EmbedderIdentity,
runtime_embedder: Option<Arc<dyn Embedder>>,
loader_info: Option<LoaderInfo>,
emit_migration_event: &mut impl FnMut(&MigrationStepReport),
initial_subscriber: Option<Arc<dyn lifecycle::Subscriber>>,
) -> Result<OpenedEngine, EngineOpenError> {
let canonical_path = canonical_database_path(&path.into())?;
let lock = acquire_lock(&canonical_path)?;
let open_result = Self::open_locked(
canonical_path.clone(),
migrations,
&embedder_identity,
emit_migration_event,
);
match open_result {
Ok((connection, readers, mut report, reader_lookaside_rcs)) => {
if let Some(info) = loader_info {
if info.download_ms.is_some() {
report.embedder_download_ms = info.download_ms;
}
if !info.events.is_empty() {
report.embedder_events = info.events;
}
}
let next_cursor = load_next_cursor(&connection);
let subscribers = Arc::new(lifecycle::SubscriberRegistry::new());
let profiling_enabled = Arc::new(AtomicBool::new(false));
let slow_threshold_ms = Arc::new(AtomicU64::new(DEFAULT_SLOW_THRESHOLD_MS));
let mut profile_contexts: Vec<Box<ProfileContext>> = Vec::new();
let projection_runtime = ProjectionRuntime::new(
canonical_path.clone(),
runtime_embedder.clone(),
embedder_identity.clone(),
report.embedder_mean_vec_pinned,
);
install_profile_callback(
&connection,
&subscribers,
&profiling_enabled,
&slow_threshold_ms,
&mut profile_contexts,
);
for reader in &readers {
install_profile_callback(
reader,
&subscribers,
&profiling_enabled,
&slow_threshold_ms,
&mut profile_contexts,
);
}
let opened = OpenedEngine {
engine: Self {
path: canonical_path.clone(),
next_cursor: AtomicU64::new(next_cursor),
closed: AtomicBool::new(false),
lock: Mutex::new(Some(lock)),
connection: Mutex::new(Some(connection)),
reader_pool: ReaderWorkerPool::new(readers),
counters: lifecycle::Counters::new(),
subscribers,
profiling_enabled,
slow_threshold_ms,
runtime_embedder,
runtime_embedder_identity: embedder_identity,
projection_runtime,
provenance_row_cap: AtomicU64::new(DEFAULT_PROVENANCE_ROW_CAP),
profile_contexts: Mutex::new(profile_contexts),
reader_lookaside_rcs,
#[cfg(debug_assertions)]
force_next_commit_failure: AtomicBool::new(false),
},
report,
};
if let Some(subscriber) = initial_subscriber {
opened.engine.subscribers.attach_persistent(subscriber);
}
if database_has_pending_projection_work(&canonical_path).unwrap_or(false) {
opened.engine.projection_runtime.notify_new_work();
}
Ok(opened)
}
Err(err) => {
if let Some(subscriber) = initial_subscriber {
emit_open_error_event(&subscriber, &err);
}
drop(lock);
Err(err)
}
}
}
fn open_locked(
path: PathBuf,
migrations: &'static [fathomdb_schema::Migration],
embedder_identity: &EmbedderIdentity,
emit_migration_event: &mut impl FnMut(&MigrationStepReport),
) -> Result<(Connection, Vec<Connection>, OpenReport, Vec<i32>), EngineOpenError> {
init_perf_experiments_runtime();
register_sqlite_vec_extension();
let mut connection = Connection::open(&path)
.map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))?;
probe_database_header(&connection)?;
probe_open_integrity(&connection)?;
probe_wal_sidecar(&path)?;
apply_perf_experiment_writer_pragmas(&connection);
connection
.pragma_update(None, "journal_mode", "WAL")
.map_err(|err| map_open_sqlite_error(err, OpenStage::WalReplay))?;
reject_legacy_shape(&connection)?;
let migration = migrate_with_event_sink(&connection, migrations, emit_migration_event)
.map_err(map_migration_error)?;
if migration.schema_version_after >= SEARCH_INDEX_TOKENIZER_SCHEMA_VERSION
&& !search_index_tokenizer_reproject_complete(&connection).map_err(|_| {
EngineOpenError::Io {
message: "could not read search_index tokenizer reproject marker".to_string(),
}
})?
{
reproject_search_index_after_tokenizer_upgrade(&connection).map_err(|_| {
EngineOpenError::Io {
message: "could not re-tokenize search_index after tokenizer upgrade"
.to_string(),
}
})?;
}
let mut embedder_mean_vec_pinned = check_embedder_profile(&connection, embedder_identity)?;
ensure_vector_partition(&mut connection, embedder_identity.dimension).map_err(|_| {
EngineOpenError::Io { message: "could not initialize vector partition".to_string() }
})?;
if identity_requires_mean_centering(embedder_identity) && !embedder_mean_vec_pinned {
let row_count: u64 = connection
.query_row("SELECT COUNT(*) FROM vector_default", [], |row| row.get(0))
.unwrap_or(0);
if row_count >= MEAN_VEC_PIN_THRESHOLD {
recover_mean_vec_pin(&mut connection, embedder_identity).map_err(|_| {
EngineOpenError::Io {
message: "could not recover mean-centering pin".to_string(),
}
})?;
embedder_mean_vec_pinned = true;
}
}
let warmup_started = Instant::now();
let embedder_mean_centering_required = embedder_identity.name == BGE_SMALL_EMBEDDER_NAME;
let report = OpenReport {
schema_version_before: migration.schema_version_before,
schema_version_after: migration.schema_version_after,
migration_steps: migration.migration_steps,
embedder_warmup_ms: u64::try_from(warmup_started.elapsed().as_millis())
.unwrap_or(u64::MAX),
query_backend: "fathomdb-query + sqlite-vec",
default_embedder: embedder_identity.clone(),
embedder_download_ms: None,
embedder_events: Vec::new(),
embedder_mean_centering_required,
embedder_mean_vec_pinned,
};
let mut readers = Vec::with_capacity(READER_POOL_SIZE);
let mut lookaside_rcs: Vec<i32> = Vec::with_capacity(READER_POOL_SIZE);
for _ in 0..READER_POOL_SIZE {
let reader = Connection::open(&path)
.map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))?;
let rc: i32 = configure_reader_lookaside(&reader);
debug_assert_eq!(
rc,
rusqlite::ffi::SQLITE_OK,
"sqlite3_db_config(LOOKASIDE) must return SQLITE_OK on a freshly opened reader",
);
lookaside_rcs.push(rc);
reader
.pragma_update(None, "journal_mode", "WAL")
.map_err(|err| map_open_sqlite_error(err, OpenStage::WalReplay))?;
reader
.pragma_update(None, "query_only", "ON")
.map_err(|err| map_open_sqlite_error(err, OpenStage::SchemaProbe))?;
apply_perf_experiment_reader_pragmas(&reader);
readers.push(reader);
}
Ok((connection, readers, report, lookaside_rcs))
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
pub fn write(&self, batch: &[PreparedWrite]) -> Result<WriteReceipt, EngineError> {
let category = if batch_is_admin(batch) {
lifecycle::EventCategory::Admin
} else {
lifecycle::EventCategory::Writer
};
self.emit_event(lifecycle::Phase::Started, category, None);
let started = Instant::now();
let outcome = self.write_inner(batch);
self.detect_slow(started, category);
match outcome {
Ok(receipt) => {
let rows = u64::try_from(batch.len()).unwrap_or(u64::MAX);
if batch_is_admin(batch) {
self.counters.record_admin();
} else {
self.counters.record_write(rows);
}
self.emit_event(lifecycle::Phase::Finished, category, None);
Ok(receipt)
}
Err(err) => {
let code = err.stable_code();
self.counters.record_error(code);
self.emit_event(lifecycle::Phase::Failed, category, Some(code));
self.emit_event(
lifecycle::Phase::Failed,
lifecycle::EventCategory::Error,
Some(code),
);
Err(err)
}
}
}
fn write_inner(&self, batch: &[PreparedWrite]) -> Result<WriteReceipt, EngineError> {
self.ensure_open()?;
if batch.is_empty() {
return Err(EngineError::WriteValidation);
}
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
let plans = validate_batch(connection, batch)?;
let projection_jobs = collect_projection_jobs(connection, batch)?;
#[cfg(debug_assertions)]
if self.force_next_commit_failure.swap(false, Ordering::SeqCst) {
return Err(EngineError::Storage);
}
let base_cursor = self.next_cursor.load(Ordering::SeqCst);
let increment = u64::try_from(batch.len()).unwrap_or(u64::MAX);
let last_cursor = base_cursor.saturating_add(increment);
let pending_projection = !projection_jobs.is_empty();
let dangling_edge_endpoints = match commit_batch(
connection,
batch,
&plans,
base_cursor,
self.provenance_row_cap.load(Ordering::Relaxed),
) {
Ok(count) => count,
Err(err) => {
self.emit_sqlite_internal_error(&err);
return Err(EngineError::Storage);
}
};
self.next_cursor.store(last_cursor, Ordering::SeqCst);
if pending_projection {
self.projection_runtime.notify_new_work();
}
let row_cursors = (0..batch.len())
.map(|i| base_cursor.saturating_add((i as u64).saturating_add(1)))
.collect();
Ok(WriteReceipt { cursor: last_cursor, row_cursors, dangling_edge_endpoints })
}
pub fn search(&self, query: &str) -> Result<SearchResult, EngineError> {
self.search_filtered(query, None)
}
pub fn search_filtered(
&self,
query: &str,
filter: Option<SearchFilter>,
) -> Result<SearchResult, EngineError> {
self.emit_event(lifecycle::Phase::Started, lifecycle::EventCategory::Search, None);
let started = Instant::now();
let outcome = self.search_inner(query, filter);
self.detect_slow(started, lifecycle::EventCategory::Search);
match outcome {
Ok(result) => {
self.counters.record_query();
self.emit_event(lifecycle::Phase::Finished, lifecycle::EventCategory::Search, None);
Ok(result)
}
Err(err) => {
let code = err.stable_code();
self.counters.record_error(code);
self.emit_event(
lifecycle::Phase::Failed,
lifecycle::EventCategory::Search,
Some(code),
);
self.emit_event(
lifecycle::Phase::Failed,
lifecycle::EventCategory::Error,
Some(code),
);
Err(err)
}
}
}
fn detect_slow(&self, started: Instant, category: lifecycle::EventCategory) {
let elapsed = started.elapsed();
let threshold = self.slow_threshold_ms.load(Ordering::Relaxed);
let threshold_duration = std::time::Duration::from_millis(threshold);
if elapsed > threshold_duration {
self.emit_event(lifecycle::Phase::Slow, category, None);
}
}
fn emit_event(
&self,
phase: lifecycle::Phase,
category: lifecycle::EventCategory,
code: Option<&'static str>,
) {
let event =
lifecycle::Event { phase, source: lifecycle::EventSource::Engine, category, code };
self.subscribers.dispatch(&event);
}
fn emit_sqlite_internal_error(&self, err: &rusqlite::Error) {
if let Some(code) = sqlite_extended_code_name(err) {
let event = lifecycle::Event {
phase: lifecycle::Phase::Failed,
source: lifecycle::EventSource::SqliteInternal,
category: lifecycle::EventCategory::Error,
code: Some(code),
};
self.subscribers.dispatch(&event);
}
}
fn search_inner(
&self,
query: &str,
filter: Option<SearchFilter>,
) -> Result<SearchResult, EngineError> {
self.ensure_open()?;
if query.trim().is_empty() {
return Err(EngineError::WriteValidation);
}
let compiled = compile_text_query(query);
let raw_query_vector =
self.runtime_embedder.as_ref().and_then(|embedder| embedder.embed(query).ok());
let query_vector_bin = match raw_query_vector.as_ref() {
Some(vector) if identity_requires_mean_centering(&self.runtime_embedder_identity) => {
let pinned = {
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
read_pinned_mean_vec(connection, self.runtime_embedder_identity.dimension)?
};
match pinned {
Some(mean) => serde_json::to_string(&subtract_mean(vector, &mean)).ok(),
None => serde_json::to_string(vector).ok(),
}
}
Some(vector) => serde_json::to_string(vector).ok(),
None => None,
};
let query_vector = raw_query_vector.and_then(|vector| serde_json::to_string(&vector).ok());
let search_limit = self
.projection_runtime
.shared
.search_limit_override
.load(Ordering::SeqCst)
.max(SEARCH_RERANK_LIMIT);
let recency_enabled =
self.projection_runtime.shared.recency_reweight_enabled.load(Ordering::SeqCst);
let vector_stage_only =
self.projection_runtime.shared.vector_stage_only_for_test.load(Ordering::SeqCst);
let (response_tx, response_rx) = mpsc::sync_channel::<ReaderResponse>(1);
let request = ReaderRequest::Search {
compiled,
query_vector,
query_vector_bin,
search_limit,
filter: filter.map(Box::new),
recency_enabled,
vector_stage_only,
respond: response_tx,
};
if self.reader_pool.dispatch(request).is_err() {
return Err(EngineError::Closing);
}
let search_result = response_rx.recv().map_err(|_| EngineError::Storage)?;
let (cursor, soft_fallback, results) = match search_result {
Ok(result) => result,
Err(err) => {
self.emit_sqlite_internal_error(&err);
return Err(EngineError::Storage);
}
};
Ok(SearchResult { projection_cursor: cursor, soft_fallback, results })
}
pub fn read_get(&self, logical_id: &str) -> Result<Option<NodeRecord>, EngineError> {
let ids = [logical_id.to_string()];
let rows = self.read_get_many(&ids)?;
Ok(rows.into_iter().next().flatten())
}
pub fn read_get_many(
&self,
logical_ids: &[String],
) -> Result<Vec<Option<NodeRecord>>, EngineError> {
self.ensure_open()?;
if logical_ids.is_empty() {
return Ok(Vec::new());
}
let (response_tx, response_rx) = mpsc::sync_channel(1);
let request =
ReaderRequest::GetById { logical_ids: logical_ids.to_vec(), respond: response_tx };
if self.reader_pool.dispatch(request).is_err() {
return Err(EngineError::Closing);
}
match response_rx.recv().map_err(|_| EngineError::Storage)? {
Ok(rows) => Ok(rows),
Err(err) => {
self.emit_sqlite_internal_error(&err);
Err(EngineError::Storage)
}
}
}
pub fn read_collection(
&self,
collection: &str,
after_id: Option<i64>,
limit: usize,
) -> Result<Vec<OpStoreRow>, EngineError> {
self.read_collection_dispatch(collection, after_id, limit)
}
pub fn read_mutations(
&self,
collection: &str,
after_id: Option<i64>,
limit: usize,
) -> Result<Vec<OpStoreRow>, EngineError> {
self.read_collection_dispatch(collection, after_id, limit)
}
fn read_collection_dispatch(
&self,
collection: &str,
after_id: Option<i64>,
limit: usize,
) -> Result<Vec<OpStoreRow>, EngineError> {
self.ensure_open()?;
let (response_tx, response_rx) = mpsc::sync_channel(1);
let request = ReaderRequest::ReadCollection {
collection: collection.to_string(),
after_id,
limit,
respond: response_tx,
};
if self.reader_pool.dispatch(request).is_err() {
return Err(EngineError::Closing);
}
match response_rx.recv().map_err(|_| EngineError::Storage)? {
Ok(rows) => Ok(rows),
Err(err) => {
self.emit_sqlite_internal_error(&err);
Err(EngineError::Storage)
}
}
}
pub fn close(&self) -> Result<(), EngineError> {
self.closed.store(true, Ordering::SeqCst);
self.projection_runtime.stop();
self.reader_pool.shutdown();
if let Ok(mut connection) = self.connection.lock() {
if let Some(conn) = connection.as_ref() {
uninstall_profile_callback(conn);
}
connection.take();
}
if let Ok(mut contexts) = self.profile_contexts.lock() {
contexts.clear();
}
if let Ok(mut lock) = self.lock.lock() {
lock.take();
}
Ok(())
}
pub fn drain(&self, timeout_ms: u64) -> Result<(), EngineError> {
self.ensure_open()?;
if self.projection_runtime.wait_for_idle(timeout_ms) {
Ok(())
} else {
Err(EngineError::Scheduler)
}
}
#[must_use]
pub fn counters(&self) -> CounterSnapshot {
self.counters.snapshot()
}
pub fn set_profiling(&self, enabled: bool) -> Result<(), EngineError> {
self.profiling_enabled.store(enabled, Ordering::Relaxed);
Ok(())
}
pub fn set_slow_threshold_ms(&self, value: u64) -> Result<(), EngineError> {
self.slow_threshold_ms.store(value, Ordering::Relaxed);
Ok(())
}
#[must_use]
pub fn subscribe(&self, subscriber: Arc<dyn lifecycle::Subscriber>) -> Subscription {
self.subscribers.attach(subscriber)
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn reader_worker_count_for_test(&self) -> usize {
self.reader_pool.worker_count()
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn live_reader_worker_count_for_test(&self) -> usize {
self.reader_pool.live_count()
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn reader_lookaside_config_rcs_for_test(&self) -> Vec<i32> {
self.reader_lookaside_rcs.clone()
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn reader_lookaside_used_per_worker_for_test(&self) -> Vec<i32> {
self.reader_pool.lookaside_used_per_worker()
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn cache_status_per_worker_for_test(&self, label: &str) -> Vec<CacheStatusReply> {
self.reader_pool.cache_status_per_worker(label)
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn force_next_commit_failure_for_test(&self) {
self.force_next_commit_failure.store(true, Ordering::SeqCst);
}
#[cfg(debug_assertions)]
#[doc(hidden)]
pub fn execute_for_test(&self, sql: &str) -> Result<(), EngineError> {
self.ensure_open()?;
let started = Instant::now();
{
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
connection.execute_batch(sql).map_err(|_| EngineError::Storage)?;
}
self.detect_slow(started, lifecycle::EventCategory::Search);
Ok(())
}
#[doc(hidden)]
#[cfg(debug_assertions)]
pub fn run_one_thread_poison_for_test(&self) -> Result<(), EngineError> {
self.ensure_open()?;
self.write(&[PreparedWrite::Node {
kind: "doc".to_string(),
body: "poison-fixture-seed".to_string(),
source_id: None,
logical_id: None,
}])?;
let poison_outcome: Mutex<Option<EngineError>> = Mutex::new(None);
let poison_thread_id: AtomicU64 = AtomicU64::new(0);
thread::scope(|scope| {
for _ in 0..4 {
scope.spawn(|| {
for _ in 0..4 {
let _ = self.search("poison-fixture-seed");
}
});
}
scope.spawn(|| {
let _ = self.write(&[PreparedWrite::Node {
kind: "doc".to_string(),
body: "writer-progress".to_string(),
source_id: None,
logical_id: None,
}]);
});
scope.spawn(|| {
poison_thread_id.store(1, Ordering::SeqCst);
if let Err(err) = self.write(&[]) {
*poison_outcome.lock().expect("poison_outcome lock") = Some(err);
}
});
});
let err = poison_outcome
.into_inner()
.expect("poison_outcome lock")
.expect("poison thread must produce a deterministic error");
let projection_state = match self.projection_status_for_test("doc") {
Ok(lifecycle::ProjectionStatus::Pending) => "Pending",
Ok(lifecycle::ProjectionStatus::Failed) => "Failed",
Ok(lifecycle::ProjectionStatus::UpToDate) => "UpToDate",
Err(_) => "UpToDate",
};
let context = lifecycle::StressFailureContext {
thread_group_id: poison_thread_id.load(Ordering::SeqCst),
op_kind: "write".to_string(),
last_error_chain: vec![err.stable_code().to_string(), err.to_string()],
projection_state: projection_state.to_string(),
};
self.subscribers.dispatch_stress_failure(&context);
Ok(())
}
#[doc(hidden)]
pub fn set_projection_scheduler_frozen_for_test(&self, frozen: bool) {
self.projection_runtime.set_frozen(frozen);
}
#[doc(hidden)]
pub fn set_projection_retry_delays_for_test(&self, delays_ms: &[u64]) {
self.projection_runtime.set_retry_delays_for_test(delays_ms);
}
#[doc(hidden)]
pub fn set_embed_timeout_ms_for_test(&self, timeout_ms: u64) {
self.projection_runtime.set_embed_timeout_ms_for_test(timeout_ms);
}
#[doc(hidden)]
pub fn set_embed_circuit_threshold_for_test(&self, threshold: u64) {
self.projection_runtime.set_embed_circuit_threshold_for_test(threshold);
}
#[doc(hidden)]
pub fn embed_circuit_open_for_test(&self) -> bool {
self.projection_runtime.embed_circuit_open_for_test()
}
#[doc(hidden)]
pub fn projection_status_for_test(
&self,
kind: &str,
) -> Result<lifecycle::ProjectionStatus, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
projection_status(connection, kind)
}
#[doc(hidden)]
pub fn has_vector_for_cursor_for_test(&self, cursor: u64) -> Result<bool, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
terminal_state_for_cursor(connection, cursor)
.map(|state| matches!(state.as_deref(), Some("up_to_date")))
.map_err(|_| EngineError::Storage)
}
#[doc(hidden)]
pub fn projection_failure_count_for_test(&self, cursor: u64) -> Result<u64, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
connection
.query_row(
"SELECT COUNT(*) FROM operational_mutations
WHERE collection_name = 'projection_failures'
AND record_key = ?1",
[cursor.to_string()],
|row| row.get::<_, u64>(0),
)
.map_err(|_| EngineError::Storage)
}
#[doc(hidden)]
pub fn set_provenance_row_cap_for_test(&self, cap: Option<u64>) {
self.provenance_row_cap.store(cap.unwrap_or(0), Ordering::Relaxed);
}
#[doc(hidden)]
pub fn provenance_row_count_for_test(&self) -> Result<u64, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
connection
.query_row("SELECT COUNT(*) FROM operational_mutations", [], |row| row.get::<_, u64>(0))
.map_err(|_| EngineError::Storage)
}
#[doc(hidden)]
pub fn oldest_provenance_record_key_for_test(
&self,
collection: &str,
) -> Result<Option<String>, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
connection
.query_row(
"SELECT record_key FROM operational_mutations
WHERE collection_name = ?1
ORDER BY id
LIMIT 1",
[collection],
|row| row.get::<_, String>(0),
)
.map(Some)
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(None),
_ => Err(EngineError::Storage),
})
}
#[doc(hidden)]
pub fn configure_vector_kind_for_test(&self, kind: &str) -> Result<(), EngineError> {
self.ensure_open()?;
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
connection
.execute(
"INSERT OR REPLACE INTO _fathomdb_vector_kinds(kind, profile, created_at)
VALUES(?1, ?2, 0)",
params![kind, DEFAULT_VECTOR_PROFILE],
)
.map_err(|_| EngineError::Storage)?;
Ok(())
}
#[doc(hidden)]
pub fn write_vector_for_test(
&self,
kind: &str,
text: &str,
) -> Result<WriteReceipt, EngineError> {
self.ensure_open()?;
let embedder =
self.runtime_embedder.as_ref().cloned().ok_or(EngineError::EmbedderNotConfigured)?;
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
if !kind_is_vector_indexed(connection, kind)? {
return Err(EngineError::KindNotVectorIndexed);
}
let expected = default_profile_dimension(connection)?;
ensure_vector_partition(connection, expected).map_err(|_| EngineError::Storage)?;
let vector = embedder.embed(text).map_err(map_runtime_embedder_error)?;
let actual = u32::try_from(vector.len()).unwrap_or(u32::MAX);
if actual != expected {
return Err(EngineError::EmbedderDimensionMismatch { expected, actual });
}
let cursor = self.next_cursor.load(Ordering::SeqCst).saturating_add(1);
let blob = encode_vector_blob(&vector);
let bin_blob = if identity_requires_mean_centering(&self.runtime_embedder_identity) {
match read_pinned_mean_vec(connection, self.runtime_embedder_identity.dimension)? {
Some(mean) => encode_vector_blob(&subtract_mean(&vector, &mean)),
None => blob.clone(),
}
} else {
blob.clone()
};
let source_type = resolve_source_type(kind)?;
let now_unix =
SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs() as i64;
let pin_event = {
let runtime = &self.projection_runtime.shared;
let mut accumulator =
runtime.mean_accumulator.lock().map_err(|_| EngineError::Storage)?;
if let Some(acc) = accumulator.as_mut() {
acc.add(&vector);
if acc.count() >= MEAN_VEC_PIN_THRESHOLD {
let mean = acc.materialize();
*accumulator = None;
Some(mean)
} else {
None
}
} else {
None
}
};
let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
tx.execute(
"INSERT INTO _fathomdb_vector_rows(rowid, kind, write_cursor) VALUES(?1, ?2, ?3)",
params![cursor, kind, cursor],
)
.map_err(|_| EngineError::Storage)?;
tx.execute(
"INSERT INTO vector_default(
rowid, embedding, embedding_bin, source_type, kind, created_at, status
) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
params![cursor, blob, bin_blob, source_type, kind, now_unix],
)
.map_err(|_| EngineError::Storage)?;
let mut emitted_event: Option<EmbedderEvent> = None;
if let Some(mean_vec) = pin_event {
let mean_bytes = encode_vector_blob(&mean_vec);
tx.execute(
"UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
params![mean_bytes],
)
.map_err(|_| EngineError::Storage)?;
let rows: Vec<(i64, Vec<u8>)> = {
let mut statement = tx
.prepare("SELECT rowid, embedding FROM vector_default ORDER BY rowid")
.map_err(|_| EngineError::Storage)?;
let mapped = statement
.query_map([], |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?)))
.map_err(|_| EngineError::Storage)?;
let mut out = Vec::new();
for r in mapped {
out.push(r.map_err(|_| EngineError::Storage)?);
}
out
};
let (doc_count, _) = run_pin_and_requantize_pass(&tx, &rows, &mean_vec)?;
emitted_event = Some(EmbedderEvent::MeanVecPinned {
dim: u32::try_from(mean_vec.len()).unwrap_or(u32::MAX),
doc_count,
});
}
tx.commit().map_err(|_| EngineError::Storage)?;
if let Some(ev) = emitted_event {
if let Ok(mut events) = self.projection_runtime.shared.pending_events.lock() {
events.push(ev);
}
}
self.next_cursor.store(cursor, Ordering::SeqCst);
Ok(WriteReceipt { cursor, row_cursors: vec![cursor], dangling_edge_endpoints: 0 })
}
#[doc(hidden)]
pub fn drain_mean_centering_events_for_test(&self) -> Result<Vec<EmbedderEvent>, EngineError> {
self.ensure_open()?;
let mut events = self
.projection_runtime
.shared
.pending_events
.lock()
.map_err(|_| EngineError::Storage)?;
let out = std::mem::take(&mut *events);
Ok(out)
}
pub fn drain_embedder_events(&self) -> Result<Vec<EmbedderEvent>, EngineError> {
self.ensure_open()?;
let mut events = self
.projection_runtime
.shared
.pending_events
.lock()
.map_err(|_| EngineError::Storage)?;
Ok(std::mem::take(&mut *events))
}
#[cfg(feature = "operator")]
pub fn recompute_mean(&self) -> Result<MeanRecomputeReport, EngineError> {
self.ensure_open()?;
let identity = self.runtime_embedder_identity.clone();
if !identity_requires_mean_centering(&identity) {
return Err(EngineError::EmbedderNotConfigured);
}
let report = {
let _gate = self
.projection_runtime
.shared
.commit_gate
.lock()
.unwrap_or_else(|p| p.into_inner());
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
#[cfg(debug_assertions)]
let fail = self
.projection_runtime
.shared
.force_recompute_failure
.swap(false, Ordering::SeqCst);
#[cfg(not(debug_assertions))]
let fail = false;
let report = recompute_mean_in_tx_inner(&tx, &identity, fail)?;
tx.commit().map_err(|_| EngineError::Storage)?;
report
};
if let Ok(mut events) = self.projection_runtime.shared.pending_events.lock() {
events.push(EmbedderEvent::MeanVecRecomputed {
dim: report.dim,
doc_count: report.doc_count_requantized,
trigger: MeanRecomputeTrigger::Manual,
});
}
Ok(report)
}
#[doc(hidden)]
pub fn set_search_limit_for_test(&self, limit: usize) {
self.projection_runtime.shared.search_limit_override.store(limit, Ordering::SeqCst);
}
#[doc(hidden)]
pub fn set_recency_reweight_enabled_for_test(&self, enabled: bool) {
self.projection_runtime.shared.recency_reweight_enabled.store(enabled, Ordering::SeqCst);
}
#[doc(hidden)]
pub fn set_vector_stage_only_for_test(&self, enabled: bool) {
self.projection_runtime.shared.vector_stage_only_for_test.store(enabled, Ordering::SeqCst);
}
#[doc(hidden)]
#[cfg(debug_assertions)]
pub fn force_next_recompute_failure_for_test(&self) {
self.projection_runtime.shared.force_recompute_failure.store(true, Ordering::SeqCst);
}
#[doc(hidden)]
pub fn vector_row_count_for_test(&self) -> Result<u64, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
connection
.query_row("SELECT COUNT(*) FROM vector_default", [], |row| row.get::<_, u64>(0))
.map_err(|_| EngineError::Storage)
}
#[doc(hidden)]
pub fn read_vector_blob_for_test(&self, rowid: i64) -> Result<Vec<u8>, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
connection
.query_row("SELECT embedding FROM vector_default WHERE rowid = ?1", [rowid], |row| {
row.get::<_, Vec<u8>>(0)
})
.map_err(|_| EngineError::Storage)
}
#[doc(hidden)]
pub fn default_embedder_profile_for_test(&self) -> Result<EmbedderIdentity, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
load_default_profile(connection).map_err(|_| EngineError::Storage)
}
#[cfg(feature = "operator")]
pub fn check_integrity(
&self,
opts: CheckIntegrityOpts,
) -> Result<IntegrityReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
Ok(IntegrityReport {
physical: physical_section(connection, opts.full),
logical: logical_section(connection),
semantic: semantic_section(connection),
})
}
#[cfg(feature = "operator")]
pub fn safe_export(
&self,
out: &Path,
manifest: &Path,
) -> Result<SafeExportArtifact, EngineError> {
self.ensure_open()?;
{
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let target = out.to_string_lossy().to_string();
connection
.execute("VACUUM INTO ?1", params![target])
.map_err(|_| EngineError::Storage)?;
}
let bytes = std::fs::read(out).map_err(|_| EngineError::Storage)?;
let digest = sha2::Sha256::digest(&bytes);
let sha256_hex = hex_encode(digest.as_slice());
let export_abs = out.canonicalize().unwrap_or_else(|_| out.to_path_buf());
let manifest_json = serde_json::json!({
"export_path": export_abs.to_string_lossy(),
"sha256": sha256_hex,
"byte_count": bytes.len() as u64,
});
let manifest_bytes =
serde_json::to_vec_pretty(&manifest_json).map_err(|_| EngineError::Storage)?;
std::fs::write(manifest, &manifest_bytes).map_err(|_| EngineError::Storage)?;
Ok(SafeExportArtifact {
export_path: out.to_path_buf(),
manifest_path: manifest.to_path_buf(),
manifest_sha256: sha256_hex,
})
}
#[cfg(feature = "operator")]
pub fn rebuild_projections(&self) -> Result<RebuildReport, EngineError> {
self.ensure_open()?;
self.run_rebuild(true, RebuildKind::Projections)
}
#[cfg(feature = "operator")]
pub fn rebuild_vec0(&self) -> Result<RebuildReport, EngineError> {
self.ensure_open()?;
self.run_rebuild(false, RebuildKind::Vec0)
}
#[cfg(feature = "operator")]
pub fn trace_source_ref(&self, source_id: &str) -> Result<TraceReport, EngineError> {
self.ensure_open()?;
if source_id.is_empty() {
return Err(EngineError::WriteValidation);
}
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let mut events: Vec<TraceEvent> = Vec::new();
let mut nodes = connection
.prepare(
"SELECT write_cursor, kind FROM canonical_nodes WHERE source_id = ?1
ORDER BY write_cursor",
)
.map_err(|_| EngineError::Storage)?;
let node_rows = nodes
.query_map([source_id], |row| {
Ok(TraceEvent {
write_cursor: row.get::<_, i64>(0)? as u64,
kind: row.get::<_, String>(1)?,
table: "canonical_nodes",
})
})
.map_err(|_| EngineError::Storage)?;
for row in node_rows {
events.push(row.map_err(|_| EngineError::Storage)?);
}
let mut edges = connection
.prepare(
"SELECT write_cursor, kind FROM canonical_edges WHERE source_id = ?1
ORDER BY write_cursor",
)
.map_err(|_| EngineError::Storage)?;
let edge_rows = edges
.query_map([source_id], |row| {
Ok(TraceEvent {
write_cursor: row.get::<_, i64>(0)? as u64,
kind: row.get::<_, String>(1)?,
table: "canonical_edges",
})
})
.map_err(|_| EngineError::Storage)?;
for row in edge_rows {
events.push(row.map_err(|_| EngineError::Storage)?);
}
events.sort_by_key(|e| e.write_cursor);
Ok(TraceReport { source_ref: source_id.to_string(), events })
}
#[cfg(feature = "operator")]
pub fn excise_source(&self, source_id: &str) -> Result<ExciseReport, EngineError> {
self.ensure_open()?;
if source_id.is_empty() {
return Err(EngineError::WriteValidation);
}
self.projection_runtime.set_frozen(true);
let drain_result = self.drain(REBUILD_DRAIN_TIMEOUT_MS);
let outcome = drain_result.and_then(|()| self.excise_source_inner(source_id));
self.projection_runtime.set_frozen(false);
outcome
}
#[cfg(feature = "operator")]
pub fn verify_embedder(
&self,
supplied_identity: &str,
supplied_dimension: u32,
) -> Result<VerifyEmbedderReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let stored = load_default_profile(connection).map_err(|_| EngineError::Storage)?;
let stored_identity = format!("{}:{}", stored.name, stored.revision);
let identity_match = stored_identity == supplied_identity;
let dimension_match = stored.dimension == supplied_dimension;
let status = match (identity_match, dimension_match) {
(true, true) => VerifyEmbedderStatus::Match,
(false, true) => VerifyEmbedderStatus::IdentityMismatch,
(true, false) => VerifyEmbedderStatus::DimensionMismatch,
(false, false) => VerifyEmbedderStatus::BothMismatch,
};
Ok(VerifyEmbedderReport {
stored_identity,
stored_dimension: stored.dimension,
supplied_identity: supplied_identity.to_string(),
supplied_dimension,
status,
})
}
#[cfg(feature = "operator")]
pub fn dump_schema(&self) -> Result<DumpSchemaReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let user_version: u32 = connection
.query_row("PRAGMA user_version", [], |row| row.get(0))
.map_err(|_| EngineError::Storage)?;
let tables = read_schema_objects(connection, "table")?;
let indexes = read_schema_objects(connection, "index")?;
Ok(DumpSchemaReport { user_version, tables: order_canonical_first(tables), indexes })
}
#[cfg(feature = "operator")]
pub fn dump_row_counts(&self) -> Result<DumpRowCountsReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let mut counts = Vec::with_capacity(CANONICAL_TABLES.len());
for name in CANONICAL_TABLES {
let rows: u64 = connection
.query_row(&format!("SELECT COUNT(*) FROM {name}"), [], |row| row.get(0))
.map_err(|_| EngineError::Storage)?;
counts.push(TableRowCount { name: (*name).to_string(), rows });
}
Ok(DumpRowCountsReport { counts })
}
#[cfg(feature = "operator")]
pub fn dump_profile(&self) -> Result<DumpProfileReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let stored = load_default_profile(connection).map_err(|_| EngineError::Storage)?;
let mut stmt = connection
.prepare("SELECT kind FROM _fathomdb_vector_kinds ORDER BY kind")
.map_err(|_| EngineError::Storage)?;
let rows =
stmt.query_map([], |row| row.get::<_, String>(0)).map_err(|_| EngineError::Storage)?;
let mut vectorized_kinds = Vec::new();
for row in rows {
vectorized_kinds.push(row.map_err(|_| EngineError::Storage)?);
}
Ok(DumpProfileReport {
embedder_identity: format!("{}:{}", stored.name, stored.revision),
embedder_dimension: stored.dimension,
vectorized_kinds,
})
}
#[cfg(feature = "operator")]
pub fn truncate_wal(&self) -> Result<TruncateWalReport, EngineError> {
self.ensure_open()?;
let connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_ref().ok_or(EngineError::Closing)?;
let (busy, log_frames, checkpointed_frames): (i64, i64, i64) = connection
.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
Ok((row.get(0)?, row.get(1)?, row.get(2)?))
})
.map_err(|_| EngineError::Storage)?;
let status = if busy == 0 { TruncateWalStatus::Done } else { TruncateWalStatus::Busy };
Ok(TruncateWalReport {
status,
busy: busy.max(0) as u32,
log_frames: log_frames.max(0) as u32,
checkpointed_frames: checkpointed_frames.max(0) as u32,
})
}
#[cfg(feature = "operator")]
fn excise_source_inner(&self, source_id: &str) -> Result<ExciseReport, EngineError> {
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
let node_cursors: Vec<i64> = {
let mut stmt = tx
.prepare("SELECT write_cursor FROM canonical_nodes WHERE source_id = ?1")
.map_err(|_| EngineError::Storage)?;
let rows = stmt
.query_map([source_id], |row| row.get::<_, i64>(0))
.map_err(|_| EngineError::Storage)?;
rows.collect::<rusqlite::Result<Vec<_>>>().map_err(|_| EngineError::Storage)?
};
let edge_cursors: Vec<i64> = {
let mut stmt = tx
.prepare("SELECT write_cursor FROM canonical_edges WHERE source_id = ?1")
.map_err(|_| EngineError::Storage)?;
let rows = stmt
.query_map([source_id], |row| row.get::<_, i64>(0))
.map_err(|_| EngineError::Storage)?;
rows.collect::<rusqlite::Result<Vec<_>>>().map_err(|_| EngineError::Storage)?
};
let mut shadow_invalidated: u64 = 0;
for cursor in node_cursors.iter().chain(edge_cursors.iter()) {
shadow_invalidated = shadow_invalidated.saturating_add(
tx.execute("DELETE FROM search_index WHERE write_cursor = ?1", [cursor])
.map_err(|_| EngineError::Storage)? as u64,
);
shadow_invalidated = shadow_invalidated.saturating_add(
tx.execute("DELETE FROM vector_default WHERE rowid = ?1", [cursor])
.map_err(|_| EngineError::Storage)? as u64,
);
shadow_invalidated = shadow_invalidated.saturating_add(
tx.execute("DELETE FROM _fathomdb_vector_rows WHERE write_cursor = ?1", [cursor])
.map_err(|_| EngineError::Storage)? as u64,
);
shadow_invalidated = shadow_invalidated.saturating_add(
tx.execute(
"DELETE FROM _fathomdb_projection_terminal WHERE write_cursor = ?1",
[cursor],
)
.map_err(|_| EngineError::Storage)? as u64,
);
}
let nodes_excised = tx
.execute("DELETE FROM canonical_nodes WHERE source_id = ?1", [source_id])
.map_err(|_| EngineError::Storage)? as u64;
let edges_excised = tx
.execute("DELETE FROM canonical_edges WHERE source_id = ?1", [source_id])
.map_err(|_| EngineError::Storage)? as u64;
let excised_at = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
let payload = serde_json::json!({
"source_id": source_id,
"excised_at": excised_at,
"nodes_excised": nodes_excised,
"edges_excised": edges_excised,
"projections_invalidated": shadow_invalidated,
})
.to_string();
let audit_cursor = self.next_cursor.load(Ordering::SeqCst).saturating_add(1);
tx.execute(
"INSERT INTO operational_mutations(
collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
) VALUES('excise_source_audit', ?1, 'append', ?2, NULL, ?3)",
params![source_id, payload, audit_cursor],
)
.map_err(|_| EngineError::Storage)?;
tx.commit().map_err(|_| EngineError::Storage)?;
self.next_cursor.store(audit_cursor, Ordering::SeqCst);
Ok(ExciseReport {
source_ref: source_id.to_string(),
nodes_excised,
edges_excised,
projections_invalidated: shadow_invalidated,
})
}
#[cfg(feature = "operator")]
fn run_rebuild(
&self,
include_fts: bool,
kind: RebuildKind,
) -> Result<RebuildReport, EngineError> {
self.projection_runtime.set_frozen(true);
let drain_result = self.drain(REBUILD_DRAIN_TIMEOUT_MS);
let result = drain_result.and_then(|()| self.rebuild_shadow_state(include_fts, kind));
self.projection_runtime.set_frozen(false);
result
}
#[cfg(feature = "operator")]
fn rebuild_shadow_state(
&self,
include_fts: bool,
kind: RebuildKind,
) -> Result<RebuildReport, EngineError> {
let mut connection = self.connection.lock().map_err(|_| EngineError::Storage)?;
let connection = connection.as_mut().ok_or(EngineError::Closing)?;
let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
let mut rows_invalidated: u64 = 0;
if include_fts {
let n = tx.execute("DELETE FROM search_index", []).map_err(|_| EngineError::Storage)?;
rows_invalidated = rows_invalidated.saturating_add(n as u64);
}
let n = tx.execute("DELETE FROM vector_default", []).map_err(|_| EngineError::Storage)?;
rows_invalidated = rows_invalidated.saturating_add(n as u64);
let n = tx
.execute("DELETE FROM _fathomdb_vector_rows", [])
.map_err(|_| EngineError::Storage)?;
rows_invalidated = rows_invalidated.saturating_add(n as u64);
let n = tx
.execute("DELETE FROM _fathomdb_projection_terminal", [])
.map_err(|_| EngineError::Storage)?;
rows_invalidated = rows_invalidated.saturating_add(n as u64);
store_projection_cursor(&tx, 0).map_err(|_| EngineError::Storage)?;
let mut rows_rebuilt: u64 = 0;
if include_fts {
for row in canonical_node_rows(&tx).map_err(|_| EngineError::Storage)? {
tx.execute(
"INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)",
params![row.body, row.kind, row.cursor],
)
.map_err(|_| EngineError::Storage)?;
rows_rebuilt = rows_rebuilt.saturating_add(1);
}
}
let projection_cursor_after =
load_projection_cursor(&tx).map_err(|_| EngineError::Storage)?;
tx.commit().map_err(|_| EngineError::Storage)?;
Ok(RebuildReport { kind, rows_invalidated, rows_rebuilt, projection_cursor_after })
}
fn ensure_open(&self) -> Result<(), EngineError> {
if self.closed.load(Ordering::SeqCst) {
return Err(EngineError::Closing);
}
Ok(())
}
}
fn batch_is_admin(batch: &[PreparedWrite]) -> bool {
!batch.is_empty() && batch.iter().all(|w| matches!(w, PreparedWrite::AdminSchema { .. }))
}
pub const TOP_K_BIT_CANDIDATES: usize = 192;
pub const MEAN_VEC_PIN_THRESHOLD: u64 = 256;
pub const SEARCH_RERANK_LIMIT: usize = 10;
#[derive(Clone, Debug)]
struct MeanAccumulator {
sum: Vec<f64>,
count: u64,
}
impl MeanAccumulator {
fn new(dim: usize) -> Self {
Self { sum: vec![0.0; dim], count: 0 }
}
fn add(&mut self, v: &[f32]) {
debug_assert_eq!(v.len(), self.sum.len(), "accumulator dim mismatch");
for (slot, value) in self.sum.iter_mut().zip(v.iter()) {
*slot += f64::from(*value);
}
self.count = self.count.saturating_add(1);
}
fn materialize(&self) -> Vec<f32> {
if self.count == 0 {
return vec![0.0; self.sum.len()];
}
let denom = self.count as f64;
self.sum.iter().map(|s| (s / denom) as f32).collect()
}
fn count(&self) -> u64 {
self.count
}
}
fn cosine_similarity(a: &[f32], b: &[f32]) -> f32 {
if a.len() != b.len() {
return 1.0;
}
let mut dot = 0.0f64;
let mut na = 0.0f64;
let mut nb = 0.0f64;
for (x, y) in a.iter().zip(b.iter()) {
dot += f64::from(*x) * f64::from(*y);
na += f64::from(*x) * f64::from(*x);
nb += f64::from(*y) * f64::from(*y);
}
if na == 0.0 || nb == 0.0 {
return 1.0;
}
(dot / (na.sqrt() * nb.sqrt())) as f32
}
fn run_pin_and_requantize_pass(
tx: &rusqlite::Transaction<'_>,
rows: &[(i64, Vec<u8>)],
mean: &[f32],
) -> Result<(u64, Vec<EmbedderEvent>), EngineError> {
let mut updated: u64 = 0;
let dim = mean.len();
for (rowid, blob) in rows {
if blob.len() != dim * 4 {
return Err(EngineError::Storage);
}
let un_centered = decode_vector_blob(blob);
let centered = subtract_mean(&un_centered, mean);
let centered_blob = encode_vector_blob(¢ered);
let (source_type, kind, created_at): (String, String, i64) = tx
.query_row(
"SELECT source_type, kind, created_at FROM vector_default WHERE rowid = ?1",
params![rowid],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
)
.map_err(|_| EngineError::Storage)?;
tx.execute("DELETE FROM vector_default WHERE rowid = ?1", params![rowid])
.map_err(|_| EngineError::Storage)?;
tx.execute(
"INSERT INTO vector_default(
rowid, embedding, embedding_bin, source_type, kind, created_at, status
) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
params![rowid, blob, centered_blob, source_type, kind, created_at],
)
.map_err(|_| EngineError::Storage)?;
updated = updated.saturating_add(1);
}
let events = vec![EmbedderEvent::MeanVecPinned {
dim: u32::try_from(dim).unwrap_or(u32::MAX),
doc_count: updated,
}];
Ok((updated, events))
}
fn run_requantize_pass(rows: &[(i64, Vec<u8>)], mean: &[f32]) -> (u64, Vec<EmbedderEvent>) {
let mut updated: u64 = 0;
let dim = mean.len();
for (_rowid, blob) in rows {
if blob.len() != dim * 4 {
continue;
}
updated = updated.saturating_add(1);
}
let events = vec![EmbedderEvent::MeanVecPinned {
dim: u32::try_from(dim).unwrap_or(u32::MAX),
doc_count: updated,
}];
(updated, events)
}
#[doc(hidden)]
pub mod mean_centering_internals_for_test {
use super::{EmbedderEvent, MeanAccumulator};
pub struct AccumulatorHandle(MeanAccumulator);
#[must_use]
pub fn new_mean_accumulator(dim: usize) -> AccumulatorHandle {
AccumulatorHandle(MeanAccumulator::new(dim))
}
pub fn accumulator_add(handle: &mut AccumulatorHandle, v: &[f32]) {
handle.0.add(v);
}
#[must_use]
pub fn accumulator_materialize(handle: &AccumulatorHandle) -> Vec<f32> {
handle.0.materialize()
}
#[must_use]
pub fn accumulator_count(handle: &AccumulatorHandle) -> u64 {
handle.0.count()
}
#[must_use]
pub fn run_requantize_pass(rows: &[(i64, Vec<u8>)], mean: &[f32]) -> (u64, Vec<EmbedderEvent>) {
super::run_requantize_pass(rows, mean)
}
}
pub const RRF_K: f64 = 60.0;
pub const RECENCY_WEIGHT: f64 = 0.5 / RRF_K;
#[doc(hidden)]
#[must_use]
pub fn fuse_rrf(vector_hits: Vec<SearchHit>, text_hits: Vec<SearchHit>) -> Vec<SearchHit> {
struct Entry {
hit: SearchHit,
score: f64,
in_vector: bool,
order: usize,
}
let mut entries: Vec<Entry> = Vec::new();
let mut accumulate = |hit: SearchHit, rank0: usize, in_vector: bool| {
let contrib = 1.0 / (RRF_K + (rank0 as f64 + 1.0));
if let Some(existing) = entries.iter_mut().find(|e| e.hit.body == hit.body) {
existing.score += contrib;
} else {
let order = entries.len();
entries.push(Entry { hit, score: contrib, in_vector, order });
}
};
for (rank0, hit) in vector_hits.into_iter().enumerate() {
accumulate(hit, rank0, true);
}
for (rank0, hit) in text_hits.into_iter().enumerate() {
accumulate(hit, rank0, false);
}
entries.sort_by(|a, b| {
b.score
.partial_cmp(&a.score)
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| b.in_vector.cmp(&a.in_vector))
.then_with(|| a.order.cmp(&b.order))
});
entries
.into_iter()
.map(|mut e| {
e.hit.score = e.score;
e.hit
})
.collect()
}
#[doc(hidden)]
#[must_use]
pub fn apply_recency_reweight(hits: Vec<SearchHit>, enabled: bool) -> Vec<SearchHit> {
if !enabled || hits.len() < 2 {
return hits;
}
let min_id = hits.iter().map(|h| h.id).min().unwrap_or(0);
let max_id = hits.iter().map(|h| h.id).max().unwrap_or(0);
if max_id == min_id {
return hits;
}
let span = (max_id - min_id) as f64;
let mut reweighted: Vec<SearchHit> = hits
.into_iter()
.map(|mut h| {
let norm = (h.id - min_id) as f64 / span;
h.score += RECENCY_WEIGHT * norm;
h
})
.collect();
reweighted.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal));
reweighted
}
#[doc(hidden)]
#[must_use]
pub fn rerank_fused(hits: Vec<SearchHit>) -> Vec<SearchHit> {
hits
}
fn vector_filter_clause(filter: Option<&SearchFilter>) -> String {
let Some(filter) = filter else {
return String::new();
};
if filter.is_unfiltered() {
return String::new();
}
let mut cols: Vec<(&str, &str)> = Vec::new();
if filter.source_type.is_some() {
cols.push(("source_type", "="));
}
if filter.kind.is_some() {
cols.push(("kind", "="));
}
if filter.created_after.is_some() {
cols.push(("created_at", ">="));
}
if filter.status.is_some() {
cols.push(("status", "="));
}
let mut clause = String::new();
for (i, (col, op)) in cols.iter().enumerate() {
clause.push_str(&format!(" AND {col}{op}?{}", i + 3));
}
clause
}
fn vector_filter_values(filter: Option<&SearchFilter>) -> Vec<rusqlite::types::Value> {
use rusqlite::types::Value;
let mut out = Vec::new();
let Some(filter) = filter else {
return out;
};
if filter.is_unfiltered() {
return out;
}
if let Some(s) = &filter.source_type {
out.push(Value::Text(s.clone()));
}
if let Some(s) = &filter.kind {
out.push(Value::Text(s.clone()));
}
if let Some(c) = filter.created_after {
out.push(Value::Integer(c));
}
if let Some(s) = &filter.status {
out.push(Value::Text(s.clone()));
}
out
}
fn build_vector_phase1_sql(filter: Option<&SearchFilter>, final_limit: usize) -> String {
let filter_clause = vector_filter_clause(filter);
format!(
"WITH candidates AS (
SELECT rowid
FROM vector_default
WHERE embedding_bin MATCH vec_quantize_binary(vec_f32(?1)){filter_clause}
ORDER BY distance
LIMIT {top_k}
)
SELECT c.rowid, vec_distance_l2(v.embedding, vec_f32(?2)) AS l2
FROM candidates c
JOIN vector_default v ON v.rowid = c.rowid
ORDER BY l2
LIMIT {final_limit}",
top_k = TOP_K_BIT_CANDIDATES,
)
}
#[doc(hidden)]
#[must_use]
pub fn vector_phase1_sql_for_test(filter: Option<&SearchFilter>) -> String {
build_vector_phase1_sql(filter, SEARCH_RERANK_LIMIT)
}
fn text_hit_passes_filter(
tx: &rusqlite::Transaction<'_>,
id: u64,
kind: &str,
filter: Option<&SearchFilter>,
) -> rusqlite::Result<bool> {
let Some(filter) = filter else {
return Ok(true);
};
if filter.is_unfiltered() {
return Ok(true);
}
if let Some(k) = &filter.kind {
if kind != k {
return Ok(false);
}
}
if let Some(st) = &filter.source_type {
match resolve_source_type(kind) {
Ok(resolved) if resolved == st.as_str() => {}
_ => return Ok(false),
}
}
if filter.created_after.is_some() || filter.status.is_some() {
let meta: Option<(i64, Option<String>)> = tx
.query_row(
"SELECT created_at, status FROM vector_default WHERE rowid = ?1 LIMIT 1",
[id as i64],
|row| Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?)),
)
.optional()?;
let Some((created_at, status)) = meta else {
return Ok(false);
};
if let Some(bound) = filter.created_after {
if created_at < bound {
return Ok(false);
}
}
if let Some(want) = &filter.status {
if status.as_deref() != Some(want.as_str()) {
return Ok(false);
}
}
}
Ok(true)
}
#[allow(clippy::too_many_arguments)]
fn read_search_in_tx(
reader: &mut Connection,
compiled: &fathomdb_query::CompiledQuery,
query_vector: Option<&str>,
query_vector_bin: Option<&str>,
final_limit: usize,
filter: Option<&SearchFilter>,
recency_enabled: bool,
vector_stage_only: bool,
) -> rusqlite::Result<(u64, Option<SoftFallback>, Vec<SearchHit>)> {
let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
let cursor = load_projection_cursor(&tx)?;
let vector_results = if let Some(query_vector) = query_vector {
let mut rowids = Vec::new();
let bin_vector = query_vector_bin.unwrap_or(query_vector);
{
let sql = build_vector_phase1_sql(filter, final_limit);
let mut params: Vec<rusqlite::types::Value> = vec![
rusqlite::types::Value::Text(bin_vector.to_string()),
rusqlite::types::Value::Text(query_vector.to_string()),
];
params.extend(vector_filter_values(filter));
let mut statement = tx.prepare(&sql)?;
let rows = statement.query_map(rusqlite::params_from_iter(params.iter()), |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, f64>(1)?))
})?;
for row in rows.flatten() {
rowids.push(row);
}
}
let mut results = Vec::new();
let mut statement =
tx.prepare("SELECT kind, body FROM canonical_nodes WHERE write_cursor = ?1 LIMIT 1")?;
for (rowid, score) in rowids {
if let Ok((kind, body)) = statement
.query_row([rowid], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)))
{
results.push(SearchHit {
id: rowid as u64,
kind,
body,
score,
branch: SoftFallbackBranch::Vector,
});
}
}
results
} else {
Vec::new()
};
let vector_rows_visible = !vector_results.is_empty();
let soft_fallback = if query_vector.is_some() && !vector_rows_visible {
tx.query_row(
"SELECT 1
FROM search_index
JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = search_index.kind
LEFT JOIN _fathomdb_projection_terminal
ON _fathomdb_projection_terminal.write_cursor = search_index.write_cursor
WHERE search_index MATCH ?1
AND _fathomdb_projection_terminal.write_cursor IS NULL
LIMIT 1",
[compiled.match_expression.as_str()],
|_row| Ok(SoftFallback { branch: SoftFallbackBranch::Vector }),
)
.ok()
} else {
None
};
let text_candidates: Vec<SearchHit> = {
let perf_limit: Option<usize> = if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_some() {
std::env::var("FATHOMDB_PERF_SEARCH_LIMIT").ok().and_then(|s| s.parse().ok())
} else {
None
};
let sql = match perf_limit {
Some(k) => format!(
"SELECT body, kind, write_cursor, bm25(search_index) FROM search_index \
WHERE search_index MATCH ?1 ORDER BY write_cursor LIMIT {k}"
),
None => "SELECT body, kind, write_cursor, bm25(search_index) FROM search_index \
WHERE search_index MATCH ?1 ORDER BY write_cursor"
.to_string(),
};
let mut statement = tx.prepare(&sql)?;
let rows = statement.query_map([compiled.match_expression.as_str()], |row| {
Ok(SearchHit {
body: row.get::<_, String>(0)?,
kind: row.get::<_, String>(1)?,
id: row.get::<_, i64>(2)? as u64,
score: row.get::<_, f64>(3)?,
branch: SoftFallbackBranch::Text,
})
})?;
rows.flatten().collect()
};
let mut text_results: Vec<SearchHit> = Vec::with_capacity(text_candidates.len());
for hit in text_candidates {
if text_hit_passes_filter(&tx, hit.id, &hit.kind, filter)? {
text_results.push(hit);
}
}
tx.commit()?;
let results = if vector_stage_only {
vector_results
} else {
rerank_fused(apply_recency_reweight(
fuse_rrf(vector_results, text_results),
recency_enabled,
))
};
Ok((cursor, soft_fallback, results))
}
const READ_COLLECTION_MAX_LIMIT: usize = 1_000_000;
fn read_get_by_id_in_tx(
reader: &mut Connection,
logical_ids: &[String],
) -> rusqlite::Result<Vec<Option<NodeRecord>>> {
if logical_ids.is_empty() {
return Ok(Vec::new());
}
let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
let mut found: HashMap<String, NodeRecord> = HashMap::new();
{
let unique: Vec<&String> = {
let mut seen = std::collections::HashSet::new();
logical_ids.iter().filter(|id| seen.insert((*id).clone())).collect()
};
let placeholders = std::iter::repeat_n("?", unique.len()).collect::<Vec<_>>().join(", ");
let sql = format!(
"SELECT logical_id, kind, body, write_cursor
FROM canonical_nodes
WHERE logical_id IN ({placeholders}) AND superseded_at IS NULL"
);
let mut statement = tx.prepare(&sql)?;
let params = rusqlite::params_from_iter(unique.iter().map(|s| s.as_str()));
let rows = statement.query_map(params, |row| {
let logical_id: String = row.get(0)?;
Ok(NodeRecord {
logical_id,
kind: row.get(1)?,
body: row.get(2)?,
write_cursor: row.get::<_, i64>(3)? as u64,
})
})?;
for row in rows {
let record = row?;
found.insert(record.logical_id.clone(), record);
}
}
let out = logical_ids.iter().map(|id| found.get(id).cloned()).collect();
Ok(out)
}
fn read_collection_in_tx(
reader: &mut Connection,
collection: &str,
after_id: Option<i64>,
limit: usize,
) -> rusqlite::Result<Vec<OpStoreRow>> {
if limit == 0 {
return Ok(Vec::new());
}
let clamped = limit.min(READ_COLLECTION_MAX_LIMIT) as i64;
let after = after_id.unwrap_or(0).max(0);
let tx = reader.transaction_with_behavior(rusqlite::TransactionBehavior::Deferred)?;
let mut statement = tx.prepare(
"SELECT id, collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
FROM operational_mutations
WHERE collection_name = ?1 AND id > ?2
ORDER BY id
LIMIT ?3",
)?;
let rows = statement.query_map(params![collection, after, clamped], |row| {
Ok(OpStoreRow {
id: row.get(0)?,
collection: row.get(1)?,
record_key: row.get(2)?,
op_kind: row.get(3)?,
payload: row.get(4)?,
schema_id: row.get(5)?,
write_cursor: row.get::<_, i64>(6)? as u64,
})
})?;
let mut out = Vec::new();
for row in rows {
out.push(row?);
}
Ok(out)
}
fn projection_dispatcher_loop(shared: Arc<ProjectionRuntimeShared>) {
let connection = match open_runtime_connection(&shared.path) {
Ok(connection) => connection,
Err(_) => return,
};
loop {
let in_flight = {
let mut state = match shared.state.lock() {
Ok(state) => state,
Err(_) => return,
};
while !state.stopping
&& (!state.pending_scan
|| state.frozen
|| state.active_jobs + state.queued_jobs >= PROJECTION_INFLIGHT_LIMIT)
{
state = match shared.state_cvar.wait(state) {
Ok(state) => state,
Err(_) => return,
};
}
if state.stopping {
return;
}
state.pending_scan = false;
state.in_flight.clone()
};
let budget = {
let state = match shared.state.lock() {
Ok(state) => state,
Err(_) => return,
};
PROJECTION_INFLIGHT_LIMIT.saturating_sub(state.active_jobs + state.queued_jobs)
};
let fetch_cap = budget.clamp(1, PROJECTION_SCAN_FETCH);
match next_pending_projection_jobs(&connection, &in_flight, fetch_cap) {
Ok(jobs) if !jobs.is_empty() => {
if let Ok(mut state) = shared.state.lock() {
state.queued_jobs = state.queued_jobs.saturating_add(jobs.len());
for job in &jobs {
state.in_flight.insert(job.cursor);
}
state.pending_scan = true;
shared.state_cvar.notify_all();
}
if let Ok(mut queue) = shared.queue.lock() {
for job in jobs {
queue.push_back(job);
}
shared.queue_cvar.notify_all();
}
}
Ok(_) => {}
Err(_) => {
if let Ok(mut state) = shared.state.lock() {
state.pending_scan = false;
shared.state_cvar.notify_all();
}
}
}
}
}
fn projection_worker_loop(shared: Arc<ProjectionRuntimeShared>) {
let mut connection = match open_runtime_connection(&shared.path) {
Ok(connection) => connection,
Err(_) => return,
};
if ensure_vector_partition(&mut connection, shared.embedder_identity.dimension).is_err() {
return;
}
loop {
let jobs = {
let mut queue = match shared.queue.lock() {
Ok(queue) => queue,
Err(_) => return,
};
loop {
let stopping = shared.state.lock().map(|state| state.stopping).unwrap_or(true);
if stopping && queue.is_empty() {
return;
}
if let Some(job) = queue.pop_front() {
let mut jobs = vec![job];
while jobs.len() < PROJECTION_COMMIT_BATCH {
let Some(job) = queue.pop_front() else {
break;
};
jobs.push(job);
}
if let Ok(mut state) = shared.state.lock() {
state.queued_jobs = state.queued_jobs.saturating_sub(jobs.len());
state.active_jobs = state.active_jobs.saturating_add(jobs.len());
shared.state_cvar.notify_all();
}
break jobs;
}
queue = match shared.queue_cvar.wait(queue) {
Ok(queue) => queue,
Err(_) => return,
};
}
};
let panicked = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
run_projection_jobs(&shared, &mut connection, &jobs);
}))
.is_err();
if panicked {
commit_projection_panic_failures(&shared, &mut connection, &jobs);
}
if let Ok(mut state) = shared.state.lock() {
state.active_jobs = state.active_jobs.saturating_sub(jobs.len());
for job in &jobs {
state.in_flight.remove(&job.cursor);
}
if !state.stopping {
state.pending_scan = true;
}
shared.state_cvar.notify_all();
}
}
}
enum ProjectionOutcome {
Success {
cursor: u64,
kind: String,
blob: Vec<u8>,
bin_blob: Vec<u8>,
},
Failure {
cursor: u64,
failure_code: &'static str,
},
}
fn run_projection_jobs(
shared: &ProjectionRuntimeShared,
connection: &mut Connection,
jobs: &[ProjectionJob],
) {
let mut outcomes = Vec::with_capacity(jobs.len());
for job in jobs {
outcomes.push(run_projection_job(shared, job));
}
let _ = commit_projection_outcomes(connection, &outcomes, shared);
}
fn commit_projection_panic_failures(
shared: &ProjectionRuntimeShared,
connection: &mut Connection,
jobs: &[ProjectionJob],
) {
let outcomes: Vec<ProjectionOutcome> = jobs
.iter()
.map(|job| ProjectionOutcome::Failure {
cursor: job.cursor,
failure_code: "ProjectionPanic",
})
.collect();
let _ = commit_projection_outcomes(connection, &outcomes, shared);
}
fn embed_with_watchdog(
embedder: &Arc<dyn Embedder>,
body: &str,
timeout: Duration,
live: &Arc<AtomicU64>,
) -> Result<Vec<f32>, RuntimeEmbedderError> {
let (tx, rx) = mpsc::channel();
let embedder = Arc::clone(embedder);
let body = body.to_string();
live.fetch_add(1, Ordering::Relaxed);
let live_thread = Arc::clone(live);
thread::spawn(move || {
let outcome =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| embedder.embed(&body)));
let _ = tx.send(outcome);
live_thread.fetch_sub(1, Ordering::Relaxed);
});
match rx.recv_timeout(timeout) {
Ok(Ok(result)) => result,
Ok(Err(panic_payload)) => std::panic::resume_unwind(panic_payload),
Err(mpsc::RecvTimeoutError::Timeout) => Err(RuntimeEmbedderError::Timeout),
Err(mpsc::RecvTimeoutError::Disconnected) => Err(RuntimeEmbedderError::Failed {
message: "embed watchdog thread dropped its result channel".to_string(),
}),
}
}
fn run_projection_job(shared: &ProjectionRuntimeShared, job: &ProjectionJob) -> ProjectionOutcome {
if shared.embed_circuit_open.load(Ordering::Relaxed) {
return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: "EmbedderError" };
}
let delays = shared.retry_delays_ms.lock().map(|delays| delays.clone()).unwrap_or_default();
let mut last_code = "EmbedderError";
for (attempt, delay_ms) in std::iter::once(0_u64).chain(delays.iter().copied()).enumerate() {
if attempt > 0 {
if shared.state.lock().map(|state| state.stopping).unwrap_or(true) {
return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code };
}
thread::sleep(Duration::from_millis(delay_ms));
}
if shared.embed_circuit_open.load(Ordering::Relaxed) {
return ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code };
}
let embed_timeout = Duration::from_millis(shared.embed_timeout_ms.load(Ordering::Relaxed));
let vector = match shared.embedder.as_ref() {
Some(embedder) => {
let _embed_permit =
shared.embed_serialize.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
let threshold = shared.embed_circuit_threshold.load(Ordering::Relaxed);
if shared.embed_circuit_open.load(Ordering::Relaxed)
|| (threshold != 0
&& shared.live_embed_threads.load(Ordering::Relaxed) >= threshold)
{
shared.embed_circuit_open.store(true, Ordering::Relaxed);
return ProjectionOutcome::Failure {
cursor: job.cursor,
failure_code: last_code,
};
}
match embed_with_watchdog(
embedder,
&job.body,
embed_timeout,
&shared.live_embed_threads,
) {
Ok(vector) => vector,
Err(RuntimeEmbedderError::Timeout) => {
last_code = "EmbedderError";
continue;
}
Err(RuntimeEmbedderError::Failed { .. }) => {
last_code = "EmbedderError";
continue;
}
}
}
None => {
last_code = "EmbedderNotConfiguredError";
continue;
}
};
if u32::try_from(vector.len()).unwrap_or(u32::MAX) != shared.embedder_identity.dimension {
last_code = "EmbedderDimensionMismatchError";
continue;
}
let blob = encode_vector_blob(&vector);
let bin_blob = blob.clone();
return ProjectionOutcome::Success {
cursor: job.cursor,
kind: job.kind.clone(),
blob,
bin_blob,
};
}
ProjectionOutcome::Failure { cursor: job.cursor, failure_code: last_code }
}
fn next_pending_projection_jobs(
connection: &Connection,
in_flight: &BTreeSet<u64>,
max_jobs: usize,
) -> rusqlite::Result<Vec<ProjectionJob>> {
if max_jobs == 0 {
return Ok(Vec::new());
}
let cursor = load_projection_cursor(connection)?;
let sql_limit = max_jobs.saturating_add(in_flight.len()).min(256);
let sql = format!(
"SELECT canonical_nodes.write_cursor, canonical_nodes.kind, canonical_nodes.body
FROM canonical_nodes
JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = canonical_nodes.kind
LEFT JOIN _fathomdb_projection_terminal
ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
WHERE canonical_nodes.write_cursor > ?1
AND _fathomdb_projection_terminal.write_cursor IS NULL
ORDER BY canonical_nodes.write_cursor
LIMIT {sql_limit}"
);
let mut statement = connection.prepare_cached(&sql)?;
let rows = statement.query_map([cursor], |row| {
Ok(ProjectionJob { cursor: row.get(0)?, kind: row.get(1)?, body: row.get(2)? })
})?;
let mut jobs = Vec::with_capacity(max_jobs);
for row in rows {
let job = row?;
if in_flight.contains(&job.cursor) {
continue;
}
jobs.push(job);
if jobs.len() >= max_jobs {
break;
}
}
Ok(jobs)
}
fn database_has_pending_projection_work(path: &Path) -> rusqlite::Result<bool> {
let connection = open_runtime_connection(path)?;
let cursor = load_projection_cursor(&connection)?;
connection
.query_row(
"SELECT 1
FROM canonical_nodes
JOIN _fathomdb_vector_kinds ON _fathomdb_vector_kinds.kind = canonical_nodes.kind
LEFT JOIN _fathomdb_projection_terminal
ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
WHERE canonical_nodes.write_cursor > ?1
AND _fathomdb_projection_terminal.write_cursor IS NULL
LIMIT 1",
[cursor],
|_row| Ok(true),
)
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(false),
_ => Err(err),
})
}
struct CanonicalNodeRow {
cursor: u64,
kind: String,
body: String,
}
fn reproject_search_index_after_tokenizer_upgrade(connection: &Connection) -> rusqlite::Result<()> {
let rows = canonical_node_rows(connection)?;
connection.execute_batch("BEGIN IMMEDIATE")?;
let result = (|| {
connection.execute("DELETE FROM search_index", [])?;
{
let mut statement = connection
.prepare("INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)")?;
for row in &rows {
statement.execute(params![row.body, row.kind, row.cursor])?;
}
}
connection.execute(
"INSERT INTO _fathomdb_open_state(key, value) VALUES(?1, ?2)
ON CONFLICT(key) DO UPDATE SET value = excluded.value",
params![SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY, "1"],
)?;
Ok(())
})();
match result {
Ok(()) => connection.execute_batch("COMMIT"),
Err(err) => {
let _ = connection.execute_batch("ROLLBACK");
Err(err)
}
}
}
fn search_index_tokenizer_reproject_complete(connection: &Connection) -> rusqlite::Result<bool> {
match connection.query_row(
"SELECT value FROM _fathomdb_open_state WHERE key = ?1",
[SEARCH_INDEX_TOKENIZER_REPROJECT_MARKER_KEY],
|row| row.get::<_, String>(0),
) {
Ok(value) => Ok(value == "1"),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(false),
Err(rusqlite::Error::SqliteFailure(_, Some(ref message)))
if message.contains("no such table") =>
{
Ok(true)
}
Err(err) => Err(err),
}
}
fn canonical_node_rows(connection: &Connection) -> rusqlite::Result<Vec<CanonicalNodeRow>> {
let mut statement = connection
.prepare("SELECT write_cursor, kind, body FROM canonical_nodes ORDER BY write_cursor")?;
let rows = statement.query_map([], |row| {
Ok(CanonicalNodeRow {
cursor: row.get::<_, u64>(0)?,
kind: row.get::<_, String>(1)?,
body: row.get::<_, String>(2)?,
})
})?;
rows.collect()
}
#[cfg(feature = "operator")]
fn hex_encode(bytes: &[u8]) -> String {
let mut out = String::with_capacity(bytes.len() * 2);
for byte in bytes {
out.push(hex_nibble(byte >> 4));
out.push(hex_nibble(byte & 0x0f));
}
out
}
#[cfg(feature = "operator")]
fn hex_nibble(value: u8) -> char {
match value {
0..=9 => (b'0' + value) as char,
10..=15 => (b'a' + value - 10) as char,
_ => unreachable!(),
}
}
#[cfg(feature = "operator")]
fn physical_section(connection: &Connection, full: bool) -> Section {
let mut findings = Vec::new();
if let Err(err) = connection.query_row("PRAGMA page_count", [], |row| row.get::<_, i64>(0)) {
findings.push(Finding {
code: "E_CORRUPT_HEADER",
stage: "PhysicalProbe",
locator: locator_from_rusqlite_error(&err),
doc_anchor: "design/recovery.md#header-malformed",
detail: format!("page_count probe failed: {err}"),
});
}
if full {
match collect_integrity_check_findings(connection) {
Ok(rows) => findings.extend(rows),
Err(err) => findings.push(Finding {
code: "E_CORRUPT_INTEGRITY_CHECK",
stage: "IntegrityCheck",
locator: locator_from_rusqlite_error(&err),
doc_anchor: "design/recovery.md#integrity-check-full-findings",
detail: format!("PRAGMA integrity_check failed: {err}"),
}),
}
}
if findings.is_empty() {
Section::Clean
} else {
Section::Findings(findings)
}
}
#[cfg(feature = "operator")]
fn logical_section(connection: &Connection) -> Section {
let mut findings = Vec::new();
if let Err(err) = connection.query_row("PRAGMA schema_version", [], |row| row.get::<_, i64>(0))
{
findings.push(Finding {
code: "E_CORRUPT_SCHEMA",
stage: "SchemaProbe",
locator: locator_from_rusqlite_error(&err),
doc_anchor: "design/recovery.md#schema-inconsistent",
detail: format!("schema_version probe failed: {err}"),
});
}
match connection.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)) {
Ok(0) => findings.push(Finding {
code: "E_CORRUPT_SCHEMA",
stage: "SchemaProbe",
locator: CorruptionLocator::MigrationStep { from: 0, to: 0 },
doc_anchor: "design/recovery.md#schema-inconsistent",
detail: "user_version is zero".to_string(),
}),
Ok(_) => {}
Err(err) => findings.push(Finding {
code: "E_CORRUPT_SCHEMA",
stage: "SchemaProbe",
locator: locator_from_rusqlite_error(&err),
doc_anchor: "design/recovery.md#schema-inconsistent",
detail: format!("user_version probe failed: {err}"),
}),
}
if findings.is_empty() {
Section::Clean
} else {
Section::Findings(findings)
}
}
#[cfg(feature = "operator")]
fn semantic_section(connection: &Connection) -> Section {
match load_default_profile(connection) {
Ok(_) => Section::Clean,
Err(rusqlite::Error::QueryReturnedNoRows) => Section::Findings(vec![Finding {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
stage: "EmbedderIdentity",
locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
doc_anchor: "design/recovery.md#embedder-identity-drift",
detail: "default embedder profile row is missing".to_string(),
}]),
Err(err) => Section::Findings(vec![Finding {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
stage: "EmbedderIdentity",
locator: locator_from_rusqlite_error(&err),
doc_anchor: "design/recovery.md#embedder-identity-drift",
detail: format!("default embedder profile probe failed: {err}"),
}]),
}
}
#[cfg(feature = "operator")]
fn collect_integrity_check_findings(connection: &Connection) -> rusqlite::Result<Vec<Finding>> {
let mut statement = connection.prepare("PRAGMA integrity_check")?;
let rows = statement.query_map([], |row| row.get::<_, String>(0))?;
let mut findings = Vec::new();
for row in rows {
let message = row?;
if message == "ok" {
continue;
}
findings.push(Finding {
code: "E_CORRUPT_INTEGRITY_CHECK",
stage: "IntegrityCheck",
locator: CorruptionLocator::OpaqueSqliteError {
sqlite_extended_code: rusqlite::ffi::SQLITE_CORRUPT,
},
doc_anchor: "design/recovery.md#integrity-check-full-findings",
detail: message,
});
}
Ok(findings)
}
#[cfg(feature = "operator")]
fn locator_from_rusqlite_error(err: &rusqlite::Error) -> CorruptionLocator {
let extended = err.sqlite_error().map(|inner| inner.extended_code).unwrap_or(0);
CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: extended }
}
fn open_runtime_connection(path: &Path) -> rusqlite::Result<Connection> {
let connection = Connection::open(path)?;
connection.pragma_update(None, "journal_mode", "WAL")?;
Ok(connection)
}
fn load_projection_cursor(connection: &Connection) -> rusqlite::Result<u64> {
connection
.query_row(
"SELECT value FROM _fathomdb_open_state WHERE key = ?1",
[PROJECTION_CURSOR_KEY],
|row| row.get::<_, String>(0),
)
.map(|value| value.parse::<u64>().unwrap_or(0))
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(0),
_ => Err(err),
})
}
fn store_projection_cursor(connection: &Connection, cursor: u64) -> rusqlite::Result<()> {
connection.execute(
"INSERT INTO _fathomdb_open_state(key, value) VALUES(?1, ?2)
ON CONFLICT(key) DO UPDATE SET value = excluded.value",
params![PROJECTION_CURSOR_KEY, cursor.to_string()],
)?;
Ok(())
}
fn record_projection_terminal(
connection: &Connection,
cursor: u64,
state: &str,
) -> rusqlite::Result<()> {
connection.execute(
"INSERT OR IGNORE INTO _fathomdb_projection_terminal(write_cursor, state) VALUES(?1, ?2)",
params![cursor, state],
)?;
Ok(())
}
fn terminal_state_for_cursor(
connection: &Connection,
cursor: u64,
) -> rusqlite::Result<Option<String>> {
connection
.query_row(
"SELECT state FROM _fathomdb_projection_terminal WHERE write_cursor = ?1",
[cursor],
|row| row.get::<_, String>(0),
)
.map(Some)
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(None),
_ => Err(err),
})
}
fn advance_projection_cursor(connection: &Connection) -> rusqlite::Result<u64> {
let mut cursor = load_projection_cursor(connection)?;
loop {
let next = cursor.saturating_add(1);
if terminal_state_for_cursor(connection, next)?.is_some() {
cursor = next;
} else {
break;
}
}
store_projection_cursor(connection, cursor)?;
Ok(cursor)
}
fn commit_projection_outcomes(
connection: &mut Connection,
outcomes: &[ProjectionOutcome],
shared: &ProjectionRuntimeShared,
) -> rusqlite::Result<()> {
let embedder_identity = &shared.embedder_identity;
let mc = identity_requires_mean_centering(embedder_identity);
let _gate = shared.commit_gate.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
let tx = connection.transaction()?;
let mut current_mean: Option<Vec<f32>> = if mc {
tx.query_row(
"SELECT mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
[],
|row| row.get::<_, Option<Vec<u8>>>(0),
)
.ok()
.flatten()
.map(|bytes| decode_vector_blob(&bytes))
} else {
None
};
let mut staged_events: Vec<EmbedderEvent> = Vec::new();
for outcome in outcomes {
match outcome {
ProjectionOutcome::Success { cursor, kind, blob, bin_blob } => {
if terminal_state_for_cursor(&tx, *cursor)?.is_some() {
continue;
}
let pin_mean: Option<Vec<f32>> = if mc && current_mean.is_none() {
let mut acc = shared.mean_accumulator.lock().unwrap_or_else(|p| p.into_inner());
match acc.as_mut() {
Some(a) => {
a.add(&decode_vector_blob(bin_blob));
if a.count() >= MEAN_VEC_PIN_THRESHOLD {
let mean = a.materialize();
*acc = None;
Some(mean)
} else {
None
}
}
None => None,
}
} else {
None
};
let source_type = resolve_source_type(kind).map_err(|_| {
rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
Some(format!("unknown kind for source_type mapping: {kind}")),
)
})?;
let now_unix =
SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs()
as i64;
tx.execute(
"INSERT OR IGNORE INTO _fathomdb_vector_rows(rowid, kind, write_cursor) VALUES(?1, ?2, ?3)",
params![cursor, kind, cursor],
)?;
let centered_blob: Vec<u8> = match ¤t_mean {
Some(mean) if mean.len() * 4 == bin_blob.len() => {
encode_vector_blob(&subtract_mean(&decode_vector_blob(bin_blob), mean))
}
_ => bin_blob.clone(),
};
tx.execute(
"INSERT OR IGNORE INTO vector_default(
rowid, embedding, embedding_bin, source_type, kind, created_at, status
) VALUES(?1, ?2, vec_quantize_binary(?3), ?4, ?5, ?6, '')",
params![cursor, blob, centered_blob, source_type, kind, now_unix],
)?;
record_projection_terminal(&tx, *cursor, "up_to_date")?;
if let Some(mean) = pin_mean {
tx.execute(
"UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
params![encode_vector_blob(&mean)],
)?;
let rows: Vec<(i64, Vec<u8>)> = {
let mut statement = tx.prepare(
"SELECT rowid, embedding FROM vector_default ORDER BY rowid",
)?;
let mapped = statement.query_map([], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?))
})?;
let mut out = Vec::new();
for r in mapped {
out.push(r?);
}
out
};
let (doc_count, _) =
run_pin_and_requantize_pass(&tx, &rows, &mean).map_err(|_| {
rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_ERROR),
Some("mean-centering re-quantize pass failed".to_string()),
)
})?;
staged_events.push(EmbedderEvent::MeanVecPinned {
dim: u32::try_from(mean.len()).unwrap_or(u32::MAX),
doc_count,
});
current_mean = Some(mean);
}
}
ProjectionOutcome::Failure { cursor, failure_code } => {
if terminal_state_for_cursor(&tx, *cursor)?.is_some() {
continue;
}
let existing: u64 = tx.query_row(
"SELECT COUNT(*) FROM operational_mutations
WHERE collection_name = 'projection_failures'
AND json_extract(payload_json, '$.write_cursor') = ?1",
[cursor],
|row| row.get(0),
)?;
if existing == 0 {
let payload = format!(
r#"{{"write_cursor":{cursor},"failure_code":"{failure_code}","recorded_at":0}}"#
);
tx.execute(
"INSERT INTO operational_mutations(
collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
) VALUES('projection_failures', ?1, 'append', ?2, NULL, ?3)",
params![cursor.to_string(), payload, cursor],
)?;
}
record_projection_terminal(&tx, *cursor, "failed")?;
}
}
}
advance_projection_cursor(&tx)?;
tx.commit()?;
if !staged_events.is_empty() {
if let Ok(mut events) = shared.pending_events.lock() {
events.extend(staged_events);
}
}
Ok(())
}
fn recover_mean_vec_pin(
connection: &mut Connection,
identity: &EmbedderIdentity,
) -> Result<(), EngineError> {
let tx = connection.transaction().map_err(|_| EngineError::Storage)?;
recompute_mean_in_tx(&tx, identity)?;
tx.commit().map_err(|_| EngineError::Storage)?;
Ok(())
}
fn recompute_mean_in_tx(
tx: &rusqlite::Transaction<'_>,
identity: &EmbedderIdentity,
) -> Result<MeanRecomputeReport, EngineError> {
recompute_mean_in_tx_inner(tx, identity, false)
}
fn recompute_mean_in_tx_inner(
tx: &rusqlite::Transaction<'_>,
identity: &EmbedderIdentity,
fail_after_mean_update: bool,
) -> Result<MeanRecomputeReport, EngineError> {
let started = Instant::now();
let dim = identity.dimension as usize;
let old_mean = read_pinned_mean_vec(tx, identity.dimension)?;
let rows: Vec<(i64, Vec<u8>)> = {
let mut statement = tx
.prepare("SELECT rowid, embedding FROM vector_default ORDER BY rowid")
.map_err(|_| EngineError::Storage)?;
let mapped = statement
.query_map([], |row| Ok((row.get::<_, i64>(0)?, row.get::<_, Vec<u8>>(1)?)))
.map_err(|_| EngineError::Storage)?;
let mut out = Vec::new();
for r in mapped {
out.push(r.map_err(|_| EngineError::Storage)?);
}
out
};
let mut accumulator = MeanAccumulator::new(dim);
for (_rowid, blob) in &rows {
if blob.len() != dim * 4 {
return Err(EngineError::Storage);
}
accumulator.add(&decode_vector_blob(blob));
}
let old_doc_count = accumulator.count();
let mean = accumulator.materialize();
let drift_cos_before = match &old_mean {
Some(old) => cosine_similarity(&mean, old),
None => 1.0,
};
tx.execute(
"UPDATE _fathomdb_embedder_profiles SET mean_vec = ?1 WHERE profile = 'default'",
params![encode_vector_blob(&mean)],
)
.map_err(|_| EngineError::Storage)?;
if fail_after_mean_update {
return Err(EngineError::Storage);
}
let (doc_count, _) = run_pin_and_requantize_pass(tx, &rows, &mean)?;
Ok(MeanRecomputeReport {
dim: u32::try_from(dim).unwrap_or(u32::MAX),
old_doc_count,
doc_count_requantized: doc_count,
drift_cos_before,
mean_was_pinned: old_mean.is_some(),
elapsed_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
})
}
fn enforce_provenance_retention(connection: &Connection, cap: u64) -> rusqlite::Result<()> {
if cap == 0 {
return Ok(());
}
let slack = cap.max(20) / 20;
let upper = cap.saturating_add(slack.max(1));
let count: u64 =
connection.query_row("SELECT COUNT(*) FROM operational_mutations", [], |row| row.get(0))?;
if count <= upper {
return Ok(());
}
let to_delete = count.saturating_sub(cap);
connection.execute(
"DELETE FROM operational_mutations
WHERE id IN (
SELECT id FROM operational_mutations
ORDER BY id
LIMIT ?1
)",
[to_delete],
)?;
Ok(())
}
fn projection_status(
connection: &Connection,
kind: &str,
) -> Result<lifecycle::ProjectionStatus, EngineError> {
let latest = connection
.query_row(
"SELECT COALESCE(MAX(write_cursor), 0) FROM canonical_nodes WHERE kind = ?1",
[kind],
|row| row.get::<_, u64>(0),
)
.map_err(|_| EngineError::Storage)?;
if latest == 0 {
return Ok(lifecycle::ProjectionStatus::UpToDate);
}
let pending: u64 = connection
.query_row(
"SELECT COUNT(*)
FROM canonical_nodes
LEFT JOIN _fathomdb_projection_terminal
ON _fathomdb_projection_terminal.write_cursor = canonical_nodes.write_cursor
WHERE canonical_nodes.kind = ?1
AND _fathomdb_projection_terminal.write_cursor IS NULL",
[kind],
|row| row.get(0),
)
.map_err(|_| EngineError::Storage)?;
if pending > 0 {
return Ok(lifecycle::ProjectionStatus::Pending);
}
match terminal_state_for_cursor(connection, latest).map_err(|_| EngineError::Storage)? {
Some(state) if state == "failed" => Ok(lifecycle::ProjectionStatus::Failed),
_ => Ok(lifecycle::ProjectionStatus::UpToDate),
}
}
fn canonical_database_path(path: &Path) -> Result<PathBuf, EngineOpenError> {
let parent = path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
let canonical_parent = parent.canonicalize().map_err(|_| EngineOpenError::Io {
message: "database parent directory is not accessible".to_string(),
})?;
let file_name = path.file_name().ok_or_else(|| EngineOpenError::Io {
message: "database path has no file name".to_string(),
})?;
Ok(canonical_parent.join(file_name))
}
fn acquire_lock(path: &Path) -> Result<File, EngineOpenError> {
let lock_path = lock_path(path);
let mut options = OpenOptions::new();
options.read(true).write(true).create(true);
#[cfg(unix)]
options.mode(0o600);
let mut file = options.open(&lock_path).map_err(|_| EngineOpenError::Io {
message: "could not open database lock file".to_string(),
})?;
match file.try_lock() {
Ok(()) => {
let pid = std::process::id().to_string();
let _ = file.set_len(0);
let _ = file.seek(SeekFrom::Start(0));
let _ = file.write_all(pid.as_bytes());
Ok(file)
}
Err(std::fs::TryLockError::WouldBlock) => {
Err(EngineOpenError::DatabaseLocked { holder_pid: read_holder_pid(&lock_path) })
}
Err(_) => {
Err(EngineOpenError::Io { message: "could not acquire database lock".to_string() })
}
}
}
fn lock_path(path: &Path) -> PathBuf {
let mut lock_path = path.as_os_str().to_os_string();
lock_path.push(LOCK_SUFFIX);
PathBuf::from(lock_path)
}
fn read_holder_pid(path: &Path) -> Option<u32> {
std::fs::read_to_string(path).ok()?.trim().parse().ok()
}
fn map_migration_error(err: SchemaMigrationError) -> EngineOpenError {
match err {
SchemaMigrationError::IncompatibleSchemaVersion { seen, supported } => {
EngineOpenError::IncompatibleSchemaVersion { seen, supported }
}
SchemaMigrationError::MigrationError(report) => EngineOpenError::MigrationError {
schema_version_before: report.schema_version_before,
schema_version_current: report.schema_version_current,
step_id: report.migration_steps.last().map_or(0, |step| step.step_id),
},
SchemaMigrationError::Storage { message } => {
EngineOpenError::Io { message: message.to_string() }
}
}
}
fn init_perf_experiments_runtime() {
static INIT: Once = Once::new();
INIT.call_once(|| {
if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
return;
}
let memstatus_off =
std::env::var_os("FATHOMDB_PERF_SQLITE_MEMSTATUS_OFF").is_some_and(|v| v == "1");
let pagecache = std::env::var("FATHOMDB_PERF_SQLITE_PAGECACHE").ok();
let pcache2_on =
std::env::var_os("FATHOMDB_PERF_SQLITE_PCACHE2").is_some_and(|v| v == "1");
if !memstatus_off && pagecache.is_none() && !pcache2_on {
return;
}
unsafe {
let rc_shutdown = rusqlite::ffi::sqlite3_shutdown();
let rc_memstatus = if memstatus_off {
rusqlite::ffi::sqlite3_config(rusqlite::ffi::SQLITE_CONFIG_MEMSTATUS, 0_i32)
} else {
-1
};
let rc_pagecache = if let Some(spec) = pagecache.as_ref() {
let mut parts = spec.split(':');
let sz = parts.next().and_then(|s| s.parse::<i32>().ok()).unwrap_or(0);
let n = parts.next().and_then(|s| s.parse::<i32>().ok()).unwrap_or(0);
if sz > 0 && n > 0 {
rusqlite::ffi::sqlite3_config(
7, std::ptr::null_mut::<std::ffi::c_void>(),
sz,
n,
)
} else {
eprintln!(
"perf-experiment: bad FATHOMDB_PERF_SQLITE_PAGECACHE spec '{spec}' (expect '<bytes>:<count>')"
);
-1
}
} else {
-1
};
let rc_pcache2 = if pcache2_on {
rusqlite::ffi::sqlite3_config(
rusqlite::ffi::SQLITE_CONFIG_PCACHE2,
&raw const pcache2::PCACHE2_METHODS.0,
)
} else {
-1
};
let rc_init = rusqlite::ffi::sqlite3_initialize();
eprintln!(
"perf-experiment: runtime-config rcs shutdown={rc_shutdown} \
memstatus={rc_memstatus} pagecache={rc_pagecache} pcache2={rc_pcache2} \
initialize={rc_init} (0=SQLITE_OK; 21=SQLITE_MISUSE; -1=not configured)"
);
}
});
}
fn register_sqlite_vec_extension() {
static REGISTER: Once = Once::new();
REGISTER.call_once(|| unsafe {
let entrypoint: unsafe extern "C" fn(
*mut rusqlite::ffi::sqlite3,
*mut *const std::os::raw::c_char,
*const rusqlite::ffi::sqlite3_api_routines,
) -> std::os::raw::c_int = std::mem::transmute(sqlite3_vec_init as *const ());
rusqlite::ffi::sqlite3_auto_extension(Some(entrypoint));
});
}
fn probe_open_integrity(connection: &Connection) -> Result<(), EngineOpenError> {
connection
.query_row("SELECT COUNT(*) FROM sqlite_schema", [], |row| row.get::<_, i64>(0))
.map(|_| ())
.map_err(|err| map_open_sqlite_error(err, OpenStage::SchemaProbe))
}
fn probe_database_header(connection: &Connection) -> Result<(), EngineOpenError> {
connection
.query_row("PRAGMA application_id", [], |row| row.get::<_, i64>(0))
.map(|_| ())
.map_err(|err| map_open_sqlite_error(err, OpenStage::HeaderProbe))
}
fn probe_wal_sidecar(db_path: &Path) -> Result<(), EngineOpenError> {
let mut wal_path = db_path.as_os_str().to_owned();
wal_path.push("-wal");
let wal_path = PathBuf::from(wal_path);
use std::io::Read;
let mut file = match std::fs::File::open(&wal_path) {
Ok(file) => file,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(_) => return Ok(()),
};
let mut bytes = [0u8; 32];
if file.read_exact(&mut bytes).is_err() {
return Ok(());
}
let magic = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
let page_size = u32::from_be_bytes([bytes[8], bytes[9], bytes[10], bytes[11]]);
const WAL_MAGIC_MASK: u32 = 0xFFFF_FFFE;
const WAL_MAGIC: u32 = 0x377F_0682;
const SQLITE_MAX_PAGE_SIZE: u32 = 65536;
let magic_ok = (magic & WAL_MAGIC_MASK) == WAL_MAGIC;
let page_size_ok =
page_size.is_power_of_two() && (512..=SQLITE_MAX_PAGE_SIZE).contains(&page_size);
if magic_ok && page_size_ok {
return Ok(());
}
Err(EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::WalReplayFailure,
stage: OpenStage::WalReplay,
locator: CorruptionLocator::FileOffset { offset: if !magic_ok { 0 } else { 8 } },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_WAL_REPLAY",
doc_anchor: "design/recovery.md#wal-replay-failures",
},
}))
}
fn reject_legacy_shape(connection: &Connection) -> Result<(), EngineOpenError> {
let has_legacy_table = table_exists(connection, "fathom_nodes")
|| table_exists(connection, "fathom_edges")
|| table_exists(connection, "fathom_chunks");
if !has_legacy_table {
return Ok(());
}
let seen =
connection.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)).unwrap_or(0);
Err(EngineOpenError::IncompatibleSchemaVersion { seen, supported: SCHEMA_VERSION })
}
fn table_exists(connection: &Connection, table: &str) -> bool {
connection
.query_row(
"SELECT 1 FROM sqlite_schema WHERE type = 'table' AND name = ?1",
[table],
|_row| Ok(()),
)
.is_ok()
}
#[cfg(feature = "operator")]
fn read_schema_objects(
connection: &Connection,
obj_type: &str,
) -> Result<Vec<SchemaObject>, EngineError> {
let mut stmt = connection
.prepare(
"SELECT name, sql FROM sqlite_schema
WHERE type = ?1 AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL
ORDER BY name",
)
.map_err(|_| EngineError::Storage)?;
let rows = stmt
.query_map([obj_type], |row| {
Ok(SchemaObject { name: row.get::<_, String>(0)?, sql: row.get::<_, String>(1)? })
})
.map_err(|_| EngineError::Storage)?;
let mut out = Vec::new();
for row in rows {
out.push(row.map_err(|_| EngineError::Storage)?);
}
Ok(out)
}
#[cfg(feature = "operator")]
fn order_canonical_first(mut objects: Vec<SchemaObject>) -> Vec<SchemaObject> {
let mut canonical: Vec<SchemaObject> = Vec::new();
for name in CANONICAL_TABLES {
if let Some(pos) = objects.iter().position(|o| o.name == *name) {
canonical.push(objects.remove(pos));
}
}
canonical.extend(objects);
canonical
}
fn load_default_profile(connection: &Connection) -> rusqlite::Result<EmbedderIdentity> {
connection.query_row(
"SELECT name, revision, dimension FROM _fathomdb_embedder_profiles WHERE profile = ?1",
[DEFAULT_VECTOR_PROFILE],
|row| {
Ok(EmbedderIdentity::new(
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, u32>(2)?,
))
},
)
}
fn default_profile_dimension(connection: &Connection) -> Result<u32, EngineError> {
load_default_profile(connection)
.map(|identity| identity.dimension)
.map_err(|_| EngineError::Storage)
}
fn kind_is_vector_indexed(connection: &Connection, kind: &str) -> Result<bool, EngineError> {
connection
.query_row("SELECT 1 FROM _fathomdb_vector_kinds WHERE kind = ?1", [kind], |_row| Ok(()))
.map(|_| true)
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(false),
_ => Err(EngineError::Storage),
})
}
fn ensure_vector_partition(connection: &mut Connection, dimension: u32) -> rusqlite::Result<()> {
let existing_sql: Option<String> = connection
.query_row(
"SELECT sql FROM sqlite_master WHERE type='table' AND name=?1",
[DEFAULT_VECTOR_PARTITION],
|row| row.get::<_, String>(0),
)
.optional()?;
match existing_sql {
None => create_vector_partition(connection, dimension),
Some(sql) if sql.contains("status") => Ok(()),
Some(sql) if sql.contains("embedding_bin") => {
migrate_vector_partition_pack1_to_pack2(connection, dimension)
}
Some(_) => migrate_vector_partition_to_pack1(connection, dimension),
}
}
fn vector_partition_create_sql(dimension: u32, if_not_exists: bool) -> String {
let guard = if if_not_exists { "IF NOT EXISTS " } else { "" };
format!(
"CREATE VIRTUAL TABLE {guard}{DEFAULT_VECTOR_PARTITION} USING vec0(\
embedding float[{dimension}],\
embedding_bin bit[{dimension}],\
source_type TEXT partition key,\
kind TEXT,\
created_at INTEGER,\
status TEXT\
)"
)
}
fn create_vector_partition(connection: &Connection, dimension: u32) -> rusqlite::Result<()> {
connection.execute_batch(&vector_partition_create_sql(dimension, true))
}
fn migrate_vector_partition_pack1_to_pack2(
connection: &mut Connection,
dimension: u32,
) -> rusqlite::Result<()> {
let tx = connection.transaction()?;
tx.execute_batch(
"CREATE TABLE _fathomdb_vector_pack2_stage (
rowid INTEGER PRIMARY KEY,
embedding BLOB NOT NULL,
embedding_bin BLOB NOT NULL,
source_type TEXT,
kind TEXT,
created_at INTEGER
);
INSERT INTO _fathomdb_vector_pack2_stage(
rowid, embedding, embedding_bin, source_type, kind, created_at
)
SELECT rowid, embedding, embedding_bin, source_type, kind, created_at
FROM vector_default;
DROP TABLE vector_default;",
)?;
tx.execute_batch(&vector_partition_create_sql(dimension, false))?;
tx.execute_batch(
"INSERT INTO vector_default(
rowid, embedding, embedding_bin, source_type, kind, created_at, status
)
SELECT rowid, embedding, vec_bit(embedding_bin), source_type, kind, created_at, ''
FROM _fathomdb_vector_pack2_stage;
DROP TABLE _fathomdb_vector_pack2_stage;",
)?;
tx.commit()
}
const KIND_TO_SOURCE_TYPE_CASE_SQL: &str = "CASE s.kind
WHEN 'email' THEN 'email'
WHEN 'article' THEN 'article'
WHEN 'paper' THEN 'paper'
WHEN 'meeting' THEN 'meeting'
WHEN 'note' THEN 'note'
WHEN 'todo' THEN 'todo'
WHEN 'doc' THEN 'article'
ELSE 'article'
END";
fn migrate_vector_partition_to_pack1(
connection: &mut Connection,
dimension: u32,
) -> rusqlite::Result<()> {
let tx = connection.transaction()?;
tx.execute_batch(
"CREATE TABLE _fathomdb_vector_migration_v0_7_0 (
rowid INTEGER PRIMARY KEY,
embedding BLOB NOT NULL,
kind TEXT NOT NULL
);
INSERT INTO _fathomdb_vector_migration_v0_7_0(rowid, embedding, kind)
SELECT v.rowid, v.embedding, r.kind
FROM vector_default v
JOIN _fathomdb_vector_rows r ON r.rowid = v.rowid;
DROP TABLE vector_default;",
)?;
tx.execute_batch(&vector_partition_create_sql(dimension, false))?;
let repopulate_sql = format!(
"INSERT INTO vector_default(
rowid, embedding, embedding_bin, source_type, kind, created_at, status
)
SELECT
s.rowid,
s.embedding,
vec_quantize_binary(s.embedding),
{KIND_TO_SOURCE_TYPE_CASE_SQL},
s.kind,
strftime('%s', 'now'),
''
FROM _fathomdb_vector_migration_v0_7_0 s;
DROP TABLE _fathomdb_vector_migration_v0_7_0;"
);
tx.execute_batch(&repopulate_sql)?;
tx.commit()
}
fn encode_vector_blob(vector: &[f32]) -> Vec<u8> {
vector.iter().flat_map(|value| value.to_le_bytes()).collect()
}
fn decode_vector_blob(bytes: &[u8]) -> Vec<f32> {
debug_assert_eq!(bytes.len() % 4, 0, "f32 BLOB length must be multiple of 4");
bytes.chunks_exact(4).map(|c| f32::from_le_bytes([c[0], c[1], c[2], c[3]])).collect()
}
fn identity_requires_mean_centering(identity: &EmbedderIdentity) -> bool {
identity.name == BGE_SMALL_EMBEDDER_NAME
}
fn read_pinned_mean_vec(
connection: &Connection,
dimension: u32,
) -> Result<Option<Vec<f32>>, EngineError> {
let bytes: Option<Vec<u8>> = connection
.query_row(
"SELECT mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
[],
|row| row.get::<_, Option<Vec<u8>>>(0),
)
.or_else(|err| match err {
rusqlite::Error::QueryReturnedNoRows => Ok(None),
other => Err(other),
})
.map_err(|_| EngineError::Storage)?;
let Some(bytes) = bytes else { return Ok(None) };
let expected_len = (dimension as usize).saturating_mul(4);
if bytes.len() != expected_len {
return Err(EngineError::Storage);
}
let mut out = Vec::with_capacity(dimension as usize);
for chunk in bytes.chunks_exact(4) {
let arr = [chunk[0], chunk[1], chunk[2], chunk[3]];
out.push(f32::from_le_bytes(arr));
}
Ok(Some(out))
}
fn subtract_mean(v: &[f32], mean: &[f32]) -> Vec<f32> {
debug_assert_eq!(v.len(), mean.len(), "subtract_mean dim mismatch");
v.iter().zip(mean.iter()).map(|(a, b)| *a - *b).collect()
}
fn resolve_source_type(kind: &str) -> Result<&'static str, EngineError> {
Ok(match kind {
"email" => "email",
"article" => "article",
"paper" => "paper",
"meeting" => "meeting",
"note" => "note",
"todo" => "todo",
"doc" => "article",
_ => return Err(EngineError::Storage),
})
}
fn map_runtime_embedder_error(err: RuntimeEmbedderError) -> EngineError {
match err {
RuntimeEmbedderError::Failed { .. } | RuntimeEmbedderError::Timeout => {
EngineError::Embedder
}
}
}
fn default_embedder_identity() -> EmbedderIdentity {
EmbedderIdentity::new(
DEFAULT_EMBEDDER_NAME,
DEFAULT_EMBEDDER_REVISION,
DEFAULT_EMBEDDER_DIMENSION,
)
}
fn check_embedder_profile(
connection: &Connection,
supplied: &EmbedderIdentity,
) -> Result<bool, EngineOpenError> {
let mut statement = match connection.prepare(
"SELECT name, revision, dimension, mean_vec FROM _fathomdb_embedder_profiles WHERE profile = 'default'",
) {
Ok(statement) => statement,
Err(_) => return Ok(false),
};
let mut rows = statement.query([]).map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?;
let Some(row) = rows.next().map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::OpaqueSqliteError { sqlite_extended_code: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?
else {
connection
.execute(
"INSERT INTO _fathomdb_embedder_profiles(profile, name, revision, dimension)
VALUES(?1, ?2, ?3, ?4)",
params![
DEFAULT_VECTOR_PROFILE,
supplied.name,
supplied.revision,
supplied.dimension
],
)
.map_err(|_| EngineOpenError::Io {
message: "could not persist embedder profile".to_string(),
})?;
return Ok(false);
};
let stored_name = row.get::<_, String>(0).map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?;
let stored_revision = row.get::<_, String>(1).map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?;
let dimension = row.get::<_, u32>(2).map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?;
let stored = EmbedderIdentity::new(stored_name, stored_revision, dimension);
if stored.name != supplied.name || stored.revision != supplied.revision {
return Err(EngineOpenError::EmbedderIdentityMismatch {
stored,
supplied: supplied.clone(),
});
}
if dimension != supplied.dimension {
return Err(EngineOpenError::EmbedderDimensionMismatch {
stored: dimension,
supplied: supplied.dimension,
});
}
let mean_vec: Option<Vec<u8>> = row.get::<_, Option<Vec<u8>>>(3).map_err(|_| {
EngineOpenError::Corruption(CorruptionDetail {
kind: CorruptionKind::EmbedderIdentityDrift,
stage: OpenStage::EmbedderIdentity,
locator: CorruptionLocator::TableRow { table: "_fathomdb_embedder_profiles", rowid: 0 },
recovery_hint: RecoveryHint {
code: "E_CORRUPT_EMBEDDER_IDENTITY",
doc_anchor: "design/recovery.md#embedder-identity-drift",
},
})
})?;
let pinned = match mean_vec {
Some(bytes) => {
let expected_len = (dimension as usize).saturating_mul(4);
if bytes.len() != expected_len {
return Err(EngineOpenError::EmbedderIdentityMismatch {
stored,
supplied: supplied.clone(),
});
}
true
}
None => false,
};
Ok(pinned)
}
#[derive(Clone, Debug, Eq, PartialEq)]
enum WritePlan {
Node,
Edge,
AppendOnlyLog,
LatestState,
AdminSchema,
}
fn validate_batch(
connection: &Connection,
batch: &[PreparedWrite],
) -> Result<Vec<WritePlan>, EngineError> {
batch.iter().map(|write| validate_write(connection, write)).collect()
}
fn collect_projection_jobs(
connection: &Connection,
batch: &[PreparedWrite],
) -> Result<Vec<ProjectionJob>, EngineError> {
let mut jobs = Vec::new();
for write in batch {
if let PreparedWrite::Node { kind, body, .. } = write {
if kind_is_vector_indexed(connection, kind)? {
jobs.push(ProjectionJob { cursor: 0, kind: kind.clone(), body: body.clone() });
}
}
}
Ok(jobs)
}
fn validate_write(
connection: &Connection,
write: &PreparedWrite,
) -> Result<WritePlan, EngineError> {
match write {
PreparedWrite::Node { kind, body, source_id, logical_id } => {
if kind.trim().is_empty() || body.trim().is_empty() {
return Err(EngineError::WriteValidation);
}
if let Some(source_id) = source_id {
if source_id.is_empty() {
return Err(EngineError::WriteValidation);
}
}
if let Some(logical_id) = logical_id {
if logical_id.is_empty() {
return Err(EngineError::WriteValidation);
}
}
Ok(WritePlan::Node)
}
PreparedWrite::Edge { kind, from, to, source_id, logical_id } => {
if kind.trim().is_empty() || from.trim().is_empty() || to.trim().is_empty() {
return Err(EngineError::WriteValidation);
}
if let Some(source_id) = source_id {
if source_id.is_empty() {
return Err(EngineError::WriteValidation);
}
}
if let Some(logical_id) = logical_id {
if logical_id.is_empty() {
return Err(EngineError::WriteValidation);
}
}
Ok(WritePlan::Edge)
}
PreparedWrite::AdminSchema { name, kind, schema_json, retention_json } => {
if name.trim().is_empty()
|| !matches!(kind.as_str(), "append_only_log" | "latest_state")
|| serde_json::from_str::<Value>(schema_json).is_err()
|| serde_json::from_str::<Value>(retention_json).is_err()
|| contains_external_ref(schema_json)
{
return Err(EngineError::SchemaValidation);
}
Ok(WritePlan::AdminSchema)
}
PreparedWrite::OpStore { collection, record_key, schema_id, body } => {
if collection.trim().is_empty() || record_key.trim().is_empty() {
return Err(EngineError::WriteValidation);
}
let (kind, schema_json) = collection_metadata(connection, collection)?;
if let Some(schema_id) = schema_id {
if schema_id != collection {
return Err(EngineError::SchemaValidation);
}
validate_payload(&schema_json, body)?;
} else if serde_json::from_str::<Value>(body).is_err() {
return Err(EngineError::SchemaValidation);
}
match kind.as_str() {
"append_only_log" => Ok(WritePlan::AppendOnlyLog),
"latest_state" => Ok(WritePlan::LatestState),
_ => Err(EngineError::OpStore),
}
}
}
}
fn collection_metadata(
connection: &Connection,
collection: &str,
) -> Result<(String, String), EngineError> {
connection
.query_row(
"SELECT kind, schema_json FROM operational_collections WHERE name = ?1",
[collection],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
)
.map_err(|_| EngineError::OpStore)
}
fn validate_payload(schema_json: &str, body: &str) -> Result<(), EngineError> {
let schema =
serde_json::from_str::<Value>(schema_json).map_err(|_| EngineError::SchemaValidation)?;
let payload = serde_json::from_str::<Value>(body).map_err(|_| EngineError::SchemaValidation)?;
let compiled = JSONSchema::compile(&schema).map_err(|_| EngineError::SchemaValidation)?;
compiled.validate(&payload).map_err(|_| EngineError::SchemaValidation)?;
Ok(())
}
fn contains_external_ref(schema_json: &str) -> bool {
let Ok(value) = serde_json::from_str::<Value>(schema_json) else {
return false;
};
value_contains_external_ref(&value)
}
fn value_contains_external_ref(value: &Value) -> bool {
match value {
Value::Object(object) => object.iter().any(|(key, value)| {
if key == "$ref" {
return value.as_str().is_some_and(|uri| !uri.starts_with('#'));
}
value_contains_external_ref(value)
}),
Value::Array(values) => values.iter().any(value_contains_external_ref),
_ => false,
}
}
fn commit_batch(
connection: &mut Connection,
batch: &[PreparedWrite],
plans: &[WritePlan],
base_cursor: u64,
provenance_row_cap: u64,
) -> rusqlite::Result<u64> {
let tx = connection.transaction()?;
for (i, (write, plan)) in batch.iter().zip(plans).enumerate() {
let cursor = base_cursor.saturating_add((i as u64).saturating_add(1));
match (write, plan) {
(PreparedWrite::Node { kind, body, source_id, logical_id }, WritePlan::Node) => {
if let Some(logical_id) = logical_id {
tx.execute(
"UPDATE canonical_nodes SET superseded_at = ?1
WHERE logical_id = ?2 AND superseded_at IS NULL",
params![cursor, logical_id],
)?;
}
tx.execute(
"INSERT INTO canonical_nodes(write_cursor, kind, body, source_id, logical_id)
VALUES(?1, ?2, ?3, ?4, ?5)",
params![cursor, kind, body, source_id, logical_id],
)?;
tx.execute(
"INSERT INTO search_index(body, kind, write_cursor) VALUES(?1, ?2, ?3)",
params![body, kind, cursor],
)?;
if kind_is_vector_indexed(&tx, kind).unwrap_or(false) {
tx.execute(
"INSERT INTO _fathomdb_projection_state(kind, last_enqueued_cursor, updated_at)
VALUES(?1, ?2, 0)
ON CONFLICT(kind) DO UPDATE SET last_enqueued_cursor = excluded.last_enqueued_cursor",
params![kind, cursor],
)?;
} else {
record_projection_terminal(&tx, cursor, "up_to_date")?;
}
}
(PreparedWrite::Edge { kind, from, to, source_id, logical_id }, WritePlan::Edge) => {
if let Some(logical_id) = logical_id {
tx.execute(
"UPDATE canonical_edges SET superseded_at = ?1
WHERE logical_id = ?2 AND superseded_at IS NULL",
params![cursor, logical_id],
)?;
}
tx.execute(
"INSERT INTO canonical_edges(write_cursor, kind, from_id, to_id, source_id, logical_id)
VALUES(?1, ?2, ?3, ?4, ?5, ?6)",
params![cursor, kind, from, to, source_id, logical_id],
)?;
record_projection_terminal(&tx, cursor, "up_to_date")?;
}
(
PreparedWrite::AdminSchema { name, kind, schema_json, retention_json },
WritePlan::AdminSchema,
) => {
tx.execute(
"INSERT INTO operational_collections(
name, kind, schema_json, retention_json, format_version, created_at
) VALUES(?1, ?2, ?3, ?4, 1, 0)
ON CONFLICT(name) DO UPDATE SET
schema_json = excluded.schema_json,
retention_json = excluded.retention_json",
params![name, kind, schema_json, retention_json],
)?;
record_projection_terminal(&tx, cursor, "up_to_date")?;
}
(
PreparedWrite::OpStore { collection, record_key, schema_id, body },
WritePlan::AppendOnlyLog,
) => {
tx.execute(
"INSERT INTO operational_mutations(
collection_name, record_key, op_kind, payload_json, schema_id, write_cursor
) VALUES(?1, ?2, 'append', ?3, ?4, ?5)",
params![collection, record_key, body, schema_id, cursor],
)?;
record_projection_terminal(&tx, cursor, "up_to_date")?;
}
(
PreparedWrite::OpStore { collection, record_key, schema_id, body },
WritePlan::LatestState,
) => {
tx.execute(
"INSERT INTO operational_state(
collection_name, record_key, payload_json, schema_id, write_cursor
) VALUES(?1, ?2, ?3, ?4, ?5)
ON CONFLICT(collection_name, record_key) DO UPDATE SET
payload_json = excluded.payload_json,
schema_id = excluded.schema_id,
write_cursor = excluded.write_cursor",
params![collection, record_key, body, schema_id, cursor],
)?;
record_projection_terminal(&tx, cursor, "up_to_date")?;
}
_ => return Err(rusqlite::Error::InvalidQuery),
}
}
let dangling_edge_endpoints = {
let mut last_index: HashMap<&str, usize> = HashMap::new();
for (i, write) in batch.iter().enumerate() {
if let PreparedWrite::Edge { logical_id: Some(lid), .. } = write {
last_index.insert(lid.as_str(), i);
}
}
let mut probe = tx.prepare(
"SELECT 1 FROM canonical_nodes WHERE logical_id = ?1 AND superseded_at IS NULL LIMIT 1",
)?;
let mut count: u64 = 0;
for (i, write) in batch.iter().enumerate() {
if let PreparedWrite::Edge { from, to, logical_id, .. } = write {
if let Some(lid) = logical_id {
let superseded_in_batch =
last_index.get(lid.as_str()).is_some_and(|&last| last > i);
if superseded_in_batch {
continue;
}
}
for endpoint in [from, to] {
if !probe.exists(params![endpoint])? {
count = count.saturating_add(1);
}
}
}
}
count
};
enforce_provenance_retention(&tx, provenance_row_cap)?;
advance_projection_cursor(&tx)?;
tx.commit()?;
Ok(dangling_edge_endpoints)
}
fn load_next_cursor(connection: &Connection) -> u64 {
let nodes = max_cursor(connection, "canonical_nodes").unwrap_or(0);
let edges = max_cursor(connection, "canonical_edges").unwrap_or(0);
let mutations = max_cursor(connection, "operational_mutations").unwrap_or(0);
let state = max_cursor(connection, "operational_state").unwrap_or(0);
nodes.max(edges).max(mutations).max(state)
}
fn max_cursor(connection: &Connection, table: &str) -> rusqlite::Result<u64> {
let sql = format!("SELECT COALESCE(MAX(write_cursor), 0) FROM {table}");
connection.query_row(&sql, [], |row| row.get::<_, u64>(0))
}
fn sqlite_extended_code_name(err: &rusqlite::Error) -> Option<&'static str> {
let sqlite_error = err.sqlite_error()?;
let extended = sqlite_error.extended_code;
Some(match extended {
rusqlite::ffi::SQLITE_SCHEMA => "SQLITE_SCHEMA",
rusqlite::ffi::SQLITE_BUSY => "SQLITE_BUSY",
rusqlite::ffi::SQLITE_LOCKED => "SQLITE_LOCKED",
rusqlite::ffi::SQLITE_CORRUPT => "SQLITE_CORRUPT",
rusqlite::ffi::SQLITE_NOTADB => "SQLITE_NOTADB",
rusqlite::ffi::SQLITE_IOERR => "SQLITE_IOERR",
rusqlite::ffi::SQLITE_FULL => "SQLITE_FULL",
rusqlite::ffi::SQLITE_READONLY => "SQLITE_READONLY",
rusqlite::ffi::SQLITE_CONSTRAINT => "SQLITE_CONSTRAINT",
rusqlite::ffi::SQLITE_MISUSE => "SQLITE_MISUSE",
rusqlite::ffi::SQLITE_INTERRUPT => "SQLITE_INTERRUPT",
rusqlite::ffi::SQLITE_NOMEM => "SQLITE_NOMEM",
rusqlite::ffi::SQLITE_PERM => "SQLITE_PERM",
rusqlite::ffi::SQLITE_ABORT => "SQLITE_ABORT",
rusqlite::ffi::SQLITE_PROTOCOL => "SQLITE_PROTOCOL",
rusqlite::ffi::SQLITE_RANGE => "SQLITE_RANGE",
rusqlite::ffi::SQLITE_TOOBIG => "SQLITE_TOOBIG",
rusqlite::ffi::SQLITE_MISMATCH => "SQLITE_MISMATCH",
rusqlite::ffi::SQLITE_AUTH => "SQLITE_AUTH",
rusqlite::ffi::SQLITE_NOTFOUND => "SQLITE_NOTFOUND",
rusqlite::ffi::SQLITE_CANTOPEN => "SQLITE_CANTOPEN",
_ => "SQLITE_UNKNOWN",
})
}
fn sqlite_extended_code_name_from_int(extended: i32) -> &'static str {
match extended {
rusqlite::ffi::SQLITE_SCHEMA => "SQLITE_SCHEMA",
rusqlite::ffi::SQLITE_BUSY => "SQLITE_BUSY",
rusqlite::ffi::SQLITE_LOCKED => "SQLITE_LOCKED",
rusqlite::ffi::SQLITE_CORRUPT => "SQLITE_CORRUPT",
rusqlite::ffi::SQLITE_NOTADB => "SQLITE_NOTADB",
rusqlite::ffi::SQLITE_IOERR => "SQLITE_IOERR",
rusqlite::ffi::SQLITE_FULL => "SQLITE_FULL",
rusqlite::ffi::SQLITE_READONLY => "SQLITE_READONLY",
rusqlite::ffi::SQLITE_CONSTRAINT => "SQLITE_CONSTRAINT",
rusqlite::ffi::SQLITE_MISUSE => "SQLITE_MISUSE",
rusqlite::ffi::SQLITE_INTERRUPT => "SQLITE_INTERRUPT",
rusqlite::ffi::SQLITE_NOMEM => "SQLITE_NOMEM",
rusqlite::ffi::SQLITE_PERM => "SQLITE_PERM",
rusqlite::ffi::SQLITE_ABORT => "SQLITE_ABORT",
rusqlite::ffi::SQLITE_PROTOCOL => "SQLITE_PROTOCOL",
rusqlite::ffi::SQLITE_RANGE => "SQLITE_RANGE",
rusqlite::ffi::SQLITE_TOOBIG => "SQLITE_TOOBIG",
rusqlite::ffi::SQLITE_MISMATCH => "SQLITE_MISMATCH",
rusqlite::ffi::SQLITE_AUTH => "SQLITE_AUTH",
rusqlite::ffi::SQLITE_NOTFOUND => "SQLITE_NOTFOUND",
rusqlite::ffi::SQLITE_CANTOPEN => "SQLITE_CANTOPEN",
_ => "SQLITE_UNKNOWN",
}
}
fn map_open_sqlite_error(err: rusqlite::Error, stage: OpenStage) -> EngineOpenError {
let Some(sqlite_error) = err.sqlite_error() else {
return EngineOpenError::Io { message: "could not open database".to_string() };
};
match sqlite_error.extended_code {
rusqlite::ffi::SQLITE_CORRUPT | rusqlite::ffi::SQLITE_NOTADB => {
EngineOpenError::Corruption(CorruptionDetail {
kind: match stage {
OpenStage::WalReplay => CorruptionKind::WalReplayFailure,
OpenStage::HeaderProbe => CorruptionKind::HeaderMalformed,
OpenStage::SchemaProbe => CorruptionKind::SchemaInconsistent,
OpenStage::EmbedderIdentity => CorruptionKind::EmbedderIdentityDrift,
},
stage,
locator: CorruptionLocator::OpaqueSqliteError {
sqlite_extended_code: sqlite_error.extended_code,
},
recovery_hint: RecoveryHint {
code: match stage {
OpenStage::WalReplay => "E_CORRUPT_WAL_REPLAY",
OpenStage::HeaderProbe => "E_CORRUPT_HEADER",
OpenStage::SchemaProbe => "E_CORRUPT_SCHEMA",
OpenStage::EmbedderIdentity => "E_CORRUPT_EMBEDDER_IDENTITY",
},
doc_anchor: match stage {
OpenStage::WalReplay => "design/recovery.md#wal-replay-failures",
OpenStage::HeaderProbe => "design/recovery.md#header-malformed",
OpenStage::SchemaProbe => "design/recovery.md#schema-inconsistent",
OpenStage::EmbedderIdentity => "design/recovery.md#embedder-identity-drift",
},
},
})
}
_ => EngineOpenError::Io { message: "could not open database".to_string() },
}
}
fn emit_open_error_event(subscriber: &Arc<dyn lifecycle::Subscriber>, err: &EngineOpenError) {
if let EngineOpenError::Corruption(detail) = err {
let code = match detail.locator {
CorruptionLocator::OpaqueSqliteError { sqlite_extended_code } => {
Some(sqlite_extended_code_name_from_int(sqlite_extended_code))
}
_ => None,
};
let event = lifecycle::Event {
phase: lifecycle::Phase::Failed,
source: lifecycle::EventSource::SqliteInternal,
category: lifecycle::EventCategory::Corruption,
code,
};
subscriber.on_event(&event);
}
}
#[allow(clippy::vec_box)]
fn install_profile_callback(
connection: &Connection,
subscribers: &Arc<lifecycle::SubscriberRegistry>,
profiling_enabled: &Arc<AtomicBool>,
slow_threshold_ms: &Arc<AtomicU64>,
contexts: &mut Vec<Box<ProfileContext>>,
) {
let mut ctx = Box::new(ProfileContext {
subscribers: Arc::clone(subscribers),
profiling_enabled: Arc::clone(profiling_enabled),
slow_threshold_ms: Arc::clone(slow_threshold_ms),
});
let ctx_ptr: *mut ProfileContext = &mut *ctx;
unsafe {
rusqlite::ffi::sqlite3_profile(
connection.handle(),
Some(profile_callback_trampoline),
ctx_ptr.cast::<std::ffi::c_void>(),
);
}
contexts.push(ctx);
}
fn uninstall_profile_callback(connection: &Connection) {
unsafe {
rusqlite::ffi::sqlite3_profile(connection.handle(), None, std::ptr::null_mut());
}
}
fn apply_perf_experiment_writer_pragmas(connection: &Connection) {
if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
return;
}
let raw = match std::env::var("FATHOMDB_PERF_WRITER_PRAGMAS") {
Ok(s) if !s.is_empty() => s,
_ => return,
};
for entry in raw.split(',') {
let entry = entry.trim();
if entry.is_empty() {
continue;
}
let (name, value) = match entry.split_once('=') {
Some((n, v)) => (n.trim(), v.trim()),
None => {
eprintln!("perf-experiment: bad writer pragma entry (expect name=value): {entry}");
continue;
}
};
if name.is_empty() {
eprintln!("perf-experiment: empty pragma name in writer entry: {entry}");
continue;
}
match connection.pragma_update(None, name, value) {
Ok(()) => {
eprintln!(
"perf-experiment: applied PRAGMA {name}={value} on writer (pre-migration)"
);
}
Err(err) => {
eprintln!("perf-experiment: writer PRAGMA {name}={value} failed: {err}");
}
}
}
}
fn apply_perf_experiment_reader_pragmas(connection: &Connection) {
if std::env::var_os("FATHOMDB_PERF_EXPERIMENTS").is_none() {
return;
}
let raw = match std::env::var("FATHOMDB_PERF_READER_PRAGMAS") {
Ok(s) if !s.is_empty() => s,
_ => return,
};
for entry in raw.split(',') {
let entry = entry.trim();
if entry.is_empty() {
continue;
}
let (name, value) = match entry.split_once('=') {
Some((n, v)) => (n.trim(), v.trim()),
None => {
eprintln!("perf-experiment: bad pragma entry (expect name=value): {entry}");
continue;
}
};
if name.is_empty() {
eprintln!("perf-experiment: empty pragma name in entry: {entry}");
continue;
}
match connection.pragma_update(None, name, value) {
Ok(()) => {
eprintln!("perf-experiment: applied PRAGMA {name}={value} on reader");
}
Err(err) => {
eprintln!("perf-experiment: PRAGMA {name}={value} failed: {err}");
}
}
}
}
fn configure_reader_lookaside(connection: &Connection) -> std::os::raw::c_int {
unsafe {
rusqlite::ffi::sqlite3_db_config(
connection.handle(),
rusqlite::ffi::SQLITE_DBCONFIG_LOOKASIDE,
std::ptr::null_mut::<std::ffi::c_void>(),
READER_LOOKASIDE_SLOT_SIZE,
READER_LOOKASIDE_SLOT_COUNT,
)
}
}
#[cfg(debug_assertions)]
fn read_lookaside_used_hiwtr(connection: &Connection) -> std::os::raw::c_int {
let mut current: std::os::raw::c_int = 0;
let mut hiwtr: std::os::raw::c_int = 0;
unsafe {
rusqlite::ffi::sqlite3_db_status(
connection.handle(),
rusqlite::ffi::SQLITE_DBSTATUS_LOOKASIDE_USED,
&mut current,
&mut hiwtr,
0,
);
}
hiwtr
}
#[cfg(debug_assertions)]
fn read_cache_status(
connection: &Connection,
) -> (std::os::raw::c_int, std::os::raw::c_int, std::os::raw::c_int) {
let mut hit_current: std::os::raw::c_int = 0;
let mut hit_hiwtr: std::os::raw::c_int = 0;
let mut miss_current: std::os::raw::c_int = 0;
let mut miss_hiwtr: std::os::raw::c_int = 0;
let mut used_current: std::os::raw::c_int = 0;
let mut used_hiwtr: std::os::raw::c_int = 0;
unsafe {
rusqlite::ffi::sqlite3_db_status(
connection.handle(),
rusqlite::ffi::SQLITE_DBSTATUS_CACHE_HIT,
&mut hit_current,
&mut hit_hiwtr,
0,
);
rusqlite::ffi::sqlite3_db_status(
connection.handle(),
rusqlite::ffi::SQLITE_DBSTATUS_CACHE_MISS,
&mut miss_current,
&mut miss_hiwtr,
0,
);
rusqlite::ffi::sqlite3_db_status(
connection.handle(),
rusqlite::ffi::SQLITE_DBSTATUS_CACHE_USED,
&mut used_current,
&mut used_hiwtr,
0,
);
}
(hit_current, miss_current, used_current)
}
unsafe extern "C" fn profile_callback_trampoline(
user_data: *mut std::ffi::c_void,
sql: *const std::os::raw::c_char,
nanoseconds: u64,
) {
if user_data.is_null() || sql.is_null() {
return;
}
let ctx = unsafe { &*(user_data.cast::<ProfileContext>()) };
let sql_text = match unsafe { std::ffi::CStr::from_ptr(sql) }.to_str() {
Ok(s) => s,
Err(_) => return,
};
let wall_clock_ms = nanoseconds / 1_000_000;
if ctx.profiling_enabled.load(Ordering::Relaxed) {
let record = lifecycle::ProfileRecord {
wall_clock_ms,
step_count: 0,
cache_delta: 0,
};
ctx.subscribers.dispatch_profile(&record);
}
let threshold = ctx.slow_threshold_ms.load(Ordering::Relaxed);
if wall_clock_ms > threshold {
let signal = lifecycle::SlowStatement { statement: sql_text.to_string(), wall_clock_ms };
ctx.subscribers.dispatch_slow_statement(&signal);
}
}
#[cfg(test)]
mod tests {
use super::{resolve_source_type, Engine, PreparedWrite, KIND_TO_SOURCE_TYPE_CASE_SQL};
use rusqlite::Connection;
use tempfile::TempDir;
#[test]
fn resolve_source_type_drift_check() {
let kinds = ["email", "article", "paper", "meeting", "note", "todo", "doc"];
let want: &[(&str, &str)] = &[
("email", "email"),
("article", "article"),
("paper", "paper"),
("meeting", "meeting"),
("note", "note"),
("todo", "todo"),
("doc", "article"),
];
for (kind, expected) in want {
let got = resolve_source_type(kind).unwrap_or_else(|_| {
panic!("resolve_source_type({kind}) returned Err; want Ok({expected})")
});
assert_eq!(got, *expected, "Rust helper drift for kind={kind}");
}
assert!(
resolve_source_type("banana").is_err(),
"unknown kind must surface as writer error"
);
let conn = Connection::open_in_memory().expect("in-memory sqlite");
conn.execute_batch("CREATE TABLE s(kind TEXT NOT NULL)").expect("create s");
for kind in &kinds {
conn.execute("INSERT INTO s(kind) VALUES (?1)", [kind]).expect("insert kind");
}
let sql = format!("SELECT s.kind, {KIND_TO_SOURCE_TYPE_CASE_SQL} FROM s");
let mut stmt = conn.prepare(&sql).expect("prepare CASE");
let rows: Vec<(String, String)> = stmt
.query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)))
.expect("query")
.map(|r| r.expect("row"))
.collect();
assert_eq!(rows.len(), kinds.len(), "row count drift");
for (kind, sql_result) in &rows {
let rust_result = resolve_source_type(kind).expect("known kind");
assert_eq!(
sql_result, rust_result,
"SQL CASE vs Rust helper drift for kind={kind}: SQL={sql_result}, Rust={rust_result}"
);
}
}
#[test]
fn write_advances_cursor() {
let dir = TempDir::new().unwrap();
let opened = Engine::open(dir.path().join("rewrite.sqlite")).expect("engine should open");
let receipt = opened
.engine
.write(&[PreparedWrite::Node {
kind: "doc".to_string(),
body: "hello".to_string(),
source_id: None,
logical_id: None,
}])
.expect("write should succeed");
assert_eq!(receipt.cursor, 1);
}
}