use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Mutex;
use rete_core::{
build_pyramid_meta, eval_query, validate_shacl, write_file, CountingReader, DataGraph,
DictionaryBuilder, GraphIndexBuilder, QueryOutput, RangeReader, Rete, ReteGraph, ShaclShapes,
SliceReader, SummaryView, DEFAULT_TILE_BUDGET,
};
struct RecordingReader {
data: Vec<u8>,
reads: Mutex<Vec<(u64, u64)>>,
fail: AtomicBool,
}
impl RecordingReader {
fn new(data: Vec<u8>) -> Self {
Self {
data,
reads: Mutex::new(Vec::new()),
fail: AtomicBool::new(false),
}
}
fn reads(&self) -> Vec<(u64, u64)> {
self.reads.lock().unwrap().clone()
}
fn bytes_read(&self) -> u64 {
self.reads.lock().unwrap().iter().map(|&(_, l)| l).sum()
}
fn fail_from_now(&self) {
self.fail.store(true, Ordering::Relaxed);
}
fn recover(&self) {
self.fail.store(false, Ordering::Relaxed);
}
}
impl RangeReader for RecordingReader {
fn len(&self) -> u64 {
self.data.len() as u64
}
fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
if self.fail.load(Ordering::Relaxed) {
return Err(std::io::Error::other("simulated network failure"));
}
self.reads.lock().unwrap().push((offset, len));
let start = offset as usize;
let end = start
.checked_add(len as usize)
.filter(|&e| e <= self.data.len())
.ok_or_else(|| std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "oob"))?;
Ok(self.data[start..end].to_vec())
}
}
fn image_with_pyramid() -> Vec<u8> {
let node = |n: u32| format!("<http://ex/n{n}>");
let knows = "<http://ex/knows>".to_string();
let mut edges: Vec<(u32, u32)> = Vec::new();
for c in 0..3u32 {
let base = c * 10;
for i in 0..10u32 {
for j in 0..10u32 {
if i != j {
edges.push((base + i, base + j));
}
}
}
}
edges.push((0, 10)); edges.push((10, 20));
let mut db = DictionaryBuilder::new();
for &(s, o) in &edges {
db.observe(&node(s), &knows, &node(o));
}
let dict = db.build();
let ids: Vec<_> = edges
.iter()
.map(|&(s, o)| dict.encode(&node(s), &knows, &node(o)).unwrap())
.collect();
let mut ib = GraphIndexBuilder::new();
for &t in &ids {
ib.push(t);
}
let (meta, levels) = build_pyramid_meta(&dict, &ids, DEFAULT_TILE_BUDGET);
write_file(&dict, &ib.build(), false, &meta, levels)
}
fn index_region(image: &[u8]) -> (u64, u64) {
let h = Rete::open(image).unwrap().header().clone();
(h.root_dir_offset, h.root_dir_offset + h.root_dir_len)
}
fn overlaps(r: (u64, u64), region: (u64, u64)) -> bool {
let (start, len) = r;
let end = start + len;
start < region.1 && region.0 < end
}
#[test]
fn summary_view_never_reads_the_index() {
let image = image_with_pyramid();
let (idx_start, idx_end) = index_region(&image);
assert!(
idx_end > idx_start,
"test graph should have a non-empty index"
);
let reader = RecordingReader::new(image.clone());
let view = SummaryView::open_ranged(&reader)
.unwrap()
.expect("has pyramid");
let reads = reader.reads();
assert_eq!(
reads.len(),
3,
"summary path should be 3 range reads, got {reads:?}"
);
for r in &reads {
assert!(
!overlaps(*r, (idx_start, idx_end)),
"summary read {r:?} overlapped the index region [{idx_start}, {idx_end})"
);
}
assert!(
reader.bytes_read() < image.len() as u64,
"summary read {} of {} bytes — expected a strict subset",
reader.bytes_read(),
image.len()
);
let full = Rete::open(&image).unwrap();
let knows_full = full.query(None, Some("<http://ex/knows>"), None).len() as u32;
assert_eq!(view.predicate_total("<http://ex/knows>"), knows_full);
}
#[test]
fn summary_works_with_index_zeroed() {
let mut image = image_with_pyramid();
let (idx_start, idx_end) = index_region(&image);
for b in &mut image[idx_start as usize..idx_end as usize] {
*b = 0;
}
let view = SummaryView::open_ranged(&SliceReader::new(&image))
.unwrap()
.expect("has pyramid");
let intact = Rete::open(&image_with_pyramid()).unwrap();
let knows_intact = intact.query(None, Some("<http://ex/knows>"), None).len() as u32;
assert_eq!(view.predicate_total("<http://ex/knows>"), knows_intact);
assert!(knows_intact > 0);
}
#[test]
fn full_open_is_bounded_not_a_scan() {
let image = image_with_pyramid();
let reader = RecordingReader::new(image.clone());
let rete = Rete::open_ranged(&reader).unwrap();
let reads = reader.reads();
assert!(
reads.len() <= 4,
"full ranged open should be ≤4 reads, got {} ({reads:?})",
reads.len()
);
let plain = Rete::open(&image).unwrap();
assert_eq!(rete.dump(None).len(), plain.dump(None).len());
}
#[test]
fn routed_pattern_query_fetches_only_the_selected_permutation() {
let image = image_with_pyramid();
let (idx_start, idx_end) = index_region(&image);
let plain = Rete::open(&image).unwrap();
let expected = plain.query(Some("<http://ex/n0>"), Some("<http://ex/knows>"), None);
assert!(!expected.is_empty());
let full_reader = RecordingReader::new(image.clone());
let _ = Rete::open_ranged(&full_reader).unwrap();
let reader = RecordingReader::new(image.clone());
let got = Rete::query_ranged(
&reader,
Some("<http://ex/n0>"),
Some("<http://ex/knows>"),
None,
)
.unwrap();
assert_eq!(got, expected);
assert!(
reader.bytes_read() < full_reader.bytes_read(),
"routed pattern query read {} bytes; full ranged open read {} bytes",
reader.bytes_read(),
full_reader.bytes_read()
);
assert!(
reader
.reads()
.iter()
.any(|r| overlaps(*r, (idx_start, idx_end))),
"the routed query should fetch the selected index permutation"
);
assert!(
!reader.reads().contains(&(idx_start, idx_end - idx_start)),
"routed query must not fetch the whole index container"
);
}
fn mt_node(n: u32) -> String {
format!(
"<http://ex/n/{:08x}/{:08x}/{n:05}>",
n.wrapping_mul(0x9E37_79B9),
n.wrapping_mul(0x85EB_CA6B) ^ 0x5151_5151
)
}
fn multi_tile_image() -> Vec<u8> {
let node = mt_node;
let knows = "<http://ex/knows>".to_string();
let mut db = DictionaryBuilder::new();
let edges: Vec<(u32, u32)> = (0..24000u32)
.flat_map(|i| [(i, (i * 7 + 1) % 24000), (i, (i * 13 + 5) % 24000)])
.collect();
for &(s, o) in &edges {
db.observe(&node(s), &knows, &node(o));
}
let dict = db.build();
let mut ib = GraphIndexBuilder::new().with_tile_budget(256);
for &(s, o) in &edges {
ib.push(dict.encode(&node(s), &knows, &node(o)).unwrap());
}
write_file(&dict, &ib.build(), false, &[], 0)
}
#[test]
fn routed_pattern_query_fetches_only_matching_tiles() {
let image = multi_tile_image();
let (idx_start, idx_end) = index_region(&image);
let index_len = idx_end - idx_start;
let plain = Rete::open(&image).unwrap();
let n7 = mt_node(7);
let expected = plain.query(Some(&n7), None, None);
assert_eq!(expected.len(), 2);
let reader = RecordingReader::new(image.clone());
let got = Rete::query_ranged(&reader, Some(&n7), None, None).unwrap();
assert_eq!(got, expected);
let index_bytes_read: u64 = reader
.reads()
.iter()
.filter(|r| overlaps(**r, (idx_start, idx_end)))
.map(|&(_, l)| l)
.sum();
assert!(
index_bytes_read < index_len / 6,
"tile-routed query read {index_bytes_read} of {index_len} index bytes — \
expected directory + one tile, a small fraction of one section"
);
}
#[test]
fn lazy_sparql_open_fetches_only_touched_tiles() {
let image = multi_tile_image();
let (idx_start, idx_end) = index_region(&image);
let index_len = idx_end - idx_start;
let q = format!("SELECT ?o WHERE {{ {} <http://ex/knows> ?o }}", mt_node(7));
let plain = Rete::open(&image).unwrap();
let expected = match eval_query(&plain, &q).unwrap() {
QueryOutput::Select(_, rows) => rows,
other => panic!("unexpected output {other:?}"),
};
assert_eq!(expected.len(), 2);
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let got = match eval_query(&rete, &q).unwrap() {
QueryOutput::Select(_, rows) => rows,
other => panic!("unexpected output {other:?}"),
};
assert_eq!(got, expected);
assert!(!rete.index_incomplete());
let index_bytes_read: u64 = reader
.reads()
.iter()
.filter(|r| overlaps(**r, (idx_start, idx_end)))
.map(|&(_, l)| l)
.sum();
assert!(
index_bytes_read < index_len / 4,
"lazy SPARQL read {index_bytes_read} of {index_len} index bytes — \
expected tile directories plus the touched tiles only"
);
let h = plain.header().clone();
let (dict_start, dict_end) = (h.dictionary_offset, h.dictionary_offset + h.dictionary_len);
let dict_bytes_read: u64 = reader
.reads()
.iter()
.filter(|r| overlaps(**r, (dict_start, dict_end)))
.map(|&(_, l)| l)
.sum();
assert!(
dict_bytes_read < h.dictionary_len / 2,
"lazy SPARQL read {dict_bytes_read} of {} dictionary bytes — \
expected directories plus a few chunks only",
h.dictionary_len
);
}
#[test]
fn lazy_open_defers_the_pyramid_until_needed() {
let image = image_with_pyramid();
let h = Rete::open(&image).unwrap().header().clone();
let pyr = (
h.pyramid_meta_offset,
h.pyramid_meta_offset + h.pyramid_meta_len,
);
assert!(pyr.1 > pyr.0, "fixture should carry a pyramid");
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let q = "SELECT ?o WHERE { <http://ex/n0> <http://ex/knows> ?o }";
let _ = eval_query(&rete, q).unwrap();
assert!(
!reader.reads().iter().any(|r| overlaps(*r, pyr)),
"lazy SPARQL open/eval fetched the pyramid region {pyr:?}: {:?}",
reader.reads()
);
let got = rete.pyramid().expect("pyramid faults in on demand");
let want = Rete::open(&image).unwrap();
assert_eq!(got.summary.len(), want.pyramid().unwrap().summary.len());
assert!(!got.summary.is_empty(), "fixture pyramid has super-edges");
assert!(
reader.reads().iter().any(|r| overlaps(*r, pyr)),
"pyramid() should have fetched the pyramid region {pyr:?}"
);
}
#[test]
fn full_scan_coalesces_tile_fetches() {
let image = multi_tile_image();
let (idx_start, idx_end) = index_region(&image);
let plain = Rete::open(&image).unwrap();
let spo_tiles = plain.default_index().tile_sections()[0].len();
assert!(
spo_tiles > 100,
"expected a many-tile fixture, got {spo_tiles}"
);
let q = "SELECT ?s ?p ?o WHERE { ?s ?p ?o }";
let expected = match eval_query(&plain, q).unwrap() {
QueryOutput::Select(_, rows) => rows.len(),
other => panic!("unexpected output {other:?}"),
};
assert!(expected > 40_000);
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let got = match eval_query(&rete, q).unwrap() {
QueryOutput::Select(_, rows) => rows.len(),
other => panic!("unexpected output {other:?}"),
};
assert_eq!(got, expected);
assert!(!rete.index_incomplete());
let index_reads = reader
.reads()
.iter()
.filter(|r| overlaps(**r, (idx_start, idx_end)))
.count();
assert!(
index_reads < 64,
"full scan issued {index_reads} index-region reads over {spo_tiles} tiles — \
expected the six tile directories plus O(log n) coalesced batch reads"
);
}
#[test]
fn small_limit_does_not_fetch_the_whole_index() {
let image = multi_tile_image();
let (idx_start, idx_end) = index_region(&image);
let index_len = idx_end - idx_start;
let spo_tiles = Rete::open(&image).unwrap().default_index().tile_sections()[0].len();
assert!(
spo_tiles > 100,
"expected a many-tile fixture, got {spo_tiles}"
);
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let q = "SELECT ?s ?p ?o WHERE { ?s ?p ?o } LIMIT 1";
let rows = match eval_query(&rete, q).unwrap() {
QueryOutput::Select(_, rows) => rows,
other => panic!("unexpected output {other:?}"),
};
assert_eq!(rows.len(), 1);
assert!(!rete.index_incomplete());
let index_bytes_read: u64 = reader
.reads()
.iter()
.filter(|r| overlaps(**r, (idx_start, idx_end)))
.map(|&(_, l)| l)
.sum();
assert!(
index_bytes_read < index_len / 4,
"LIMIT 1 read {index_bytes_read} of {index_len} index bytes over \
{spo_tiles} tiles — expected only the first prefetch window"
);
}
#[test]
fn multi_term_output_coalesces_dictionary_faults() {
let image = multi_tile_image();
let h = Rete::open(&image).unwrap().header().clone();
let dict = (h.dictionary_offset, h.dictionary_offset + h.dictionary_len);
let q = "SELECT ?s ?o WHERE { ?s <http://ex/knows> ?o } LIMIT 400";
let plain = Rete::open(&image).unwrap();
let expected = match eval_query(&plain, q).unwrap() {
QueryOutput::Select(_, rows) => rows,
other => panic!("unexpected output {other:?}"),
};
assert_eq!(expected.len(), 400);
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let got = match eval_query(&rete, q).unwrap() {
QueryOutput::Select(_, rows) => rows,
other => panic!("unexpected output {other:?}"),
};
assert_eq!(got, expected);
assert!(!rete.index_incomplete());
let dict_reads = reader
.reads()
.iter()
.filter(|r| overlaps(**r, dict))
.count();
assert!(
dict_reads < 24,
"resolving 800 output terms issued {dict_reads} dictionary reads — \
expected a few coalesced chunk batches, not one per term/chunk"
);
}
#[test]
fn dump_over_lazy_open_coalesces_fetches() {
let image = multi_tile_image();
let plain = Rete::open(&image).unwrap();
let expected = plain.dump(None);
let reader = std::sync::Arc::new(RecordingReader::new(image.clone()));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let got = rete.dump(None);
assert_eq!(got, expected);
assert!(!rete.index_incomplete());
let total_reads = reader.reads().len();
assert!(
total_reads < 80,
"lazy dump issued {total_reads} range reads — expected the six permutation \
directories plus coalesced chunk batches and O(log n) tile-prefetch batches"
);
}
#[test]
fn lazy_sparql_open_surfaces_failed_tile_fetches() {
let image = multi_tile_image();
let reader = std::sync::Arc::new(RecordingReader::new(image));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
assert!(!rete.index_incomplete());
reader.fail_from_now();
let q = format!("SELECT ?o WHERE {{ {} <http://ex/knows> ?o }}", mt_node(7));
let _ = eval_query(&rete, &q).unwrap(); assert!(
rete.index_incomplete(),
"a failed tile fetch must set the incomplete flag"
);
}
#[test]
fn reset_load_failures_makes_failed_fetches_retryable() {
let image = multi_tile_image();
let expected = {
let healthy =
Rete::open_ranged_lazy(std::sync::Arc::new(RecordingReader::new(image.clone())))
.unwrap();
let q = format!("SELECT ?o WHERE {{ {} <http://ex/knows> ?o }}", mt_node(7));
match eval_query(&healthy, &q).unwrap() {
QueryOutput::Select(_, rows) => rows.len(),
_ => unreachable!(),
}
};
assert!(expected > 0, "fixture must produce rows");
let reader = std::sync::Arc::new(RecordingReader::new(image));
let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
reader.fail_from_now();
let q = format!("SELECT ?o WHERE {{ {} <http://ex/knows> ?o }}", mt_node(7));
let _ = eval_query(&rete, &q).unwrap();
assert!(rete.index_incomplete(), "outage query must be flagged");
reader.recover();
rete.reset_load_failures();
assert!(!rete.index_incomplete(), "reset must clear the verdict");
let rows = match eval_query(&rete, &q).unwrap() {
QueryOutput::Select(_, rows) => rows.len(),
_ => unreachable!(),
};
assert!(
!rete.index_incomplete(),
"recovered query must not re-flag — its fetches succeeded"
);
assert_eq!(
rows, expected,
"the retried query must return the COMPLETE answer, not a poisoned empty tile"
);
}
#[test]
fn routed_pattern_query_with_unknown_term_skips_the_index() {
let image = image_with_pyramid();
let (idx_start, idx_end) = index_region(&image);
let reader = RecordingReader::new(image);
let got = Rete::query_ranged(
&reader,
Some("<http://ex/missing>"),
Some("<http://ex/knows>"),
None,
)
.unwrap();
assert!(got.is_empty());
for r in reader.reads() {
assert!(
!overlaps(r, (idx_start, idx_end)),
"unknown-term query read index range {r:?}"
);
}
}
#[test]
fn shacl_over_lazy_open_matches_eager_and_fetches_only_targets() {
const TYPE: &str = "<http://www.w3.org/1999/02/22-rdf-syntax-ns#type>";
const PERSON: &str = "<http://ex/Person>";
const EMAIL: &str = "<http://ex/email>";
let mut triples: Vec<(String, String, String)> = Vec::new();
for i in 0..6 {
let p = format!("<http://ex/p{i}>");
triples.push((p.clone(), TYPE.to_string(), PERSON.to_string()));
let email = if i == 3 {
"\"bad\"".to_string() } else {
format!("\"p{i}@ex.org\"")
};
triples.push((p, EMAIL.to_string(), email));
}
for i in 0..3000u32 {
triples.push((
format!("<http://ex/n{i:05}>"),
"<http://ex/knows>".to_string(),
format!("<http://ex/n{:05}>", (i + 1) % 3000),
));
}
let mut db = DictionaryBuilder::new();
for (s, p, o) in &triples {
db.observe(s, p, o);
}
let dict = db.build();
let ids: Vec<(u32, u32, u32)> = triples
.iter()
.map(|(s, p, o)| dict.encode(s, p, o).unwrap())
.collect();
let mut ib = GraphIndexBuilder::new().with_tile_budget(2048);
for &t in &ids {
ib.push(t);
}
let (meta, levels) = build_pyramid_meta(&dict, &ids, DEFAULT_TILE_BUDGET);
let image = write_file(&dict, &ib.build(), false, &meta, levels);
let shapes = ShaclShapes::parse_turtle(
r#"
@prefix ex: <http://ex/> .
@prefix sh: <http://www.w3.org/ns/shacl#> .
ex:PersonShape a sh:NodeShape ;
sh:targetClass ex:Person ;
sh:property [ sh:path ex:email ; sh:minCount 1 ; sh:pattern "^[^@]+@[^@]+$" ] .
"#,
)
.unwrap();
let eager_rete = Rete::open(&image).unwrap();
let eager = validate_shacl(&DataGraph::from_rete(&eager_rete, None), &shapes);
let leaked: &'static [u8] = Box::leak(image.clone().into_boxed_slice());
let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
let lazy_rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
let before = reader.bytes_read();
let lazy = validate_shacl(&ReteGraph::new(&lazy_rete), &shapes);
let pulled = reader.bytes_read() - before;
assert!(!lazy_rete.index_incomplete(), "no lazy fetch failed");
assert!(!eager.conforms, "the bad email must fail the pattern");
assert_eq!(eager.conforms, lazy.conforms);
assert_eq!(eager.results.len(), lazy.results.len());
assert!(
(pulled as usize) < image.len() / 3,
"lazy SHACL pulled {pulled} B of a {} B file — expected only the type/email tiles",
image.len()
);
}
#[test]
fn lazy_bound_object_lookup_reaches_late_slices_of_a_split_group() {
let node = mt_node;
let cites = "<http://ex/cites>".to_string();
let mut db = DictionaryBuilder::new();
let edges: Vec<(u32, u32)> = (0..8000u32).map(|i| (i % 400, 10_000 + i)).collect();
for &(s, o) in &edges {
db.observe(&node(s), &cites, &node(o));
}
let dict = db.build();
let mut ib = GraphIndexBuilder::new().with_tile_budget(256);
for &(s, o) in &edges {
ib.push(dict.encode(&node(s), &cites, &node(o)).unwrap());
}
let image = write_file(&dict, &ib.build(), false, &[], 0);
let plain = Rete::open(&image).unwrap();
let lazy =
Rete::open_ranged_lazy(std::sync::Arc::new(RecordingReader::new(image.clone()))).unwrap();
for i in [0u32, 4_000, 7_999] {
let q = format!(
"SELECT ?s WHERE {{ ?s <http://ex/cites> {} }}",
node(10_000 + i)
);
let want = match eval_query(&plain, &q).unwrap() {
QueryOutput::Select(_, r) => r.len(),
_ => unreachable!(),
};
assert_eq!(want, 1, "plain open must find o=10{i}");
let got = match eval_query(&lazy, &q).unwrap() {
QueryOutput::Select(_, r) => r.len(),
_ => unreachable!(),
};
assert!(!lazy.index_incomplete());
assert_eq!(got, want, "LAZY open must find o offset {i} in its slice");
}
}