#![allow(dead_code)]
#[path = "support/corpus_harness.rs"]
mod corpus_harness;
use std::sync::{Arc, Barrier, Mutex};
use std::thread;
use std::time::{Duration, Instant};
use corpus_harness::{CorpusFixture, HeldOutQuery, VaryingEmbedder, CORPUS_DIM};
use fathomdb_embedder_api::Embedder;
use fathomdb_engine::Engine;
const DEVLOOP_BUDGET_P50: Duration = Duration::from_millis(50);
const DEVLOOP_BUDGET_P99: Duration = Duration::from_millis(150);
const DEVLOOP_RECALL_FLOOR: f64 = 0.85;
const DEVLOOP_CATASTROPHIC_MULT: u32 = 10;
const DEVLOOP_SAMPLES: usize = 100;
const DEVLOOP_SEED: u64 = 0x0DEF_1009_0DEF_1009;
const DEVLOOP_AC019_THREADS: usize = 4;
const DEVLOOP_AC019_QUERIES_PER_THREAD: usize = 50;
fn real_embedder_requested() -> bool {
matches!(std::env::var("DEVLOOP_REAL_EMBEDDER"), Ok(v) if v == "1" || v.eq_ignore_ascii_case("true"))
}
fn devloop_fixture() -> (CorpusFixture, &'static str) {
if real_embedder_requested() {
(CorpusFixture::medium().with_real_embedder(), "real")
} else {
(CorpusFixture::medium().with_synthetic_embedder(), "synthetic")
}
}
#[cfg(feature = "default-embedder")]
fn ground_truth_embedder(real: bool) -> Arc<dyn Embedder> {
if real {
Arc::new(fathomdb_embedder::CandleBgeEmbedder::new().expect("construct real bge embedder"))
} else {
Arc::new(VaryingEmbedder::new(CORPUS_DIM))
}
}
#[cfg(not(feature = "default-embedder"))]
fn ground_truth_embedder(_real: bool) -> Arc<dyn Embedder> {
Arc::new(VaryingEmbedder::new(CORPUS_DIM))
}
fn percentile_ceil(samples: &[Duration], numerator: usize, denominator: usize) -> Duration {
if samples.is_empty() {
return Duration::ZERO;
}
let mut sorted = samples.to_vec();
sorted.sort_unstable();
let rank = (sorted.len() * numerator).div_ceil(denominator);
let idx = rank.saturating_sub(1).min(sorted.len() - 1);
sorted[idx]
}
fn cache_field(report: &corpus_harness::IngestReport, embedder: &str) -> String {
if embedder == "synthetic" {
return "na".to_string();
}
match &report.cache_miss_reason {
None => "hit".to_string(),
Some(reason) => reason.replace(' ', "_"),
}
}
fn emit_numbers(
ac: &str,
n: usize,
samples: usize,
(p50, p99): (Duration, Duration),
recall: Option<f64>,
cache: &str,
embedder: &str,
) {
let recall_field = match recall {
Some(r) => format!("{r:.4}"),
None => "NA".to_string(),
};
eprintln!(
"DEVLOOP_NUMBERS ac={ac} n={n} samples={samples} p50_ms={p50} p99_ms={p99} \
recall_at_10={recall_field} cache={cache} embedder={embedder}",
p50 = p50.as_millis(),
p99 = p99.as_millis(),
);
}
fn warn_if_over_latency(ac: &str, p50: Duration, p99: Duration) {
if p50 > DEVLOOP_BUDGET_P50 {
eprintln!(
"DEVLOOP_PERF_WARN ac={ac} metric=p50 value_ms={} budget_ms={} status=OVER",
p50.as_millis(),
DEVLOOP_BUDGET_P50.as_millis(),
);
}
if p99 > DEVLOOP_BUDGET_P99 {
eprintln!(
"DEVLOOP_PERF_WARN ac={ac} metric=p99 value_ms={} budget_ms={} status=OVER",
p99.as_millis(),
DEVLOOP_BUDGET_P99.as_millis(),
);
}
}
fn warn_if_under_recall(ac: &str, recall: f64) {
if recall < DEVLOOP_RECALL_FLOOR {
eprintln!(
"DEVLOOP_PERF_WARN ac={ac} metric=recall_at_10 value={recall:.4} floor={floor:.2} status=UNDER",
floor = DEVLOOP_RECALL_FLOOR,
);
}
}
fn enforce_catastrophic_ceiling(ac: &str, p50: Duration, p99: Duration) {
let ceil_p50 = DEVLOOP_BUDGET_P50 * DEVLOOP_CATASTROPHIC_MULT;
let ceil_p99 = DEVLOOP_BUDGET_P99 * DEVLOOP_CATASTROPHIC_MULT;
assert!(
p50 <= ceil_p50,
"DEVLOOP CATASTROPHIC ({ac}): p50={p50:?} > {ceil_p50:?} (10× soft budget) — \
an orders-of-magnitude latency regression (e.g. projection-scanner throughput, \
see 53a270d), not sample noise"
);
assert!(
p99 <= ceil_p99,
"DEVLOOP CATASTROPHIC ({ac}): p99={p99:?} > {ceil_p99:?} (10× soft budget) — \
an orders-of-magnitude latency regression, not sample noise"
);
}
fn setup() -> Option<(
tempfile::TempDir,
Engine,
Vec<HeldOutQuery>,
corpus_harness::IngestReport,
&'static str,
)> {
let (fx, embedder) = devloop_fixture();
let (dir, engine) = fx.open_or_skip()?;
let report = fx.ingest_into(&engine);
engine.drain(15_000).expect("drain after ingest");
fx.assert_vec0_row_count_matches_ingest(&engine);
fx.assert_fts_index_populated(&engine);
let queries = fx.query_set(DEVLOOP_SAMPLES, DEVLOOP_SEED);
assert!(!queries.is_empty(), "devloop query_set produced no queries");
Some((dir, engine, queries, report, embedder))
}
#[test]
fn ac_013_devloop() {
let Some((_dir, engine, queries, report, embedder)) = setup() else {
return;
};
for q in &queries {
let _ = engine.search(&q.text).expect("warmup search");
}
let mut samples = Vec::with_capacity(queries.len());
for q in &queries {
let started = Instant::now();
let _ = engine.search(&q.text).expect("measure search");
samples.push(started.elapsed());
}
let p50 = percentile_ceil(&samples, 50, 100);
let p99 = percentile_ceil(&samples, 99, 100);
emit_numbers(
"013",
report.nodes,
samples.len(),
(p50, p99),
None,
&cache_field(&report, embedder),
embedder,
);
if embedder == "synthetic" {
warn_if_over_latency("013", p50, p99); enforce_catastrophic_ceiling("013", p50, p99); } else {
eprintln!(
"DEVLOOP_PERF_INFO ac=013 metric=latency disposition=report_only embedder={embedder}"
);
}
}
#[test]
fn ac_013b_devloop() {
let Some((_dir, engine, queries, report, embedder)) = setup() else {
return;
};
let gt_embedder = ground_truth_embedder(embedder == "real");
let db_path = engine.path().to_path_buf();
let conn = rusqlite::Connection::open(&db_path).expect("raw ground-truth conn");
conn.pragma_update(None, "query_only", "ON").ok();
let mut total_hits = 0usize;
let mut total_queries = 0usize;
for q in &queries {
let vector = gt_embedder.embed(&q.text).expect("embed gt");
let vector_json = serde_json::to_string(&vector).expect("json");
let mut gt_rowid_stmt = conn
.prepare(
"SELECT rowid FROM vector_default WHERE embedding MATCH vec_f32(?1) \
ORDER BY distance LIMIT 10",
)
.expect("prepare gt rowid");
let gt_rowids: Vec<i64> = gt_rowid_stmt
.query_map([&vector_json], |row| row.get::<_, i64>(0))
.expect("gt rowid query")
.filter_map(Result::ok)
.collect();
let mut body_stmt = conn
.prepare("SELECT body FROM canonical_nodes WHERE write_cursor = ?1 LIMIT 1")
.expect("prepare body");
let mut gt_bodies = Vec::with_capacity(gt_rowids.len());
for rowid in >_rowids {
if let Ok(body) = body_stmt.query_row([rowid], |row| row.get::<_, String>(0)) {
gt_bodies.push(body);
}
}
let prod: Vec<String> = engine
.search(&q.text)
.expect("measure search")
.results
.iter()
.map(|h| h.body.clone())
.collect();
let gt_set: std::collections::HashSet<&String> = gt_bodies.iter().collect();
total_hits += prod.iter().filter(|b| gt_set.contains(b)).count();
total_queries += 1;
}
let recall = total_hits as f64 / (10.0 * total_queries.max(1) as f64);
emit_numbers(
"013b",
report.nodes,
total_queries,
(Duration::ZERO, Duration::ZERO),
Some(recall),
&cache_field(&report, embedder),
embedder,
);
if embedder == "real" {
warn_if_under_recall("013b", recall);
} else {
eprintln!("DEVLOOP_PERF_INFO ac=013b metric=recall_at_10 disposition=report_only embedder=synthetic");
}
}
#[test]
fn ac_019_devloop() {
let Some((_dir, engine, queries, report, embedder)) = setup() else {
return;
};
for q in &queries {
let _ = engine.search(&q.text).expect("baseline warmup");
}
let mut baseline = Vec::with_capacity(queries.len());
for q in &queries {
let started = Instant::now();
let _ = engine.search(&q.text).expect("baseline measure");
baseline.push(started.elapsed());
}
let baseline_p99 = percentile_ceil(&baseline, 99, 100);
let query_texts: Arc<Vec<String>> = Arc::new(queries.iter().map(|q| q.text.clone()).collect());
let engine = Arc::new(engine);
let barrier = Arc::new(Barrier::new(DEVLOOP_AC019_THREADS + 1));
let sink: Arc<Mutex<Vec<Duration>>> = Arc::new(Mutex::new(Vec::new()));
let mut handles = Vec::with_capacity(DEVLOOP_AC019_THREADS);
for tid in 0..DEVLOOP_AC019_THREADS {
let engine = Arc::clone(&engine);
let barrier = Arc::clone(&barrier);
let qs = Arc::clone(&query_texts);
let sink = Arc::clone(&sink);
handles.push(thread::spawn(move || {
let mut local = Vec::with_capacity(DEVLOOP_AC019_QUERIES_PER_THREAD);
let base = tid * DEVLOOP_AC019_QUERIES_PER_THREAD;
barrier.wait();
for i in 0..DEVLOOP_AC019_QUERIES_PER_THREAD {
let q = &qs[(base + i) % qs.len()];
let started = Instant::now();
let _ = engine.search(q).expect("stress search");
local.push(started.elapsed());
}
sink.lock().unwrap().extend(local);
}));
}
let stress_started = Instant::now();
barrier.wait();
for h in handles {
h.join().expect("stress thread");
}
let stress_elapsed = stress_started.elapsed();
let stress_samples = std::mem::take(&mut *sink.lock().unwrap());
let n_stress = stress_samples.len();
let stress_p50 = percentile_ceil(&stress_samples, 50, 100);
let stress_p99 = percentile_ceil(&stress_samples, 99, 100);
emit_numbers(
"019",
report.nodes,
n_stress,
(stress_p50, stress_p99),
None,
&cache_field(&report, embedder),
embedder,
);
eprintln!(
"DEVLOOP_AC019_DETAIL ac=019 threads={threads} per_thread={per} stress_ms={se} \
baseline_p99_ms={bp} stress_p50_ms={sp50} stress_p99_ms={sp99} disposition=report_only",
threads = DEVLOOP_AC019_THREADS,
per = DEVLOOP_AC019_QUERIES_PER_THREAD,
se = stress_elapsed.as_millis(),
bp = baseline_p99.as_millis(),
sp50 = stress_p50.as_millis(),
sp99 = stress_p99.as_millis(),
);
}