use std::path::Path;
use std::sync::Arc;
use fathomdb_embedder_api::{Embedder, EmbedderError, EmbedderIdentity, Vector};
use fathomdb_engine::{Engine, PreparedWrite, RowKind};
use rusqlite::Connection;
use tempfile::TempDir;
#[derive(Debug)]
struct DeterministicEmbedder {
identity: EmbedderIdentity,
dim: u32,
}
impl DeterministicEmbedder {
fn new(dim: u32) -> Self {
Self { identity: EmbedderIdentity::new("deterministic", "exp-s-substrate", dim), dim }
}
}
impl Embedder for DeterministicEmbedder {
fn identity(&self) -> EmbedderIdentity {
self.identity.clone()
}
fn embed(&self, text: &str) -> Result<Vector, EmbedderError> {
let dim = self.dim as usize;
let mut v = vec![0.0_f32; dim];
let mut h: u64 = 0xcbf29ce4_84222325;
for &b in text.as_bytes() {
h ^= b as u64;
h = h.wrapping_mul(0x0100_0000_01b3);
}
for k in 0..6 {
let coord = ((h >> (k * 8)) as usize) % dim;
let sign = if (h >> (k * 8 + 7)) & 1 == 0 { 1.0 } else { -1.0 };
v[coord] += sign * 0.5_f32;
}
let norm: f32 = v.iter().map(|x| x * x).sum::<f32>().sqrt().max(1e-6);
for x in &mut v {
*x /= norm;
}
Ok(v)
}
}
fn fresh_engine(name: &str) -> (TempDir, Engine) {
let dir = TempDir::new().expect("tempdir");
let path = dir.path().join(format!("{name}.sqlite"));
let opened =
Engine::open_with_embedder_for_test(&path, Arc::new(DeterministicEmbedder::new(768)))
.expect("open");
opened.engine.configure_vector_kind_for_test("doc").expect("configure vector kind");
(dir, opened.engine)
}
fn write_fixture(engine: &Engine) {
let leaves: Vec<PreparedWrite> = (0..6)
.map(|i| PreparedWrite::Node {
kind: "doc".to_string(),
body: format!("leaf body {i} alpha bravo charlie token-{i}"),
source_id: fathomdb_engine::SourceId::new(format!("leaf-{i}")).expect("test source id"),
logical_id: None,
state: fathomdb_engine::InitialState::Active,
reason: None,
valid_from: None,
valid_until: None,
})
.collect();
engine.write(&leaves).expect("write leaf batch");
for i in 0..3 {
engine
.write_canonical_row_with_kind_for_test(
"doc",
&format!("coverage summary {i} delta echo"),
RowKind::Coverage,
)
.expect("write coverage row");
}
for i in 0..2 {
engine
.write_canonical_row_with_kind_for_test(
"doc",
&format!("graph structural {i} foxtrot golf"),
RowKind::Graph,
)
.expect("write graph row");
}
}
#[test]
fn r_sub_1_row_kinds_coexist_and_are_queryable_by_row_kind() {
let (_dir, engine) = fresh_engine("row_kinds");
write_fixture(&engine);
engine.drain(15_000).expect("drain");
let leaf = engine.canonical_rows_with_row_kind_for_test(RowKind::Leaf).expect("leaf rows");
let coverage =
engine.canonical_rows_with_row_kind_for_test(RowKind::Coverage).expect("coverage rows");
let graph = engine.canonical_rows_with_row_kind_for_test(RowKind::Graph).expect("graph rows");
assert_eq!(leaf.len(), 6, "expected 6 leaf rows, got {leaf:?}");
assert_eq!(coverage.len(), 3, "expected 3 coverage rows, got {coverage:?}");
assert_eq!(graph.len(), 2, "expected 2 graph rows, got {graph:?}");
for c in &coverage {
assert!(!leaf.contains(c), "coverage cursor {c} leaked into leaf set");
assert!(!graph.contains(c), "coverage cursor {c} leaked into graph set");
}
let present = [leaf.is_empty(), coverage.is_empty(), graph.is_empty()]
.iter()
.filter(|empty| !**empty)
.count();
assert!(present >= 2, "expected >=2 distinct row_kinds populated, got {present}");
}
#[test]
fn d2_index_targets_differ_by_row_kind() {
let (_dir, engine) = fresh_engine("index_targets");
write_fixture(&engine);
engine.drain(15_000).expect("drain");
let db_path = engine.path().to_path_buf();
engine.close().expect("close");
let conn = open_readonly(&db_path);
let fts_count: i64 =
conn.query_row("SELECT COUNT(*) FROM search_index", [], |r| r.get(0)).expect("fts count");
assert_eq!(fts_count, 11, "all 11 rows (leaf+coverage+graph) must be FTS-indexed");
let vec_count: i64 =
conn.query_row("SELECT COUNT(*) FROM vector_default", [], |r| r.get(0)).expect("vec count");
assert_eq!(vec_count, 9, "only leaf+coverage rows (9) must be vector-indexed, not graph");
}
fn open_readonly(path: &Path) -> Connection {
Connection::open_with_flags(
path,
rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
)
.expect("open read-only sqlite")
}
fn serialize_index_state(path: &Path) -> Vec<u8> {
let conn = open_readonly(path);
let mut out: Vec<u8> = Vec::new();
out.extend_from_slice(b"# search_index\n");
let mut stmt = conn
.prepare("SELECT write_cursor, kind, body FROM search_index ORDER BY write_cursor")
.expect("prep fts");
let rows = stmt
.query_map([], |r| {
Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?, r.get::<_, String>(2)?))
})
.expect("fts rows");
for row in rows {
let (wc, kind, body) = row.expect("fts row");
out.extend_from_slice(format!("{wc}|{kind}|{body}\n").as_bytes());
}
out.extend_from_slice(b"# search_index_v2\n");
let mut stmt = conn
.prepare(
"SELECT write_cursor, kind, body, status FROM search_index_v2 ORDER BY write_cursor",
)
.expect("prep fts v2");
let rows = stmt
.query_map([], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, String>(3)?,
))
})
.expect("fts v2 rows");
for row in rows {
let (wc, kind, body, status) = row.expect("fts v2 row");
out.extend_from_slice(format!("{wc}|{kind}|{body}|{status}\n").as_bytes());
}
out.extend_from_slice(b"# vector_default\n");
let mut stmt = conn
.prepare(
"SELECT rowid, source_type, kind, embedding_bin FROM vector_default ORDER BY rowid",
)
.expect("prep vec");
let rows = stmt
.query_map([], |r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, String>(1)?,
r.get::<_, String>(2)?,
r.get::<_, Vec<u8>>(3)?,
))
})
.expect("vec rows");
for row in rows {
let (rowid, st, kind, blob) = row.expect("vec row");
out.extend_from_slice(format!("{rowid}|{st}|{kind}|").as_bytes());
out.extend_from_slice(&blob);
out.push(b'\n');
}
out.extend_from_slice(b"# _fathomdb_vector_rows\n");
let mut stmt = conn
.prepare("SELECT rowid, kind, write_cursor FROM _fathomdb_vector_rows ORDER BY rowid")
.expect("prep vecrows");
let rows = stmt
.query_map([], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?, r.get::<_, i64>(2)?)))
.expect("vecrows rows");
for row in rows {
let (rowid, kind, wc) = row.expect("vecrows row");
out.extend_from_slice(format!("{rowid}|{kind}|{wc}\n").as_bytes());
}
out.extend_from_slice(b"# canonical_nodes.row_kind\n");
let mut stmt = conn
.prepare("SELECT write_cursor, kind, row_kind FROM canonical_nodes ORDER BY write_cursor")
.expect("prep rowkind");
let rows = stmt
.query_map([], |r| {
Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?, r.get::<_, String>(2)?))
})
.expect("rowkind rows");
for row in rows {
let (wc, kind, rk) = row.expect("rowkind row");
out.extend_from_slice(format!("{wc}|{kind}|{rk}\n").as_bytes());
}
out.extend_from_slice(b"# _fathomdb_projection_terminal\n");
let mut stmt = conn
.prepare(
"SELECT write_cursor, state FROM _fathomdb_projection_terminal ORDER BY write_cursor",
)
.expect("prep terminal");
let rows = stmt
.query_map([], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, String>(1)?)))
.expect("terminal rows");
for row in rows {
let (wc, state) = row.expect("terminal row");
out.extend_from_slice(format!("{wc}|{state}\n").as_bytes());
}
out
}
#[test]
fn r_sub_2_incremental_multi_index_write_is_deterministic() {
let (dir_a, engine_a) = fresh_engine("determinism_a");
write_fixture(&engine_a);
engine_a.drain(30_000).expect("drain A to quiescence");
let path_a = engine_a.path().to_path_buf();
engine_a.close().expect("close A");
let state_a = serialize_index_state(&path_a);
let (dir_b, engine_b) = fresh_engine("determinism_b");
write_fixture(&engine_b);
engine_b.drain(30_000).expect("drain B to quiescence");
let path_b = engine_b.path().to_path_buf();
engine_b.close().expect("close B");
let state_b = serialize_index_state(&path_b);
assert!(!state_a.is_empty(), "serialized index state must be non-empty");
assert_eq!(
state_a.len(),
state_b.len(),
"index-state byte length differs across runs ({} vs {})",
state_a.len(),
state_b.len()
);
assert!(
state_a == state_b,
"incremental multi-index write is NOT deterministic: index state differs byte-for-byte \
across two fresh DBs on the same CPU"
);
drop(dir_a);
drop(dir_b);
}