use super::*;
fn tempdir() -> tempfile::TempDir {
tempfile::tempdir().expect("tempdir")
}
#[tokio::test]
async fn data_dir_matches_the_daemon_era_layout() {
let root = std::path::Path::new("/data/root");
assert_eq!(
data_dir_for_palace(root, "my-palace"),
std::path::Path::new("/data/root/my-palace/bm25")
);
let lane = Bm25Lane::with_limits(root.to_path_buf(), 3, None);
assert_eq!(
lane.data_dir_for_palace("my-palace"),
data_dir_for_palace(root, "my-palace"),
"the method and the free function must not drift"
);
lane.shutdown().await;
}
#[tokio::test]
async fn default_cap_is_three() {
assert_eq!(DEFAULT_MAX_RESIDENT, 3);
}
#[test]
#[serial_test::serial]
fn max_resident_honours_env_override() {
let _env = crate::commands::env_test_lock().blocking_lock();
let prev = std::env::var(ENV_MAX_PALACES).ok();
unsafe { std::env::set_var(ENV_MAX_PALACES, "7") };
assert_eq!(max_resident_from_env(), 7);
unsafe { std::env::set_var(ENV_MAX_PALACES, "not-a-number") };
assert_eq!(max_resident_from_env(), DEFAULT_MAX_RESIDENT);
unsafe { std::env::set_var(ENV_MAX_PALACES, "0") };
assert_eq!(
max_resident_from_env(),
DEFAULT_MAX_RESIDENT,
"zero is a typo, not a request to evict everything"
);
match prev {
Some(v) => unsafe { std::env::set_var(ENV_MAX_PALACES, v) },
None => unsafe { std::env::remove_var(ENV_MAX_PALACES) },
}
}
#[test]
#[serial_test::serial]
fn text_budget_honours_env_override() {
let _env = crate::commands::env_test_lock().blocking_lock();
let prev = std::env::var(ENV_TEXT_BUDGET_MB).ok();
unsafe { std::env::set_var(ENV_TEXT_BUDGET_MB, "64") };
assert_eq!(text_budget_from_env(), Some(64));
unsafe { std::env::set_var(ENV_TEXT_BUDGET_MB, "0") };
assert_eq!(text_budget_from_env(), None, "0 disables enforcement");
unsafe { std::env::set_var(ENV_TEXT_BUDGET_MB, "garbage") };
assert_eq!(text_budget_from_env(), Some(DEFAULT_TEXT_BUDGET_MB));
match prev {
Some(v) => unsafe { std::env::set_var(ENV_TEXT_BUDGET_MB, v) },
None => unsafe { std::env::remove_var(ENV_TEXT_BUDGET_MB) },
}
}
#[tokio::test]
async fn cap_is_clamped_to_at_least_one() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 0, None);
assert_eq!(lane.max_resident(), 1);
lane.index("p", "d", "alpha").await.unwrap();
assert_eq!(lane.resident_count().await, 1);
lane.shutdown().await;
}
#[tokio::test]
async fn index_then_search_finds_the_document() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "d1", "authentication login password")
.await
.unwrap();
lane.index("alpha", "d2", "rendering ui components")
.await
.unwrap();
let hits = lane.search("alpha", "authentication", 5).await.unwrap();
assert_eq!(hits.len(), 1, "got: {hits:?}");
assert_eq!(hits[0].doc_id, "d1");
lane.shutdown().await;
}
#[tokio::test]
async fn palaces_do_not_share_a_corpus() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "shared-id", "kangaroo").await.unwrap();
lane.index("beta", "shared-id", "platypus").await.unwrap();
assert_eq!(lane.search("alpha", "kangaroo", 5).await.unwrap().len(), 1);
assert!(lane
.search("alpha", "platypus", 5)
.await
.unwrap()
.is_empty());
assert_eq!(lane.search("beta", "platypus", 5).await.unwrap().len(), 1);
assert!(lane.search("beta", "kangaroo", 5).await.unwrap().is_empty());
lane.shutdown().await;
}
#[tokio::test]
async fn delete_removes_the_document() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "d1", "ephemeral note").await.unwrap();
assert_eq!(lane.search("alpha", "ephemeral", 5).await.unwrap().len(), 1);
lane.delete("alpha", "d1").await.unwrap();
assert!(lane
.search("alpha", "ephemeral", 5)
.await
.unwrap()
.is_empty());
lane.delete("alpha", "never-existed").await.unwrap();
lane.shutdown().await;
}
#[tokio::test]
async fn stats_report_docs_and_bytes() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
let empty = lane.stats("alpha").await.unwrap();
assert_eq!(empty.doc_count, 0);
assert_eq!(empty.total_text_bytes, 0);
lane.index("alpha", "a", "hello").await.unwrap();
lane.index("alpha", "b", "world!").await.unwrap();
let stats = lane.stats("alpha").await.unwrap();
assert_eq!(stats.doc_count, 2);
assert_eq!(stats.total_text_bytes, 11);
lane.shutdown().await;
}
#[tokio::test]
async fn missing_docs_answers_by_identity() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "a", "alpha").await.unwrap();
lane.index("alpha", "stale", "a forgotten drawer")
.await
.unwrap();
let asked = vec!["a".to_string(), "b".to_string()];
let cov = lane.missing_docs("alpha", &asked).await.unwrap();
assert_eq!(cov.checked, 2);
assert_eq!(cov.missing, vec!["b".to_string()]);
lane.index("alpha", "b", "beta").await.unwrap();
assert!(lane
.missing_docs("alpha", &asked)
.await
.unwrap()
.missing
.is_empty());
lane.shutdown().await;
}
#[tokio::test]
async fn a_write_reaches_disk_without_an_explicit_flush() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "d1", "persisted by the ticker")
.await
.unwrap();
let snapshot = lane
.data_dir_for_palace("alpha")
.join(crate::bm25_index::SNAPSHOT_FILENAME);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while !snapshot.exists() {
assert!(
std::time::Instant::now() < deadline,
"the flush ticker never wrote {}",
snapshot.display()
);
tokio::time::sleep(FLUSH_INTERVAL).await;
}
let raw = std::fs::read_to_string(&snapshot).unwrap();
assert!(raw.contains("persisted by the ticker"), "got: {raw}");
lane.shutdown().await;
}
#[tokio::test]
async fn flush_persists_a_pending_write() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "d1", "flushed on demand")
.await
.unwrap();
lane.flush("alpha").await.unwrap();
let snapshot = lane
.data_dir_for_palace("alpha")
.join(crate::bm25_index::SNAPSHOT_FILENAME);
let raw =
std::fs::read_to_string(&snapshot).expect("snapshot must exist after an explicit flush");
assert!(raw.contains("flushed on demand"), "got: {raw}");
lane.flush("never-touched").await.unwrap();
lane.shutdown().await;
}
#[tokio::test]
async fn shutdown_flushes_and_is_idempotent() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
lane.index("alpha", "d1", "written just before exit")
.await
.unwrap();
lane.shutdown().await;
let snapshot = lane
.data_dir_for_palace("alpha")
.join(crate::bm25_index::SNAPSHOT_FILENAME);
let raw = std::fs::read_to_string(&snapshot).expect("shutdown must flush");
assert!(raw.contains("written just before exit"), "got: {raw}");
lane.shutdown().await;
let reopened = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
let hits = reopened.search("alpha", "exit", 5).await.unwrap();
assert_eq!(hits.len(), 1, "got: {hits:?}");
reopened.shutdown().await;
}
#[tokio::test]
async fn eviction_flushes_before_dropping_the_index() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 1, None);
lane.index("alpha", "d1", "evicted but not lost")
.await
.unwrap();
lane.index("beta", "d2", "second palace").await.unwrap();
assert_eq!(lane.resident_count().await, 1, "cap of 1 must hold");
assert_eq!(lane.evicted_count(), 1);
let hits = lane.search("alpha", "evicted", 5).await.unwrap();
assert_eq!(hits.len(), 1, "the evicted write must survive: {hits:?}");
lane.shutdown().await;
}
#[tokio::test]
async fn over_budget_evicts_the_coldest() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 8, Some(0));
assert_eq!(lane.text_budget_bytes(), Some(0));
lane.index("alpha", "d1", "alpha text").await.unwrap();
lane.index("beta", "d2", "beta text").await.unwrap();
lane.enforce_text_budget().await;
assert_eq!(
lane.resident_count().await,
1,
"the budget must evict down to one, and stop there"
);
assert!(
lane.evicted_count() >= 1,
"a zero budget over two palaces must have evicted"
);
assert_eq!(lane.search("beta", "beta", 5).await.unwrap().len(), 1);
assert_eq!(lane.search("alpha", "alpha", 5).await.unwrap().len(), 1);
lane.shutdown().await;
}
#[tokio::test]
async fn a_disabled_budget_never_evicts() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 8, None);
assert_eq!(lane.text_budget_bytes(), None);
lane.index("alpha", "d1", "alpha text").await.unwrap();
lane.index("beta", "d2", "beta text").await.unwrap();
lane.enforce_text_budget().await;
assert_eq!(lane.resident_count().await, 2);
assert_eq!(lane.evicted_count(), 0);
lane.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_concurrent_fanout_never_exceeds_the_cap() {
const CAP: usize = 3;
const PALACES: usize = 12;
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), CAP, None);
let mut tasks = Vec::new();
for i in 0..PALACES {
for doc in 0..2 {
let lane = Arc::clone(&lane);
tasks.push(tokio::spawn(async move {
lane.index(
&format!("palace-{i}"),
&format!("doc-{doc}"),
&format!("unique-token-{i}-{doc}"),
)
.await
}));
}
}
for t in tasks {
t.await.expect("task joined").expect("index succeeded");
}
assert!(
lane.resident_count().await <= CAP,
"resident={} exceeded cap={CAP}",
lane.resident_count().await
);
assert!(
lane.evicted_count() > 0,
"a {PALACES}-palace fanout under a cap of {CAP} must have evicted something"
);
for i in 0..PALACES {
for doc in 0..2 {
let hits = lane
.search(
&format!("palace-{i}"),
&format!("unique-token-{i}-{doc}"),
5,
)
.await
.unwrap();
assert_eq!(
hits.first().map(|h| h.doc_id.as_str()),
Some(format!("doc-{doc}").as_str()),
"palace-{i}/doc-{doc} was lost across evictions: {hits:?}"
);
}
}
lane.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_callers_for_one_palace_share_one_index() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
let mut tasks = Vec::new();
for i in 0..16 {
let lane = Arc::clone(&lane);
tasks.push(tokio::spawn(async move {
lane.index("hot", &format!("doc-{i}"), &format!("token{i}"))
.await
}));
}
for t in tasks {
t.await.expect("task joined").expect("index succeeded");
}
assert_eq!(
lane.loaded_count(),
1,
"a single palace must be loaded exactly once no matter how many callers race"
);
let stats = lane.stats("hot").await.unwrap();
assert_eq!(
stats.doc_count, 16,
"every concurrent write must have landed"
);
lane.shutdown().await;
}
#[tokio::test]
async fn a_cold_load_failure_propagates() {
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 3, None);
let palace_dir = dir.path().join("blocked");
std::fs::create_dir_all(&palace_dir).unwrap();
std::fs::write(palace_dir.join("bm25"), b"i am a file, not a directory").unwrap();
let err = lane
.search("blocked", "anything", 5)
.await
.expect_err("a palace whose bm25 dir cannot be created must error");
assert!(
format!("{err:#}").contains("blocked"),
"the error must name the palace: {err:#}"
);
lane.shutdown().await;
}
#[tokio::test]
#[cfg(unix)]
async fn the_budget_keeps_a_palace_whose_snapshot_cannot_be_flushed() {
use std::os::unix::fs::PermissionsExt;
if unsafe { libc::geteuid() } == 0 {
return;
}
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 8, Some(0));
lane.shutdown().await;
lane.index("alpha", "d1", "unflushable-but-not-lost")
.await
.unwrap();
lane.index("beta", "d2", "beta text").await.unwrap();
let alpha_dir = lane.data_dir_for_palace("alpha");
let alpha_snapshot = alpha_dir.join(crate::bm25_index::SNAPSHOT_FILENAME);
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o500)).unwrap();
lane.enforce_text_budget().await;
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o755)).unwrap();
assert!(
!alpha_snapshot.exists(),
"the sealed directory must have failed alpha's flush, but {} exists",
alpha_snapshot.display()
);
assert_eq!(
lane.resident_count().await,
1,
"a zero budget over two palaces must evict exactly one"
);
let hits = lane
.search("alpha", "unflushable-but-not-lost", 5)
.await
.unwrap();
assert_eq!(
hits.len(),
1,
"the unflushable palace's write was dropped with the index: {hits:?}"
);
assert_eq!(
lane.loaded_count(),
2,
"alpha was evicted and reloaded — the lane dropped an index it could not flush"
);
}
#[tokio::test]
#[cfg(unix)]
async fn a_cold_load_refuses_to_evict_an_unflushable_victim() {
use std::os::unix::fs::PermissionsExt;
if unsafe { libc::geteuid() } == 0 {
return;
}
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 1, None);
lane.shutdown().await;
lane.index("alpha", "d1", "unflushable-but-not-lost")
.await
.unwrap();
let alpha_dir = lane.data_dir_for_palace("alpha");
let alpha_snapshot = alpha_dir.join(crate::bm25_index::SNAPSHOT_FILENAME);
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o500)).unwrap();
let result = lane.index("beta", "d2", "beta text").await;
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o755)).unwrap();
assert!(
!alpha_snapshot.exists(),
"the sealed directory must have failed alpha's flush, but {} exists",
alpha_snapshot.display()
);
let err = result.expect_err("loading beta must fail rather than drop alpha's unflushed write");
assert!(
format!("{err:#}").contains("beta"),
"the error must name the palace that could not be loaded: {err:#}"
);
assert_eq!(
lane.evicted_count(),
0,
"nothing may be evicted when no resident snapshot could be flushed"
);
assert_eq!(
lane.loaded_count(),
1,
"a load discarded by a failed eviction must not count as resident"
);
assert_eq!(lane.resident_count().await, 1);
let hits = lane
.search("alpha", "unflushable-but-not-lost", 5)
.await
.unwrap();
assert_eq!(
hits.len(),
1,
"alpha's write was lost to a failed eviction: {hits:?}"
);
}
#[tokio::test]
#[cfg(unix)]
async fn the_budget_evicts_past_a_palace_it_cannot_flush() {
use std::os::unix::fs::PermissionsExt;
if unsafe { libc::geteuid() } == 0 {
return;
}
let dir = tempdir();
let lane = Bm25Lane::with_limits(dir.path().to_path_buf(), 8, Some(0));
lane.shutdown().await;
lane.index("alpha", "d1", "unflushable-but-not-lost")
.await
.unwrap();
for (palace, doc, text) in [
("beta", "d2", "beta text"),
("gamma", "d3", "gamma text"),
("delta", "d4", "delta text"),
] {
lane.index(palace, doc, text).await.unwrap();
}
assert_eq!(
lane.resident_count().await,
4,
"cap of 8 must hold all four"
);
let alpha_dir = lane.data_dir_for_palace("alpha");
let alpha_snapshot = alpha_dir.join(crate::bm25_index::SNAPSHOT_FILENAME);
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o500)).unwrap();
lane.enforce_text_budget().await;
std::fs::set_permissions(&alpha_dir, std::fs::Permissions::from_mode(0o755)).unwrap();
assert!(
!alpha_snapshot.exists(),
"the sealed directory must have failed alpha's flush, but {} exists",
alpha_snapshot.display()
);
assert_eq!(
lane.resident_count().await,
1,
"one unwritable palace must not stop the other three being evicted"
);
assert_eq!(
lane.evicted_count(),
3,
"beta, gamma and delta were all flushable and must all have been evicted"
);
let hits = lane
.search("alpha", "unflushable-but-not-lost", 5)
.await
.unwrap();
assert_eq!(
hits.len(),
1,
"the surviving palace must be alpha, with its write intact: {hits:?}"
);
assert_eq!(
lane.loaded_count(),
4,
"alpha was evicted and reloaded — the lane dropped an index it could not flush"
);
}
#[test]
fn bm25_hit_round_trips() {
let h = BM25Hit {
doc_id: "drawer-1".into(),
score: 0.42,
};
let s = serde_json::to_string(&h).unwrap();
let back: BM25Hit = serde_json::from_str(&s).unwrap();
assert_eq!(back.doc_id, "drawer-1");
assert!((back.score - 0.42).abs() < 1e-6);
}