use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
use fathomdb_engine::{Engine, PreparedWrite};
use fathomdb_schema::SQLITE_SUFFIX;
use tempfile::TempDir;
const READER_POOL_SIZE: usize = 8;
fn fresh_engine(name: &str) -> (TempDir, fathomdb_engine::OpenedEngine) {
let dir = TempDir::new().unwrap();
let path = dir.path().join(format!("{name}{SQLITE_SUFFIX}"));
let opened = Engine::open(path).expect("engine open");
(dir, opened)
}
#[test]
fn reader_pool_spawns_eight_worker_threads_at_open() {
let (_dir, opened) = fresh_engine("reader_pool_open");
assert_eq!(opened.engine.reader_worker_count_for_test(), READER_POOL_SIZE);
let deadline = Instant::now() + Duration::from_secs(2);
while opened.engine.live_reader_worker_count_for_test() < READER_POOL_SIZE
&& Instant::now() < deadline
{
thread::sleep(Duration::from_millis(5));
}
assert_eq!(opened.engine.live_reader_worker_count_for_test(), READER_POOL_SIZE);
}
#[test]
fn reader_workers_exit_on_close_and_drop_connections() {
let (_dir, opened) = fresh_engine("reader_pool_close");
for _ in 0..32 {
let _ = opened.engine.search("ping").err();
}
opened.engine.close().expect("close");
let deadline = Instant::now() + Duration::from_secs(2);
while opened.engine.live_reader_worker_count_for_test() != 0 && Instant::now() < deadline {
thread::sleep(Duration::from_millis(5));
}
assert_eq!(
opened.engine.live_reader_worker_count_for_test(),
0,
"reader workers did not exit on close"
);
let err = opened.engine.search("ping").expect_err("search after close");
let msg = err.to_string();
assert!(
msg.contains("closing") || msg.contains("storage"),
"unexpected post-close search error: {msg}"
);
}
#[test]
fn reader_workers_exit_on_drop_without_explicit_close() {
let live = {
let (_dir, opened) = fresh_engine("reader_pool_drop");
let deadline = Instant::now() + Duration::from_secs(2);
while opened.engine.live_reader_worker_count_for_test() < READER_POOL_SIZE
&& Instant::now() < deadline
{
thread::sleep(Duration::from_millis(5));
}
let live = opened.engine.live_reader_worker_count_for_test();
drop(opened);
live
};
assert_eq!(live, READER_POOL_SIZE);
assert_eq!(live, READER_POOL_SIZE);
}
#[test]
fn concurrent_searches_route_to_workers_without_loss_or_duplication() {
let (_dir, opened) = fresh_engine("reader_pool_routing");
opened
.engine
.write(&[PreparedWrite::Node {
kind: "doc".to_string(),
body: "hello world".to_string(),
source_id: fathomdb_engine::SourceId::new("test:fixture").expect("test source id"),
logical_id: None,
state: fathomdb_engine::InitialState::Active,
reason: None,
valid_from: None,
valid_until: None,
}])
.expect("seed write");
const CALLERS: usize = 32;
const PER_CALLER: usize = 25;
let engine = Arc::new(opened.engine);
let success = Arc::new(AtomicUsize::new(0));
let failure = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::with_capacity(CALLERS);
for _ in 0..CALLERS {
let engine = Arc::clone(&engine);
let success = Arc::clone(&success);
let failure = Arc::clone(&failure);
handles.push(thread::spawn(move || {
for _ in 0..PER_CALLER {
match engine.search("hello") {
Ok(_) => {
success.fetch_add(1, Ordering::SeqCst);
}
Err(_) => {
failure.fetch_add(1, Ordering::SeqCst);
}
}
}
}));
}
for handle in handles {
handle.join().expect("caller thread");
}
assert_eq!(failure.load(Ordering::SeqCst), 0, "no search may fail");
assert_eq!(success.load(Ordering::SeqCst), CALLERS * PER_CALLER);
assert_eq!(engine.live_reader_worker_count_for_test(), READER_POOL_SIZE);
}
#[test]
fn reader_workers_have_lookaside_configured_with_ok_rc() {
let (_dir, opened) = fresh_engine("reader_pool_lookaside_rc");
let rcs = opened.engine.reader_lookaside_config_rcs_for_test();
assert_eq!(rcs.len(), READER_POOL_SIZE);
for (idx, rc) in rcs.iter().enumerate() {
assert_eq!(*rc, 0, "worker {idx} sqlite3_db_config(LOOKASIDE) rc must be SQLITE_OK");
}
}
#[test]
fn reader_workers_consume_lookaside_slots_after_warmup_read() {
let (_dir, opened) = fresh_engine("reader_pool_lookaside_used");
for _ in 0..(READER_POOL_SIZE * 8) {
let _ = opened.engine.search("warmup").err();
}
let used = opened.engine.reader_lookaside_used_per_worker_for_test();
eprintln!("LOOKASIDE_USED_HIWTR_PER_WORKER={used:?}");
assert_eq!(used.len(), READER_POOL_SIZE);
for (idx, slots) in used.iter().enumerate() {
assert!(
*slots > 0,
"worker {idx} SQLITE_DBSTATUS_LOOKASIDE_USED must be >0 after warmup; got {slots}",
);
}
}