use super::*;
#[test]
fn record_in_normal_state_takes_sync_path() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert_eq!(
db.write_router.state(),
crate::engine::write_router::RouterState::Normal
);
let rid = db
.record(
"sync path test",
"episodic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 1,
"sync-path write must land in memories immediately"
);
}
#[test]
fn record_in_queueing_state_routes_to_oplog_does_not_touch_memories() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.write_router.switch_to_queueing();
assert_eq!(
db.write_router.state(),
crate::engine::write_router::RouterState::Queueing
);
let mem_before: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |row| row.get(0))
.unwrap();
let oplog_before: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE applied = 0 AND op_type = 'record'",
[],
|row| row.get(0),
)
.unwrap();
let rid = db
.record(
"queued path test",
"episodic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let mem_after: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |row| row.get(0))
.unwrap();
assert_eq!(
mem_after, mem_before,
"queued path must NOT write to memories table during reembed cutover \
(brainstorm-2/3 invariant 1: queued-after-barrier writes are replayed \
by post-swap materializer, not committed to old generation)"
);
let oplog_after: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE applied = 0 AND op_type = 'record' \
AND target_rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(
oplog_after,
oplog_before + 1,
"queued path must write the record op to oplog with applied=0"
);
let applied_gen: Option<i64> = db
.conn()
.query_row(
"SELECT applied_generation FROM oplog WHERE target_rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert!(
applied_gen.is_none(),
"applied_generation must be NULL for queued ops; got Some({applied_gen:?})"
);
db.write_router.switch_to_normal();
}
#[test]
fn record_guard_drops_inflight_counter_panic_safe_via_raii() {
let db = std::sync::Arc::new(YantrikDB::new(":memory:", 8).unwrap());
assert_eq!(db.write_router.inflight(), 0);
let db_panic = std::sync::Arc::clone(&db);
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _guard = db_panic
.write_router
.try_enter_sync_writer()
.expect("Normal state must yield guard");
assert_eq!(db_panic.write_router.inflight(), 1);
panic!("simulated mid-write panic");
}));
assert!(result.is_err(), "panic must propagate up");
assert_eq!(
db.write_router.inflight(),
0,
"panic-safe inflight counter (RAII Drop) — required for reembed cutover correctness"
);
}
mod mode_test_embedders {
use crate::types::Embedder;
pub struct FakeEmbedder {
pub dim: usize,
pub fp: Option<String>,
pub name: Option<String>,
pub sentinel: f32,
}
impl Embedder for FakeEmbedder {
fn embed(
&self,
_text: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let mut v = vec![0.0_f32; self.dim];
if !v.is_empty() {
v[0] = self.sentinel;
}
Ok(v)
}
fn dim(&self) -> usize {
self.dim
}
fn fingerprint(&self) -> Option<String> {
self.fp.clone()
}
fn name(&self) -> Option<String> {
self.name.clone()
}
}
}
#[test]
fn set_embedder_test_1_same_dim_different_digest_on_populated_db_rejected() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:embedder_A".to_string()),
name: Some("embedder_A".to_string()),
sentinel: 1.0,
}))
.unwrap();
let _ = db
.record(
"first memory",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let err = db
.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:embedder_B".to_string()),
name: Some("embedder_B".to_string()),
sentinel: 2.0,
}))
.unwrap_err();
assert!(
matches!(
err,
crate::error::YantrikDbError::ChangeEmbedderDigestRequiresReembed { .. }
),
"same-dim-different-digest on Known-provenance populated DB must \
return ChangeEmbedderDigestRequiresReembed (silent-corruption \
prevention invariant from brainstorm-3); got {err:?}"
);
}
#[test]
fn set_embedder_test_2_different_dim_on_populated_db_rejected() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:fp64".to_string()),
name: None,
sentinel: 1.0,
}))
.unwrap();
let _ = db
.record(
"m",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let err = db
.set_embedder(Box::new(FakeEmbedder {
dim: 128,
fp: Some("sha256:fp128".to_string()),
name: None,
sentinel: 2.0,
}))
.unwrap_err();
assert!(
matches!(
err,
crate::error::YantrikDbError::ChangeEmbedderDimensionRequiresReembed { .. }
),
"dim change on populated DB must return \
ChangeEmbedderDimensionRequiresReembed; got {err:?}"
);
}
#[test]
fn set_embedder_test_3_empty_db_with_fingerprint_upgrades_provenance_to_known() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 128).unwrap();
assert!(matches!(
db.search_state.load().index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
));
db.set_embedder(Box::new(FakeEmbedder {
dim: 128,
fp: Some("sha256:initial".to_string()),
name: Some("initial".to_string()),
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
match &s.index_embedding {
crate::engine::reembed::EmbeddingProvenance::Known { name, digest, dim } => {
assert_eq!(name.as_deref(), Some("initial"));
assert_eq!(digest, "sha256:initial");
assert_eq!(*dim, 128);
}
other => panic!("expected Known provenance after attach on empty DB, got {other:?}"),
}
}
#[test]
fn set_embedder_test_4_empty_db_no_fingerprint_stays_external_or_unknown() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: None,
name: None,
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
assert!(
matches!(
s.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { dim: 64 }
),
"no-fingerprint embedder on empty DB must keep ExternalOrUnknown provenance, \
got {:?}",
s.index_embedding
);
assert!(s.has_runtime_embedder());
}
#[test]
fn set_embedder_test_5_same_digest_replacement_does_not_bump_generation() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:same".to_string()),
name: Some("same".to_string()),
sentinel: 1.0,
}))
.unwrap();
let gen_before = db.search_state.load().generation;
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:same".to_string()),
name: Some("same".to_string()),
sentinel: 2.0,
}))
.unwrap();
let gen_after = db.search_state.load().generation;
assert_eq!(
gen_after, gen_before,
"same-digest replacement must NOT bump generation (no coherent-bundle change)"
);
let v = db.embed("anything").unwrap();
assert!(
(v[0] - 2.0).abs() < 1e-6,
"runtime Arc must have been replaced; expected sentinel 2.0, got {}",
v[0]
);
}
#[test]
fn set_embedder_test_6_external_or_unknown_compat_attach_does_not_claim_provenance() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 128).unwrap();
let _ = db
.record(
"external vec",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 128),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
assert!(matches!(
db.search_state.load().index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
));
db.set_embedder(Box::new(FakeEmbedder {
dim: 128,
fp: Some("sha256:attached".to_string()),
name: Some("attached".to_string()),
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
assert!(
matches!(
s.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
),
"compat-attach must NOT upgrade ExternalOrUnknown provenance to Known; \
existing vectors weren't built with this embedder. Got {:?}",
s.index_embedding
);
assert!(s.has_runtime_embedder());
assert_eq!(
s.runtime_embedder_digest.as_deref(),
Some("sha256:attached")
);
}
#[test]
fn set_embedder_test_7_has_embedder_derives_from_search_state() {
let mut db = YantrikDB::new(":memory:", 384).unwrap();
assert!(!db.has_embedder(), "fresh engine: no embedder");
use mode_test_embedders::FakeEmbedder;
db.set_embedder(Box::new(FakeEmbedder {
dim: 384,
fp: Some("sha256:x".to_string()),
name: None,
sentinel: 0.42,
}))
.unwrap();
assert!(db.has_embedder());
let v = db.embed("anything").unwrap();
assert!((v[0] - 0.42).abs() < 1e-6);
}
#[test]
fn record_text_revalidates_generation_and_retries_after_swap() {
use std::sync::mpsc::channel;
use std::sync::{Arc, Mutex};
use std::time::Duration;
let (started_tx, started_rx) = channel::<()>();
let (release_tx, release_rx) = channel::<()>();
let call_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
struct SharedBlocking {
dim: usize,
fp: Option<String>,
name: Option<String>,
sentinel: f32,
started_tx: Mutex<Option<std::sync::mpsc::Sender<()>>>,
release_rx: Mutex<Option<std::sync::mpsc::Receiver<()>>>,
call_count: Arc<std::sync::atomic::AtomicUsize>,
}
impl crate::types::Embedder for SharedBlocking {
fn embed(
&self,
_text: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let n = self
.call_count
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if n == 0 {
if let Some(tx) = self.started_tx.lock().unwrap().take() {
let _ = tx.send(());
}
if let Some(rx) = self.release_rx.lock().unwrap().take() {
let _ = rx.recv();
}
}
let mut v = vec![0.0_f32; self.dim];
if !v.is_empty() {
v[0] = self.sentinel;
}
Ok(v)
}
fn dim(&self) -> usize {
self.dim
}
fn fingerprint(&self) -> Option<String> {
self.fp.clone()
}
fn name(&self) -> Option<String> {
self.name.clone()
}
}
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(SharedBlocking {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: Some("blocking-initial".to_string()),
sentinel: 0.42,
started_tx: Mutex::new(Some(started_tx)),
release_rx: Mutex::new(Some(release_rx)),
call_count: Arc::clone(&call_count),
}))
.unwrap();
let gen_before = db.search_state.load().generation;
let arc_db = Arc::new(db);
let worker_db = Arc::clone(&arc_db);
let worker = std::thread::spawn(move || {
worker_db
.record_text(
"hello",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap()
});
started_rx
.recv_timeout(Duration::from_secs(5))
.expect("embed should have started within 5s");
let old_state = arc_db.search_state.load_full();
let new_state = crate::engine::reembed::SearchState {
index_embedding: crate::engine::reembed::EmbeddingProvenance::Known {
name: Some("blocking-rotated".to_string()),
digest: "sha256:rotated".to_string(),
dim: 64,
},
embedder: old_state.embedder.clone(),
runtime_embedder_name: Some("blocking-rotated".to_string()),
runtime_embedder_digest: Some("sha256:rotated".to_string()),
generation: old_state.generation + 1,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: Arc::clone(&old_state.vec_index),
};
arc_db.search_state.store(Arc::new(new_state));
release_tx.send(()).unwrap();
let rid = worker.join().expect("record_text must complete");
let n_calls = call_count.load(std::sync::atomic::Ordering::SeqCst);
assert!(
n_calls >= 2,
"record_text must re-embed after SearchState swap; got {n_calls} calls"
);
assert!(!rid.is_empty(), "record_text returns a valid rid");
let gen_after = arc_db.search_state.load().generation;
assert!(
gen_after > gen_before,
"test must observe a generation advance: before={gen_before} after={gen_after}"
);
let conn = arc_db.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 1, "the retried record_text must be durably stored");
}
#[test]
fn record_text_routes_to_queued_when_router_is_queueing() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: Some("initial".to_string()),
sentinel: 0.5,
}))
.unwrap();
db.write_router.switch_to_queueing();
let rid = db
.record_text(
"hello-queued",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(!rid.is_empty());
let conn = db.conn();
let memories_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
memories_count, 0,
"queued path must NOT write to memories table"
);
let oplog_row: (i64, Option<String>) = conn
.query_row(
"SELECT applied, embedding_model FROM oplog WHERE target_rid = ?1 AND op_type = 'record'",
[&rid],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(oplog_row.0, 0, "queued op must be applied=0");
assert_eq!(
oplog_row.1.as_deref(),
Some("initial"),
"queued op carries the current runtime embedder name for post-swap re-encode"
);
}
#[test]
fn log_op_stamps_applied_generation_from_active_search_state() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial_generation: i64 = db.search_state.load().generation as i64;
let op_id = db
.log_op("test_event", None, &serde_json::json!({"x": 1}), None)
.unwrap();
let applied_generation: Option<i64> = {
let conn = db.conn();
conn.query_row(
"SELECT applied_generation FROM oplog WHERE op_id = ?1",
[&op_id],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(
applied_generation,
Some(initial_generation),
"log_op must stamp applied_generation with the current SearchState generation"
);
let old_state = db.search_state.load_full();
let bumped = crate::engine::reembed::SearchState {
index_embedding: old_state.index_embedding.clone(),
embedder: old_state.embedder.clone(),
runtime_embedder_name: old_state.runtime_embedder_name.clone(),
runtime_embedder_digest: old_state.runtime_embedder_digest.clone(),
generation: old_state.generation + 1,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&old_state.vec_index),
};
db.search_state.store(std::sync::Arc::new(bumped));
let op_id2 = db
.log_op(
"test_event_after_bump",
None,
&serde_json::json!({"x": 2}),
None,
)
.unwrap();
let applied_generation2: Option<i64> = {
let conn = db.conn();
conn.query_row(
"SELECT applied_generation FROM oplog WHERE op_id = ?1",
[&op_id2],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(
applied_generation2,
Some(initial_generation + 1),
"log_op picks up the new generation after search_state.store"
);
}
#[test]
fn set_embedder_test_8_atomic_publication_no_partial_state() {
use mode_test_embedders::FakeEmbedder;
use std::sync::Arc;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: None,
sentinel: 1.0,
}))
.unwrap();
let state = db.search_state.load_full();
assert_eq!(
state.embedder.is_some(),
state.runtime_embedder_digest.is_some(),
"embedder Some <=> digest Some must hold for any consistent snapshot"
);
for sentinel in [2.0_f32, 3.0, 4.0, 5.0] {
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: None,
sentinel,
}))
.unwrap();
let s = db.search_state.load_full();
assert_eq!(
s.embedder.is_some(),
s.runtime_embedder_digest.is_some(),
"consistency must hold across replacements (no partial state)"
);
let _arc_held: Arc<crate::engine::reembed::SearchState> = s;
}
}
#[test]
fn search_state_initial_on_fresh_engine() {
let db = YantrikDB::new(":memory:", 384).unwrap();
let state = db.search_state.load_full();
assert_eq!(
state.dim(),
384,
"initial dim must match constructor parameter"
);
assert!(matches!(
state.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { dim: 384 }
));
assert!(
!state.has_runtime_embedder(),
"fresh engine must have no runtime embedder until set_embedder*"
);
assert_eq!(state.generation, 0);
assert_eq!(state.covers_through_seq, 0);
}
#[test]
fn schema_v27_fresh_install_has_reembed_surfaces() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let memories_cols = table_columns(&conn, "memories");
for required in ["embedding_new", "embedding_new_model"] {
assert!(
memories_cols.iter().any(|c| c == required),
"v27: fresh-install memories table missing column {required}, got: {memories_cols:?}"
);
}
let oplog_cols = table_columns(&conn, "oplog");
assert!(
oplog_cols.iter().any(|c| c == "embedding_model"),
"v27: fresh-install oplog table missing column embedding_model, got: {oplog_cols:?}"
);
assert!(
oplog_cols.iter().any(|c| c == "applied_generation"),
"v27: fresh-install oplog table missing column applied_generation \
(brainstorm-2 correction \u{2014} per-generation application tracking \
replaces boolean `applied` as truth), got: {oplog_cols:?}"
);
let events_cols = table_columns(&conn, "reembed_events");
for required in ["generation", "phase", "timestamp", "payload_json"] {
assert!(
events_cols.iter().any(|c| c == required),
"v27: fresh-install reembed_events missing column {required}, got: {events_cols:?}"
);
}
for required_idx in [
"idx_reembed_events_generation",
"idx_oplog_applied_generation",
] {
assert!(
index_exists(&conn, required_idx),
"v27: fresh-install missing index {required_idx}"
);
}
}
#[test]
fn schema_v27_migration_from_v26_is_additive_only() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000c027";
let planted_embedding = vec![0_u8; 8 * std::mem::size_of::<f32>()];
{
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
conn.execute(
"INSERT INTO memories (rid, type, text, embedding, created_at, updated_at, last_access, source) \
VALUES (?1, 'episodic', 'planted under v27 schema', ?2, 0.0, 0.0, 0.0, 'user')",
params![planted_rid, planted_embedding],
)
.unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '26')",
[],
)
.unwrap();
conn.execute("DROP TABLE IF EXISTS reembed_events", [])
.unwrap();
conn.execute("DROP INDEX IF EXISTS idx_reembed_events_generation", [])
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v27 migration must run cleanly against a rewound-meta v26 DB");
let conn = db.conn();
let (preserved_text, embedding_new): (String, Option<Vec<u8>>) = conn
.query_row(
"SELECT text, embedding_new FROM memories WHERE rid = ?1",
params![planted_rid],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<Vec<u8>>>(1)?)),
)
.unwrap();
assert_eq!(
preserved_text, "planted under v27 schema",
"v27 migration must NOT mutate existing memory data"
);
assert!(
embedding_new.is_none(),
"v27 migration must leave embedding_new as NULL on pre-existing rows"
);
assert!(
index_exists(&conn, "idx_reembed_events_generation"),
"v27 migration must recreate idx_reembed_events_generation"
);
conn.execute(
"INSERT INTO reembed_events (generation, phase, timestamp, payload_json) \
VALUES (?1, ?2, ?3, ?4)",
params![1_i64, "Probing", 0.0_f64, "{}"],
)
.unwrap();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE generation = 1",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 1,
"reembed_events table must be writable post-migration"
);
}
#[test]
fn schema_v27_migration_replay_is_idempotent() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let _db = YantrikDB::new(path, 8).unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '26')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v27 migration runner must heal rewound-meta deployments on a v27-schema DB");
db.record(
"post-v27-heal smoke",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
#[test]
fn try_publish_search_state_rejects_stale_generation() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
assert_eq!(initial.generation, 0, "fresh engine starts at gen 0");
let stale_proposal = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: 0,
covers_through_seq: 0,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
let advanced = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: 2,
covers_through_seq: 0,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.search_state.store(std::sync::Arc::new(advanced));
let err = db
.try_publish_search_state(stale_proposal)
.expect_err("stale-generation publish must be rejected");
match err {
crate::error::YantrikDbError::SearchStatePublishStaleGeneration {
current_generation,
attempted_generation,
} => {
assert_eq!(current_generation, 2);
assert_eq!(attempted_generation, 0);
}
other => panic!("unexpected error variant: {other:?}"),
}
assert_eq!(
db.search_state.load().generation,
2,
"rejected publish must leave search_state untouched"
);
}
#[test]
fn try_publish_search_state_accepts_equal_generation_publish() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
let same_gen = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: Some("rotated-name".to_string()),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation,
covers_through_seq: initial.covers_through_seq,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(same_gen)
.expect("equal-generation publish must be accepted");
assert_eq!(
db.search_state.load().runtime_embedder_name.as_deref(),
Some("rotated-name"),
"the equal-generation publish must have landed"
);
}
#[test]
fn try_publish_search_state_accepts_strictly_advancing_generation() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
let advanced = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 1,
covers_through_seq: initial.covers_through_seq,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(advanced)
.expect("strictly-advancing-generation publish must be accepted");
assert_eq!(
db.search_state.load().generation,
initial.generation + 1,
"the advanced publish must have landed"
);
}
#[test]
fn schema_v28_fresh_install_has_embedding_generation_and_active_generation() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let cols = table_columns(&conn, "memories");
assert!(
cols.iter().any(|c| c == "embedding_generation"),
"v28: fresh install must add memories.embedding_generation column, got: {cols:?}"
);
assert!(
index_exists(&conn, "idx_memories_embedding_generation"),
"v28: fresh install must create idx_memories_embedding_generation"
);
let active_gen: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.ok();
assert_eq!(
active_gen.as_deref(),
Some("0"),
"v28: fresh install must seed meta.active_generation = '0'"
);
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
schema_version,
crate::base::schema::SCHEMA_VERSION.to_string(),
"fresh install stamps SCHEMA_VERSION"
);
}
#[test]
fn schema_v28_migration_from_v27_is_additive_and_idempotent() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000c028";
let planted_embedding = vec![0_u8; 8 * std::mem::size_of::<f32>()];
{
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
conn.execute(
"INSERT INTO memories (rid, type, text, embedding, created_at, updated_at, last_access, source, embedding_generation) \
VALUES (?1, 'episodic', 'planted under v28 schema', ?2, 0.0, 0.0, 0.0, 'user', 42)",
params![planted_rid, planted_embedding],
)
.unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '27')",
[],
)
.unwrap();
}
let db =
YantrikDB::new(path, 8).expect("v28 migration runner must heal rewound-meta deployments");
let conn = db.conn();
let (text, gen): (String, i64) = conn
.query_row(
"SELECT text, embedding_generation FROM memories WHERE rid = ?1",
[&planted_rid],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(text, "planted under v28 schema");
assert_eq!(gen, 42, "migration must not mutate existing row data");
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
schema_version,
crate::base::schema::SCHEMA_VERSION.to_string()
);
let active_gen: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active_gen, "0");
}
#[test]
fn record_stamps_embedding_generation_from_search_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"stamped",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let stamped: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(stamped, 0, "fresh engine: state.generation = 0, stamp = 0");
drop(conn);
let old_state = db.search_state.load_full();
let advanced = crate::engine::reembed::SearchState {
index_embedding: old_state.index_embedding.clone(),
embedder: old_state.embedder.clone(),
runtime_embedder_name: old_state.runtime_embedder_name.clone(),
runtime_embedder_digest: old_state.runtime_embedder_digest.clone(),
generation: 7,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&old_state.vec_index),
};
db.try_publish_search_state(advanced).unwrap();
let rid2 = db
.record(
"stamped at gen 7",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let stamped2: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&rid2],
|r| r.get(0),
)
.unwrap();
assert_eq!(
stamped2, 7,
"record after generation advance must stamp the new generation"
);
}
#[test]
fn open_reads_durable_active_generation_into_search_state() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(db.search_state.load().generation, 0);
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('active_generation', '3')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(
db.search_state.load().generation,
3,
"open() must read meta.active_generation into SearchState.generation"
);
}
#[test]
fn set_embedder_routes_through_try_publish_search_state() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
let gen_before = db.search_state.load().generation;
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:check".to_string()),
name: Some("check".to_string()),
sentinel: 0.1,
}))
.unwrap();
let gen_after = db.search_state.load().generation;
assert_eq!(
gen_before, gen_after,
"set_embedder must preserve generation (only Phase-2 reembed advances it)"
);
assert!(db.has_embedder(), "set_embedder must have published");
}
#[test]
fn search_state_publish_is_atomic_under_concurrent_reads() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
let db = Arc::new(YantrikDB::new(":memory:", 64).unwrap());
let initial = db.search_state.load_full();
let stop = Arc::new(AtomicBool::new(false));
let mut handles = Vec::new();
for _ in 0..4 {
let db_c = Arc::clone(&db);
let stop_c = Arc::clone(&stop);
handles.push(thread::spawn(move || {
let mut observations: Vec<(u64, usize, u32)> = Vec::new();
while !stop_c.load(Ordering::Relaxed) {
let s = db_c.search_state.load_full();
observations.push((s.generation, s.dim(), s.hnsw_m));
}
observations
}));
}
for n in 1..=50u64 {
let prev = db.search_state.load_full();
let next = crate::engine::reembed::SearchState {
index_embedding: prev.index_embedding.clone(),
embedder: prev.embedder.clone(),
runtime_embedder_name: prev.runtime_embedder_name.clone(),
runtime_embedder_digest: prev.runtime_embedder_digest.clone(),
generation: prev.generation + 1,
covers_through_seq: prev.covers_through_seq + n,
hnsw_m: if n % 2 == 0 { 16 } else { 32 },
hnsw_ef_construction: prev.hnsw_ef_construction,
hnsw_ef_search: prev.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&prev.vec_index),
};
db.try_publish_search_state(next).unwrap();
}
stop.store(true, Ordering::Relaxed);
let baseline = (initial.generation, initial.dim(), initial.hnsw_m);
for handle in handles {
let observations = handle.join().unwrap();
for (gen, dim, hnsw_m) in &observations {
let consistent = (*gen, *dim, *hnsw_m) == baseline
|| (*gen >= 1
&& *dim == initial.dim()
&& ((*gen % 2 == 0 && *hnsw_m == 16) || (*gen % 2 == 1 && *hnsw_m == 32)));
assert!(
consistent,
"torn SearchState observation: gen={gen} dim={dim} hnsw_m={hnsw_m} \
(expected even-gen→16 or odd-gen→32, or baseline)"
);
}
}
}
#[test]
fn open_recovery_discards_staging_when_sql_swap_uncommitted() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000d017";
{
let db = YantrikDB::new(path, 8).unwrap();
db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
conn.execute(
"UPDATE memories SET embedding_new = X'AABBCCDD', \
embedding_new_model = 'simulated-target' WHERE rowid = 1",
[],
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('reembed_state', ?1)",
params![serde_json::json!({
"generation": 5,
"phase": "Encoding",
"old_embedder": "old",
"new_embedder_name": "simulated-target",
})
.to_string()],
)
.unwrap();
let _ = planted_rid;
}
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
let still_in_flight: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'reembed_state'",
[],
|r| r.get(0),
)
.ok();
assert!(
still_in_flight.is_none(),
"Layer 7 must clear meta.reembed_state on uncommitted-swap recovery; got: {still_in_flight:?}"
);
let staged: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE embedding_new IS NOT NULL OR \
embedding_new_model IS NOT NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
staged, 0,
"staging columns must be NULL after discard recovery"
);
let active: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active, "0");
assert_eq!(db.search_state.load().generation, 0);
let aborted_recovery: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE phase = 'Aborted' AND generation = 5 \
AND payload_json LIKE '%discarded_staging%'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(aborted_recovery, 1, "Aborted recovery event must be logged");
}
#[test]
fn open_recovery_durable_swap_resumes_at_new_generation() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let db = YantrikDB::new(path, 8).unwrap();
let _ = db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
);
let conn = db.conn();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('active_generation', '3')",
[],
)
.unwrap();
conn.execute(
"UPDATE memories SET embedding_generation = 3 WHERE embedding IS NOT NULL",
[],
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('reembed_state', ?1)",
params![serde_json::json!({
"generation": 3,
"phase": "Swapping",
"old_embedder": "old",
"new_embedder_name": "new",
})
.to_string()],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
let still_in_flight: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'reembed_state'",
[],
|r| r.get(0),
)
.ok();
assert!(
still_in_flight.is_none(),
"Layer 7 must clear meta.reembed_state on durable-swap recovery"
);
assert_eq!(
db.search_state.load().generation,
3,
"SearchState rebuilds at durable active generation"
);
let active: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active, "3");
let completed_recovery: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE phase = 'Completed' AND generation = 3 \
AND payload_json LIKE '%completed_durable%'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
completed_recovery, 1,
"Completed recovery event must be logged"
);
}
#[test]
fn open_with_uncommitted_staging_columns_stays_at_old_generation() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = {
let db = YantrikDB::new(path, 8).unwrap();
db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap()
};
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"UPDATE memories SET embedding_new = X'AABBCCDD', \
embedding_new_model = 'simulated-new-embedder' WHERE rid = ?1",
params![planted_rid],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(
db.search_state.load().generation,
0,
"open() must not promote partial staging into the active generation"
);
let conn = db.conn();
let staged_present: bool = conn
.query_row(
"SELECT embedding_new IS NOT NULL FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert!(
staged_present,
"staged column survives the open (Phase-2 resume logic decides what to do with it)"
);
let active_present: bool = conn
.query_row(
"SELECT embedding IS NOT NULL FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert!(
active_present,
"pre-reembed active embedding bytes preserved"
);
let row_gen: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
row_gen, 0,
"row's stamped generation unchanged by partial staging"
);
}
#[test]
fn covers_through_seq_is_durably_carried_on_published_search_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let initial = db.search_state.load_full();
assert_eq!(initial.covers_through_seq, 0, "fresh engine: covers 0");
let next = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 1,
covers_through_seq: 12345,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(next).unwrap();
assert_eq!(
db.search_state.load().covers_through_seq,
12345,
"published covers_through_seq must be readable from the active SearchState"
);
let bumped = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 2,
covers_through_seq: 98765,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(bumped).unwrap();
assert_eq!(
db.search_state.load().covers_through_seq,
98765,
"covers_through_seq advances per swap"
);
}
#[test]
fn record_text_round_trip_through_queue_path_under_reembed() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:queued-test".to_string()),
name: Some("queued-test-embedder".to_string()),
sentinel: 0.7,
}))
.unwrap();
db.write_router.switch_to_queueing();
let rid = db
.record_text(
"queued-round-trip-text",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"k": "v"}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let mem_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(mem_count, 0, "queued write does not touch memories table");
let (op_type, applied, applied_generation, embedding_model, payload): (
String,
i64,
Option<i64>,
Option<String>,
String,
) = conn
.query_row(
"SELECT op_type, applied, applied_generation, embedding_model, payload \
FROM oplog WHERE target_rid = ?1 AND op_type = 'record'",
[&rid],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?)),
)
.unwrap();
assert_eq!(op_type, "record");
assert_eq!(applied, 0, "queued op is applied=0");
assert_eq!(
applied_generation, None,
"queued op has applied_generation=NULL (post-swap materializer fills under new gen)"
);
assert_eq!(
embedding_model.as_deref(),
Some("queued-test-embedder"),
"queued op carries the active runtime embedder name"
);
let v: serde_json::Value = serde_json::from_str(&payload).unwrap();
assert_eq!(
v["text"].as_str(),
Some("queued-round-trip-text"),
"payload preserves the original text for re-encode"
);
}
#[test]
fn boundary_audit_pattern_detects_synthetic_violation() {
let synthetic_violation =
" let sql = \"SELECT rid, embedding FROM memories WHERE rid = ?1\";";
let lower = synthetic_violation.to_ascii_lowercase();
let patterns = [
"select embedding ",
"select embedding,",
"select embedding\"",
"select embedding\\",
", embedding ",
", embedding,",
", embedding\"",
", embedding\\",
", embedding\n",
];
let any_match = patterns.iter().any(|p| lower.contains(p));
assert!(
any_match,
"boundary audit pattern must catch the synthetic raw-SQL-embedding pattern; \
if this asserts, the audit in durable_embeddings.rs is letting violations slip"
);
let allowlist_safe = " let sql = \"SELECT rid, embedding_hash FROM memories\";";
let lower_safe = allowlist_safe.to_ascii_lowercase();
let safe_match = patterns.iter().any(|p| lower_safe.contains(p));
assert!(
!safe_match,
"audit must NOT flag the allowlist-safe `embedding_hash` pattern; \
the audit over-rejects which would prevent legitimate refactors"
);
}
fn reembed_err(db: &YantrikDB, opts: crate::engine::reembed::ReembedOptions) -> String {
match db.reembed("test-embedder", opts) {
Err(e) => e.to_string(),
Ok(_) => panic!("an unimplemented ReembedOptions knob must error, not succeed"),
}
}
#[test]
fn unimplemented_reembed_namespace_is_rejected_not_silently_global() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let opts = crate::engine::reembed::ReembedOptions {
namespace: Some("only-this-one".to_string()),
..Default::default()
};
let msg = reembed_err(&db, opts);
assert!(
msg.contains("embedding space belongs to the ENGINE"),
"namespace must be refused as an engine-scoped-embedding request; got: {msg}"
);
}
#[test]
fn unimplemented_reembed_pause_policy_is_rejected() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let opts = crate::engine::reembed::ReembedOptions {
write_policy: crate::engine::reembed::ReembedWritePolicy::Pause,
..Default::default()
};
let msg = reembed_err(&db, opts);
assert!(
msg.contains("Pause is not implemented"),
"Pause must be refused rather than silently granting Queue; got: {msg}"
);
}
#[test]
fn unimplemented_reembed_resume_from_checkpoint_is_rejected() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let opts = crate::engine::reembed::ReembedOptions {
resume_from_checkpoint: true,
..Default::default()
};
let msg = reembed_err(&db, opts);
assert!(
msg.contains("resume_from_checkpoint is not implemented"),
"resume must be refused before an interrupted run discards its work; got: {msg}"
);
}
#[test]
fn scoped_reembed_must_not_disturb_other_namespaces() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..8 {
db.record(
&format!("keep record {i} about ledgers"),
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(i as f32 + 1.0, 8),
"keep",
0.9,
"general",
"user",
None,
)
.unwrap();
db.record(
&format!("touch record {i} about deployments"),
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(i as f32 + 100.0, 8),
"touch",
0.9,
"general",
"user",
None,
)
.unwrap();
}
let snapshot = |db: &YantrikDB| -> Vec<(String, i64, usize)> {
let conn = db.conn();
let mut stmt = conn
.prepare(
"SELECT rid, COALESCE(embedding_generation,0), length(embedding) \
FROM memories WHERE namespace = 'keep' ORDER BY rid",
)
.unwrap();
let rows = stmt
.query_map([], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, i64>(1)?,
r.get::<_, i64>(2)? as usize,
))
})
.unwrap()
.map(|r| r.unwrap())
.collect();
rows
};
let chunk_count = |db: &YantrikDB| -> i64 {
db.conn()
.query_row("SELECT COUNT(*) FROM memory_chunks", [], |r| r.get(0))
.unwrap_or(0)
};
let keep_before = snapshot(&db);
let chunks_before = chunk_count(&db);
assert!(
!keep_before.is_empty(),
"fixture must have 'keep' rows to protect"
);
let opts = crate::engine::reembed::ReembedOptions {
namespace: Some("touch".to_string()),
..Default::default()
};
let _ = db.reembed("test-embedder", opts);
assert_eq!(
keep_before,
snapshot(&db),
"a namespace-scoped reembed (or its rejection) changed rows in ANOTHER \
namespace — rid/generation/embedding-length must be bit-identical"
);
assert_eq!(
chunks_before,
chunk_count(&db),
"memory_chunks was purged globally; untouched namespaces lost their \
chunk vectors (reembed.rs DELETE FROM memory_chunks has no predicate)"
);
}
#[test]
fn record_with_rid_defers_instead_of_crossing_a_reembed_cutover() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.write_router.switch_to_queueing();
let err = db
.record_with_rid(
"01900000-0000-7000-8000-0000000000ab",
"deterministic replicated fact",
"semantic",
0.6,
0.0,
1000.0,
&empty_meta(),
&vec_seed(1.0, 8),
"work",
0.9,
"general",
"inference",
None,
1_700_000_000_000_000,
&[],
"test-model",
None,
crate::provenance::WriteAdmission::Admitted,
)
.expect_err("a cutover in flight must defer, not commit against a doomed generation");
assert!(
matches!(
err,
crate::error::YantrikDbError::DeterministicWriteDeferredDuringReembed { .. }
),
"must be the typed retryable deferral, got: {err}"
);
let n: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = '01900000-0000-7000-8000-0000000000ab'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
n, 0,
"a deferred deterministic write must leave no row behind"
);
db.write_router.switch_to_normal();
db.record_with_rid(
"01900000-0000-7000-8000-0000000000ab",
"deterministic replicated fact",
"semantic",
0.6,
0.0,
1000.0,
&empty_meta(),
&vec_seed(1.0, 8),
"work",
0.9,
"general",
"inference",
None,
1_700_000_000_000_000,
&[],
"test-model",
None,
crate::provenance::WriteAdmission::Admitted,
)
.expect("must succeed once the router is Normal again");
}
#[test]
fn forget_defers_rather_than_tombstoning_into_a_doomed_index() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"a fact the user will ask to forget",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
db.write_router.switch_to_queueing();
let err = db
.forget(&rid)
.expect_err("forget during a cutover must defer, not tombstone into a doomed index");
assert!(
matches!(
err,
crate::error::YantrikDbError::ForgetDeferredDuringReembed { .. }
),
"must be the typed retryable deferral, got: {err}"
);
let status: String = db
.conn()
.query_row(
"SELECT consolidation_status FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
status, "active",
"a deferred forget must leave the row exactly as it was"
);
db.write_router.switch_to_normal();
assert!(db.forget(&rid).unwrap(), "forget must succeed post-cutover");
let status: String = db
.conn()
.query_row(
"SELECT consolidation_status FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(status, "tombstoned");
}
#[test]
fn rebuild_vec_index_discards_rather_than_mixing_embedding_spaces() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..6 {
db.record(
&format!("rebuildable record {i}"),
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(i as f32 + 1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
}
let n = db.rebuild_vec_index().expect("rebuild must work normally");
assert!(n > 0, "rebuild should install a populated cold tier");
db.write_router.switch_to_queueing();
let err = db
.rebuild_vec_index()
.expect_err("a rebuild finishing during a cutover must not install");
assert!(
matches!(
err,
crate::error::YantrikDbError::IndexRebuildDeferredDuringReembed { .. }
),
"must be the typed retryable deferral, got: {err}"
);
db.write_router.switch_to_normal();
assert!(
db.rebuild_vec_index().unwrap() > 0,
"rebuild must work again once the cutover completes"
);
}