#![cfg(feature = "default-embedder")]
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
use fathomdb_embedder::CandleBgeEmbedder;
use fathomdb_embedder_api::Embedder;
use fathomdb_engine::lifecycle::ProjectionStatus;
use fathomdb_engine::{EmbedderChoice, Engine, PreparedWrite};
use tempfile::TempDir;
#[path = "support/corpus_subset.rs"]
mod corpus_subset;
use corpus_subset::load_subset_or_skip;
const DEFAULT_SEED_N: usize = 10_000;
const WRITE_BATCH: usize = 256;
const DRAIN_TIMEOUT_MS: u64 = 90 * 60 * 1000;
fn env_usize(key: &str, default: usize) -> usize {
std::env::var(key).ok().and_then(|v| v.parse().ok()).unwrap_or(default)
}
fn cosine(a: &[f32], b: &[f32]) -> f32 {
let dot: f32 = a.iter().zip(b).map(|(x, y)| x * y).sum();
let na: f32 = a.iter().map(|x| x * x).sum::<f32>().sqrt();
let nb: f32 = b.iter().map(|x| x * x).sum::<f32>().sqrt();
dot / (na * nb).max(1e-12)
}
#[test]
fn sustained_seed_serialized_path_completes_and_is_correct() {
if std::env::var_os("AGENT_LONG").is_none() {
eprintln!("[skip] AGENT_LONG not set; PR-9 sustained-seed measurement is opt-in");
return;
}
if std::env::var("FATHOMDB_SKIP_NETWORK_TESTS").is_ok() {
eprintln!("[skip] FATHOMDB_SKIP_NETWORK_TESTS set; embedder cache unavailable");
return;
}
let Some(docs) = load_subset_or_skip(usize::MAX) else {
eprintln!("[skip] corpus not present; cannot run PR-9 concurrent-embed seed");
return;
};
let target_n = env_usize("PR9_SEED_N", DEFAULT_SEED_N);
let real_bodies: Vec<String> =
docs.iter().map(|d| d.body.clone()).filter(|b| !b.trim().is_empty()).collect();
assert!(!real_bodies.is_empty(), "corpus yielded no non-empty bodies");
let bodies: Vec<String> =
(0..target_n).map(|i| real_bodies[i % real_bodies.len()].clone()).collect();
eprintln!("PR9_SETUP target_n={target_n} real_docs={} (cycled to fill)", real_bodies.len());
let embedder: Arc<dyn Embedder> =
Arc::new(CandleBgeEmbedder::new().expect("construct real bge embedder"));
let dir = TempDir::new().expect("tempdir");
let path = dir.path().join("pr9_concurrent.sqlite");
let opened = Engine::open_with_choice(&path, EmbedderChoice::Caller(embedder.clone()))
.expect("open with real bge embedder");
let engine = Arc::new(opened.engine);
assert_eq!(
opened.report.default_embedder.name, "fathomdb-bge-small-en-v1.5",
"must run against the real bge-small identity"
);
engine.configure_vector_kind_for_test("doc").expect("configure vector kind");
let started = Instant::now();
let mut written = 0usize;
while written < bodies.len() {
let take = WRITE_BATCH.min(bodies.len() - written);
let batch: Vec<PreparedWrite> = bodies[written..written + take]
.iter()
.map(|b| PreparedWrite::Node {
kind: "doc".to_string(),
body: b.clone(),
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,
})
.collect();
engine.write(&batch).expect("seed write");
written += take;
}
eprintln!("PR9_WROTE n={written} enqueue_s={}", started.elapsed().as_secs());
let monitor_engine = engine.clone();
let monitor_done = Arc::new(AtomicBool::new(false));
let monitor_flag = monitor_done.clone();
let monitor = thread::spawn(move || {
let started = Instant::now();
let mut prev = 0u64;
let mut stalls = 0u32;
while !monitor_flag.load(Ordering::Relaxed) {
thread::sleep(Duration::from_secs(20));
let count = monitor_engine.vector_row_count_for_test().unwrap_or(u64::MAX);
let delta = count.saturating_sub(prev);
if delta == 0 {
stalls += 1;
} else {
stalls = 0;
}
eprintln!(
"PR9_PROGRESS elapsed_s={} vector_rows={count} delta={delta} stalls={stalls}",
started.elapsed().as_secs()
);
prev = count;
}
});
let drain_started = Instant::now();
let drained = engine.drain(DRAIN_TIMEOUT_MS);
monitor_done.store(true, Ordering::Relaxed);
let _ = monitor.join();
eprintln!("PR9_DRAINED ok={} drain_s={}", drained.is_ok(), drain_started.elapsed().as_secs());
assert!(
drained.is_ok(),
"serialized seed must complete drain in bounded time; got {drained:?}"
);
assert_eq!(
engine.projection_status_for_test("doc").expect("projection status"),
ProjectionStatus::UpToDate,
"projection must reach UpToDate after the concurrent seed"
);
let row_count = engine.vector_row_count_for_test().expect("vector row count");
assert_eq!(
row_count as usize, written,
"every seeded doc must produce exactly one vector row (no concurrency drops)"
);
let dim = opened.report.default_embedder.dimension as usize;
let mut checked = 0u64;
for rowid in 1..=row_count as i64 {
let blob = engine.read_vector_blob_for_test(rowid).expect("read vector blob");
let v: Vec<f32> =
blob.chunks_exact(4).map(|c| f32::from_le_bytes([c[0], c[1], c[2], c[3]])).collect();
assert_eq!(v.len(), dim, "row {rowid} vector has wrong dimension");
assert!(v.iter().all(|x| x.is_finite()), "row {rowid} vector has non-finite values");
let norm: f32 = v.iter().map(|x| x * x).sum::<f32>().sqrt();
assert!(
(norm - 1.0).abs() < 1e-3,
"row {rowid} vector not unit-norm (got {norm}); concurrency corruption suspected"
);
checked += 1;
}
eprintln!("PR9_NORM_OK checked={checked} dim={dim}");
let spot = (row_count as usize).min(16);
for i in 0..spot {
let rowid = (i + 1) as i64;
let blob = engine.read_vector_blob_for_test(rowid).expect("read vector blob");
let stored: Vec<f32> =
blob.chunks_exact(4).map(|c| f32::from_le_bytes([c[0], c[1], c[2], c[3]])).collect();
let expected = embedder.embed(&bodies[i]).expect("re-embed for spot-check");
let cos = cosine(&stored, &expected);
assert!(
cos > 0.9999,
"row {rowid} stored vector diverges from single-threaded re-embed (cos={cos}); \
concurrent-embed produced a wrong vector"
);
}
eprintln!("PR9_SPOTCHECK_OK rows={spot} verdict=SERIALIZED_SEED_COMPLETE_AND_CORRECT");
}