use std::sync::Arc;
use fathomdb_embedder_api::{Embedder, EmbedderError, EmbedderIdentity, Vector};
use fathomdb_engine::{Engine, PreparedWrite};
use fathomdb_schema::SQLITE_SUFFIX;
use sha2::{Digest, Sha256};
use tempfile::TempDir;
fn prior_stable_id_content(body: &str) -> String {
let hex: String = Sha256::digest(body.as_bytes()).iter().map(|b| format!("{b:02x}")).collect();
format!("h:{hex}")
}
#[derive(Clone, Debug)]
struct FixedEmbedder;
impl Embedder for FixedEmbedder {
fn identity(&self) -> EmbedderIdentity {
EmbedderIdentity::new("deterministic", "rev-a", 8)
}
fn embed(&self, _text: &str) -> Result<Vector, EmbedderError> {
let mut v = vec![0.0_f32; 8];
v[0] = 1.0;
Ok(v)
}
}
fn opened(name: &str) -> (TempDir, fathomdb_engine::OpenedEngine) {
let dir = TempDir::new().unwrap();
let path = dir.path().join(format!("{name}{SQLITE_SUFFIX}"));
let opened = Engine::open_with_embedder_for_test(&path, Arc::new(FixedEmbedder)).expect("open");
(dir, opened)
}
fn seed(engine: &Engine) {
engine.configure_vector_kind_for_test("doc").expect("vector kind");
for body in ["hybrid retrieval alpha", "hybrid retrieval beta"] {
engine
.write(&[PreparedWrite::Node {
kind: "doc".to_string(),
body: body.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("write");
}
engine.drain(10_000).expect("drain");
}
#[test]
fn telemetry_is_off_by_default() {
let (_dir, opened) = opened("tel_off");
seed(&opened.engine);
let _ = opened.engine.search("hybrid").expect("search");
assert_eq!(opened.engine.last_telemetry_query_id(), None);
assert!(
opened.engine.record_feedback("q0-0", &[1], &[], "agent:test").is_err(),
"record_feedback must error when telemetry is off"
);
opened.engine.close().unwrap();
}
#[test]
fn telemetry_captures_event_and_feedback_deterministically() {
let (dir, opened) = opened("tel_on");
seed(&opened.engine);
let sink = dir.path().join("telemetry.jsonl");
let sink_str = sink.to_str().unwrap();
opened.engine.enable_telemetry(sink_str).expect("enable");
let r0 = opened.engine.search("hybrid").expect("search");
assert!(!r0.results.is_empty(), "expected hits to capture");
assert_eq!(opened.engine.last_telemetry_query_id().as_deref(), Some("q0-0"));
let expected_stable_ids: Vec<String> = r0.results.iter().map(|h| h.id.to_prefixed()).collect();
let _ = opened.engine.search("retrieval").expect("search");
assert_eq!(opened.engine.last_telemetry_query_id().as_deref(), Some("q0-1"));
opened
.engine
.record_feedback("q0-0", &[r0.results[0].write_cursor], &[], "agent:test")
.expect("feedback");
opened.engine.close().unwrap();
let body = std::fs::read_to_string(&sink).expect("sink readable");
let lines: Vec<&str> = body.lines().collect();
assert_eq!(lines.len(), 3, "expected 2 events + 1 feedback, got {}", lines.len());
let ev0: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
assert_eq!(ev0["type"], "event");
assert_eq!(ev0["query_id"], "q0-0");
assert_eq!(ev0["schema_version"], 1);
assert_eq!(ev0["query_chars"], "hybrid".chars().count() as u64);
assert!(ev0["result_ids"].as_array().is_some_and(|a| !a.is_empty()));
assert!(ev0["arm_of"].is_object());
let result_ids = ev0["result_ids"].as_array().unwrap();
let stable_ids = ev0["result_stable_ids"].as_array().expect("result_stable_ids present");
assert_eq!(stable_ids.len(), result_ids.len(), "stable ids parallel result_ids 1:1");
assert!(
stable_ids.iter().all(|v| v.as_str().is_some_and(|s| s.starts_with("h:"))),
"doc-corpus hits carry content-hash stable ids: {stable_ids:?}"
);
let emitted: Vec<String> = stable_ids.iter().map(|v| v.as_str().unwrap().to_string()).collect();
assert_eq!(emitted, expected_stable_ids, "telemetry stable-id bytes == id.to_prefixed()");
let expected_prior: std::collections::HashSet<String> =
["hybrid retrieval alpha", "hybrid retrieval beta"]
.iter()
.map(|b| prior_stable_id_content(b))
.collect();
for s in &emitted {
assert!(
expected_prior.contains(s),
"emitted stable id {s} not a prior derive_stable_id byte-form"
);
}
let fb: serde_json::Value = serde_json::from_str(lines[2]).unwrap();
assert_eq!(fb["type"], "feedback");
assert_eq!(fb["query_id"], "q0-0");
assert_eq!(fb["label_source"], "agent:test");
assert!(!body.contains("hybrid"), "query text must NOT be captured");
assert!(!body.contains("retrieval"), "query text must NOT be captured");
assert!(!body.contains("source_id"), "source_id must NOT be captured");
}
#[test]
fn record_feedback_rejects_unissued_query_id_and_writes_nothing() {
let (dir, opened) = opened("tel_bogus");
seed(&opened.engine);
let sink = dir.path().join("telemetry.jsonl");
opened.engine.enable_telemetry(sink.to_str().unwrap()).expect("enable");
let r0 = opened.engine.search("hybrid").expect("search");
assert!(!r0.results.is_empty());
assert_eq!(opened.engine.last_telemetry_query_id().as_deref(), Some("q0-0"));
for bogus in ["hybrid", "q0-99", "q1-0", "not-an-id", ""] {
assert!(
opened
.engine
.record_feedback(bogus, &[r0.results[0].write_cursor], &[], "agent:test")
.is_err(),
"record_feedback must reject unissued query_id {bogus:?}"
);
}
opened
.engine
.record_feedback("q0-0", &[r0.results[0].write_cursor], &[], "agent:test")
.expect("issued id accepted");
opened.engine.close().unwrap();
let body = std::fs::read_to_string(&sink).expect("sink readable");
let lines: Vec<&str> = body.lines().collect();
assert_eq!(lines.len(), 2, "bogus feedback must not append: {body}");
let ev: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
assert_eq!(ev["type"], "event");
let fb: serde_json::Value = serde_json::from_str(lines[1]).unwrap();
assert_eq!(fb["type"], "feedback");
assert_eq!(fb["query_id"], "q0-0");
assert!(!body.contains("hybrid"), "rejected query text must NOT appear");
assert!(!body.contains("not-an-id"), "rejected string must NOT appear");
}