use fathomdb_engine::{ConsolidateAxis, Engine, EngineError, PreparedWrite};
use rusqlite::Connection;
use tempfile::TempDir;
use fathomdb_schema::SQLITE_SUFFIX;
fn fixture_dir() -> std::path::PathBuf {
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/slice15_consolidate")
}
fn consolidate_harness_cmd() -> Vec<String> {
let script = fixture_dir().join("stub_consolidate_harness.py");
assert!(script.exists(), "consolidate stub harness must exist at {}", script.display());
vec!["python3".to_string(), script.to_string_lossy().to_string()]
}
fn extract_only_harness_cmd() -> Vec<String> {
let script = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/slice15_byo_llm/stub_harness.py");
assert!(script.exists(), "extract stub harness must exist at {}", script.display());
vec!["python3".to_string(), script.to_string_lossy().to_string()]
}
fn db_path(dir: &TempDir, name: &str) -> std::path::PathBuf {
dir.path().join(format!("{name}{SQLITE_SUFFIX}"))
}
#[allow(clippy::too_many_arguments)]
fn fact_edge(
kind: &str,
from: &str,
to: &str,
logical_id: &str,
body: &str,
t_valid: i64,
confidence: f64,
) -> PreparedWrite {
PreparedWrite::Edge {
kind: kind.to_string(),
from: from.to_string(),
to: to.to_string(),
source_id: fathomdb_engine::SourceId::new(format!("doc-{to}")).expect("test source id"),
logical_id: Some(logical_id.to_string()),
body: Some(body.to_string()),
t_valid: Some(t_valid),
t_invalid: None,
confidence: Some(confidence),
extractor_model_id: Some("stub-extractor-v1".to_string()),
temporal_fallback: None,
}
}
fn seed_competing_edges(engine: &Engine) {
let older = fact_edge(
"works_for",
"bob",
"acme",
"edge-acme",
"Bob works for Acme",
1_546_300_800, 0.90,
);
let newer = fact_edge(
"works_for",
"bob",
"globex",
"edge-globex",
"Bob works for Globex",
1_640_995_200, 0.80,
);
engine.write(&[older, newer]).expect("seed two competing edges");
}
fn edge_fts_count(conn: &Connection, term: &str) -> u64 {
conn.query_row(
"SELECT COUNT(*) FROM search_index_edges WHERE search_index_edges MATCH ?1",
[term],
|r| r.get(0),
)
.expect("search_index_edges must exist and be queryable")
}
fn now_epoch() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
fn edge_row(conn: &Connection, logical_id: &str) -> (Option<String>, Option<i64>, Option<i64>) {
conn.query_row(
"SELECT body, t_invalid, superseded_at FROM canonical_edges WHERE logical_id = ?1",
[logical_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("edge row must exist")
}
#[test]
fn recency_consolidation_invalidates_older_edge() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "recency");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let cmd_strings = consolidate_harness_cmd();
let cmd_refs: Vec<&str> = cmd_strings.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened
.engine
.consolidate_with_provider(&cmd_refs, &axes)
.expect("consolidate_with_provider must succeed with stub harness");
assert_eq!(receipt.clusters_processed, 1, "one (subject, relation) cluster");
assert_eq!(receipt.edges_examined, 2, "two competing edges examined");
assert_eq!(receipt.edges_kept, 1, "newer edge kept");
assert_eq!(receipt.edges_invalidated, 1, "older edge invalidated");
assert_eq!(receipt.edges_superseded, 0, "no supersede verdicts in recency path");
let conn = Connection::open(&path).unwrap();
let (acme_body, acme_t_invalid, acme_superseded) = edge_row(&conn, "edge-acme");
assert_eq!(
acme_t_invalid,
Some(1_640_995_200), "older edge must be invalidated at the newer edge's t_valid (correct temporal bound). \
TC-33: the harness WIRE is still ISO-8601 — the engine renders the candidate's t_valid \
out as ISO and normalises the verdict's t_invalid back to epoch seconds on the way in."
);
assert_eq!(
acme_body.as_deref(),
Some("Bob works for Acme"),
"older edge body must be preserved verbatim (no content rewrite/merge)"
);
assert!(acme_superseded.is_none(), "invalidate is metadata-only: row is NOT tombstoned");
let (globex_body, globex_t_invalid, globex_superseded) = edge_row(&conn, "edge-globex");
assert!(globex_t_invalid.is_none(), "newer edge must stay live (t_invalid NULL)");
assert!(globex_superseded.is_none(), "newer edge must stay live (not superseded)");
assert_eq!(globex_body.as_deref(), Some("Bob works for Globex"), "newer edge body preserved");
let total: u64 = conn
.query_row("SELECT COUNT(*) FROM canonical_edges WHERE kind = 'works_for'", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(total, 2, "both edge rows survive consolidation");
let live_now: u64 = conn
.query_row(
"SELECT COUNT(*) FROM canonical_edges \
WHERE kind = 'works_for' AND superseded_at IS NULL \
AND (t_invalid IS NULL OR t_invalid > ?1)",
[now_epoch()],
|r| r.get(0),
)
.unwrap();
assert_eq!(live_now, 1, "exactly the newer edge is temporally live after consolidation");
}
#[test]
fn invalidate_retains_projection_terminal_while_pruning_fts() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "phantom");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let acme_cursor: i64 = {
let conn = Connection::open(&path).unwrap();
let cursor: i64 = conn
.query_row(
"SELECT write_cursor FROM canonical_edges WHERE logical_id = 'edge-acme'",
[],
|r| r.get(0),
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO _fathomdb_projection_terminal(write_cursor, state) \
VALUES(?1, 'up_to_date')",
[cursor],
)
.unwrap();
assert!(edge_fts_count(&conn, "Acme") >= 1, "setup: older edge FTS shadow present");
cursor
};
let cmd_strings = consolidate_harness_cmd();
let cmd_refs: Vec<&str> = cmd_strings.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened.engine.consolidate_with_provider(&cmd_refs, &axes).expect("consolidate");
assert_eq!(receipt.edges_invalidated, 1);
let conn = Connection::open(&path).unwrap();
let terminal_after: u64 = conn
.query_row(
"SELECT COUNT(*) FROM _fathomdb_projection_terminal WHERE write_cursor = ?1",
[acme_cursor],
|r| r.get(0),
)
.unwrap();
assert_eq!(
terminal_after, 1,
"invalidate must RETAIN the projection terminal (non-superseded)"
);
let fts_after: u64 = conn
.query_row(
"SELECT COUNT(*) FROM search_index_edges WHERE write_cursor = ?1",
[acme_cursor],
|r| r.get(0),
)
.unwrap();
assert_eq!(
fts_after, 0,
"invalidated edge's FTS shadow must be pruned (hidden from retrieval)"
);
}
#[test]
fn supersede_verdict_marks_superseded_row_survives() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "supersede");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "stub-consolidate-v1", "supported_tasks": ["consolidate"],
"max_docs_per_request": 8}), flush=True)
elif msg.get("type") == "consolidate":
edges = msg.get("cluster", {}).get("edges", [])
verdicts = []
for e in edges:
ref = e.get("edge_ref")
if ref == "edge-acme":
verdicts.append({"edge_ref": ref, "verdict": "supersede", "by": "edge-globex"})
else:
verdicts.append({"edge_ref": ref, "verdict": "keep"})
print(json.dumps({"protocol": P, "type": "result",
"request_id": msg.get("request_id"), "verdicts": verdicts}), flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened
.engine
.consolidate_with_provider(&cmd_refs, &axes)
.expect("consolidate must succeed");
assert_eq!(receipt.edges_superseded, 1, "one supersede verdict applied");
assert_eq!(receipt.edges_kept, 1, "one keep verdict");
let conn = Connection::open(&path).unwrap();
let (acme_body, _acme_t_invalid, acme_superseded) = edge_row(&conn, "edge-acme");
assert!(acme_superseded.is_some(), "superseded edge must have a non-null superseded_at");
assert_eq!(
acme_body.as_deref(),
Some("Bob works for Acme"),
"superseded edge body must be preserved (invalidate-not-delete, no rewrite)"
);
let (_g_body, _g_ti, globex_superseded) = edge_row(&conn, "edge-globex");
assert!(globex_superseded.is_none(), "kept edge must remain active");
let total: u64 = conn
.query_row("SELECT COUNT(*) FROM canonical_edges WHERE kind = 'works_for'", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(total, 2, "both rows survive supersession (no delete)");
}
#[test]
fn consolidate_refused_by_extract_only_harness() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "refused");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let cmd_strings = extract_only_harness_cmd();
let cmd_refs: Vec<&str> = cmd_strings.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let result = opened.engine.consolidate_with_provider(&cmd_refs, &axes);
assert!(
matches!(result, Err(EngineError::Consolidator)),
"an extract-only harness must refuse the consolidate task, got {result:?}"
);
}
#[test]
fn consolidate_refused_when_task_not_advertised() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "not_advertised");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "bad", "supported_tasks": ["extract"],
"max_docs_per_request": 8}), flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let result = opened.engine.consolidate_with_provider(&cmd_refs, &axes);
assert!(
matches!(result, Err(EngineError::Consolidator)),
"a harness that does not advertise 'consolidate' must be refused, got {result:?}"
);
}
#[test]
fn consolidate_rejects_out_of_cluster_verdict() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "out_of_cluster");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "stub", "supported_tasks": ["consolidate"],
"max_docs_per_request": 8}), flush=True)
elif msg.get("type") == "consolidate":
print(json.dumps({"protocol": P, "type": "result",
"request_id": msg.get("request_id"),
"verdicts": [{"edge_ref": "edge-not-in-cluster", "verdict": "keep"}]}),
flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let result = opened.engine.consolidate_with_provider(&cmd_refs, &axes);
assert!(
matches!(result, Err(EngineError::Consolidator)),
"an out-of-cluster verdict must return Err(Consolidator), got {result:?}"
);
let conn = Connection::open(&path).unwrap();
let (_b, acme_ti, acme_sup) = edge_row(&conn, "edge-acme");
assert!(
acme_ti.is_none() && acme_sup.is_none(),
"no metadata change on a rejected verdict batch"
);
}
#[test]
fn footprint_no_network_egress() {
let cmd = consolidate_harness_cmd();
assert!(cmd.len() >= 2, "command must have at least 2 parts");
let script_path = std::path::Path::new(&cmd[1]);
assert!(script_path.is_absolute(), "consolidate harness script must be an absolute path");
assert!(script_path.exists(), "consolidate harness script must exist locally");
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "no_network");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt =
opened.engine.consolidate_with_provider(&cmd_refs, &axes).expect("consolidate w/o network");
assert_eq!(receipt.clusters_processed, 1);
}
#[test]
fn empty_cluster_is_skipped() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "empty");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let cmd_strings = consolidate_harness_cmd();
let cmd_refs: Vec<&str> = cmd_strings.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "nobody".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened
.engine
.consolidate_with_provider(&cmd_refs, &axes)
.expect("consolidate must succeed");
assert_eq!(receipt.clusters_processed, 0, "empty cluster is skipped");
assert_eq!(receipt.edges_examined, 0);
}
#[test]
fn invalidate_hides_stale_edge_from_edge_fts() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "invalidate_fts");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
{
let conn = Connection::open(&path).unwrap();
assert!(edge_fts_count(&conn, "acme") > 0, "acme edge must be FTS-indexed pre-consolidate");
assert!(edge_fts_count(&conn, "globex") > 0, "globex edge must be FTS-indexed pre");
}
let cmd_strings = consolidate_harness_cmd();
let cmd_refs: Vec<&str> = cmd_strings.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened
.engine
.consolidate_with_provider(&cmd_refs, &axes)
.expect("consolidate must succeed");
assert_eq!(receipt.edges_invalidated, 1, "older edge invalidated");
let conn = Connection::open(&path).unwrap();
assert_eq!(
edge_fts_count(&conn, "acme"),
0,
"invalidated edge must NO LONGER appear in edge FTS (shadow pruned)"
);
assert!(edge_fts_count(&conn, "globex") > 0, "winning edge must still be FTS-searchable");
let (acme_body, acme_ti, acme_sup) = edge_row(&conn, "edge-acme");
assert_eq!(acme_body.as_deref(), Some("Bob works for Acme"), "body preserved (no rewrite)");
assert!(acme_ti.is_some(), "t_invalid recorded");
assert!(acme_sup.is_none(), "invalidate is not a tombstone");
}
#[test]
fn supersede_hides_loser_from_edge_fts() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "supersede_fts");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "stub-consolidate-v1", "supported_tasks": ["consolidate"],
"max_docs_per_request": 8}), flush=True)
elif msg.get("type") == "consolidate":
edges = msg.get("cluster", {}).get("edges", [])
verdicts = []
for e in edges:
ref = e.get("edge_ref")
if ref == "edge-acme":
verdicts.append({"edge_ref": ref, "verdict": "supersede", "by": "edge-globex"})
else:
verdicts.append({"edge_ref": ref, "verdict": "keep"})
print(json.dumps({"protocol": P, "type": "result",
"request_id": msg.get("request_id"), "verdicts": verdicts}), flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let receipt = opened
.engine
.consolidate_with_provider(&cmd_refs, &axes)
.expect("consolidate must succeed");
assert_eq!(receipt.edges_superseded, 1, "one supersede verdict applied");
let conn = Connection::open(&path).unwrap();
assert_eq!(
edge_fts_count(&conn, "acme"),
0,
"superseded edge must NO LONGER appear in edge FTS (shadow pruned)"
);
assert!(edge_fts_count(&conn, "globex") > 0, "kept edge must still be FTS-searchable");
let (acme_body, _ti, acme_sup) = edge_row(&conn, "edge-acme");
assert_eq!(acme_body.as_deref(), Some("Bob works for Acme"), "loser body preserved");
assert!(acme_sup.is_some(), "loser is superseded");
}
#[test]
fn consolidate_rejects_missing_verdict() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "missing_verdict");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "stub", "supported_tasks": ["consolidate"],
"max_docs_per_request": 8}), flush=True)
elif msg.get("type") == "consolidate":
print(json.dumps({"protocol": P, "type": "result",
"request_id": msg.get("request_id"),
"verdicts": [{"edge_ref": "edge-acme", "verdict": "keep"}]}),
flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let result = opened.engine.consolidate_with_provider(&cmd_refs, &axes);
assert!(
matches!(result, Err(EngineError::Consolidator)),
"an incomplete verdict set must return Err(Consolidator), got {result:?}"
);
let conn = Connection::open(&path).unwrap();
let (_b, acme_ti, acme_sup) = edge_row(&conn, "edge-acme");
assert!(acme_ti.is_none() && acme_sup.is_none(), "no metadata change on a rejected batch");
let (_b2, gx_ti, gx_sup) = edge_row(&conn, "edge-globex");
assert!(gx_ti.is_none() && gx_sup.is_none(), "no metadata change on a rejected batch");
}
#[test]
fn consolidate_rejects_duplicate_verdict() {
let dir = TempDir::new().unwrap();
let path = db_path(&dir, "duplicate_verdict");
let opened = Engine::open_without_embedder_for_test(&path).expect("open");
seed_competing_edges(&opened.engine);
let harness = r#"
import json, sys
P = "fathomdb.consolidate.v1"
for line in sys.stdin:
line = line.strip()
if not line:
continue
msg = json.loads(line)
if msg.get("type") == "hello":
print(json.dumps({"protocol": P, "type": "ready", "schema_version": 1,
"model": "stub", "supported_tasks": ["consolidate"],
"max_docs_per_request": 8}), flush=True)
elif msg.get("type") == "consolidate":
print(json.dumps({"protocol": P, "type": "result",
"request_id": msg.get("request_id"),
"verdicts": [{"edge_ref": "edge-acme", "verdict": "keep"},
{"edge_ref": "edge-acme", "verdict": "keep"}]}),
flush=True)
"#;
let cmd = ["python3".to_string(), "-c".to_string(), harness.to_string()];
let cmd_refs: Vec<&str> = cmd.iter().map(|s| s.as_str()).collect();
let axes = vec![ConsolidateAxis {
subject_logical_id: "bob".to_string(),
relation: "works_for".to_string(),
}];
let result = opened.engine.consolidate_with_provider(&cmd_refs, &axes);
assert!(
matches!(result, Err(EngineError::Consolidator)),
"a duplicate edge_ref must return Err(Consolidator), got {result:?}"
);
let conn = Connection::open(&path).unwrap();
let (_b, acme_ti, acme_sup) = edge_row(&conn, "edge-acme");
assert!(acme_ti.is_none() && acme_sup.is_none(), "no metadata change on a rejected batch");
}