1use crate::dictionary::Dictionary;
18use crate::header::{
19 Header, FLAG_HAS_QUADS, FLAG_HAS_QUOTED_TRIPLES, FLAG_TILE_SYNOPSIS, HEADER_LEN, MAGIC,
20};
21use crate::index::{GraphIndex, IndexPermutation, Pattern, PermSet, NUM_PERMS};
22use crate::meta::{ClassNode, CommunityDescriptor, LevelLinks, LevelRollup, PyramidMeta};
23use crate::pyramid::{build_dendrogram, project_graph, PyramidAlgo};
24use crate::reader::RangeReader;
25use crate::tiling::{choose_round_for_budget, summarize, SuperEdge};
26use crate::triples::Triple;
27use crate::varint::{read_uvarint, write_uvarint};
28
29pub const DEFAULT_TILE_BUDGET: usize = 64 * 1024;
31
32pub fn build_pyramid_meta(
41 dict: &Dictionary,
42 triples: &[(u32, u32, u32)],
43 budget: usize,
44) -> (Vec<u8>, u16) {
45 build_pyramid_meta_with(dict, triples, budget, None)
46}
47
48pub fn build_pyramid_meta_with(
52 dict: &Dictionary,
53 triples: &[(u32, u32, u32)],
54 budget: usize,
55 type_override: Option<&str>,
56) -> (Vec<u8>, u16) {
57 build_pyramid_meta_algo(dict, triples, budget, type_override, PyramidAlgo::Louvain)
58}
59
60pub fn build_pyramid_meta_algo(
67 dict: &Dictionary,
68 triples: &[(u32, u32, u32)],
69 budget: usize,
70 type_override: Option<&str>,
71 algo: PyramidAlgo,
72) -> (Vec<u8>, u16) {
73 let timing = std::env::var_os("RETE_BUILD_TIMING").is_some();
79 let mut t = timing.then(std::time::Instant::now);
80 let mut lap = |label: &str| {
81 if let Some(t0) = &mut t {
82 eprintln!(
83 " [pyramid] {label}: {:.0} ms",
84 t0.elapsed().as_secs_f64() * 1000.0
85 );
86 *t0 = std::time::Instant::now();
87 }
88 };
89
90 let louvain = |lap: &mut dyn FnMut(&str)| {
92 let g = project_graph(dict, triples);
93 lap("project_graph");
94 let d = build_dendrogram(&g);
95 lap("build_dendrogram (Louvain)");
96 d
97 };
98 let dend = match algo {
99 PyramidAlgo::Louvain => louvain(&mut lap),
100 PyramidAlgo::Types => {
101 match crate::schema_pyramid::build_type_dendrogram(dict, triples, type_override) {
102 Some(d) => {
103 lap("build_type_dendrogram");
104 d
105 }
106 None => {
107 eprintln!(
108 " [pyramid] --pyramid-algo types: no usable rdf:type \
109 predicate — falling back to louvain"
110 );
111 louvain(&mut lap)
112 }
113 }
114 }
115 };
116 let round = choose_round_for_budget(dict, triples, &dend, budget);
117 lap("choose_round_for_budget");
118 let summary = summarize(dict, triples, &dend, round);
119 lap("summarize");
120 let sp = crate::schema_pyramid::build_schema_pyramid_with(
125 dict,
126 triples,
127 &dend,
128 round,
129 type_override,
130 );
131 lap("build_schema_pyramid");
132 let predicate_stats = compute_predicate_stats(triples);
133 lap("compute_predicate_stats");
134 let char_sets = compute_char_sets(triples);
135 lap("compute_char_sets");
136 let label_index = compute_label_index(dict, triples);
137 lap("compute_label_index");
138 let meta = PyramidMeta::new(round as u32, summary, &[])
139 .with_schema(
140 sp.class_hierarchy,
141 sp.level_rollups,
142 sp.level_links,
143 sp.descriptors,
144 sp.subclass_cycles,
145 sp.disjoint_pairs,
146 sp.equivalent_pairs,
147 )
148 .with_predicate_stats(predicate_stats)
149 .with_char_sets(char_sets)
150 .with_label_index(label_index);
151 let out = (meta.encode(), dend.rounds() as u16);
152 lap("encode");
153 out
154}
155
156const LABEL_PREDICATES: &[&str] = &[
160 "<http://www.w3.org/2000/01/rdf-schema#label>",
161 "<http://www.w3.org/2004/02/skos/core#prefLabel>",
162 "<http://www.w3.org/2004/02/skos/core#altLabel>",
163 "<http://xmlns.com/foaf/0.1/name>",
164 "<http://purl.org/dc/terms/title>",
165 "<http://purl.org/dc/elements/1.1/title>",
166 "<http://schema.org/name>",
167];
168
169fn compute_label_index(
176 dict: &Dictionary,
177 triples: &[(u32, u32, u32)],
178) -> Vec<crate::meta::LabelEntry> {
179 use crate::terms::{is_literal, literal_lexical};
180 use std::collections::{HashMap, HashSet};
181 const MAX_LABELS: usize = 8192;
182
183 let label_pids: HashSet<u32> = LABEL_PREDICATES
185 .iter()
186 .filter_map(|p| dict.predicate_id(p))
187 .collect();
188 if label_pids.is_empty() {
189 return Vec::new();
190 }
191 let mut degree: HashMap<u32, u32> = HashMap::new();
193 for &(s, _p, _o) in triples {
194 *degree.entry(s).or_insert(0) += 1;
195 }
196 let mut seen: HashSet<(u32, String)> = HashSet::new();
198 let mut candidates: Vec<(u32, String, u32)> = Vec::new(); for &(s, p, o) in triples {
200 if !label_pids.contains(&p) {
201 continue;
202 }
203 let Some(term) = dict.object_term(o) else {
204 continue;
205 };
206 if !is_literal(&term) {
207 continue;
208 }
209 let Some(label) = literal_lexical(&term) else {
210 continue;
211 };
212 if label.is_empty() {
213 continue;
214 }
215 if seen.insert((s, label.to_lowercase())) {
216 candidates.push((*degree.get(&s).unwrap_or(&0), label, s));
217 }
218 }
219 if candidates.len() > MAX_LABELS {
222 candidates.sort_by(|a, b| {
223 b.0.cmp(&a.0)
224 .then_with(|| a.2.cmp(&b.2))
225 .then_with(|| a.1.cmp(&b.1))
226 });
227 candidates.truncate(MAX_LABELS);
228 }
229 candidates.sort_by(|a, b| {
231 a.1.to_lowercase()
232 .cmp(&b.1.to_lowercase())
233 .then_with(|| a.1.cmp(&b.1))
234 .then_with(|| a.2.cmp(&b.2))
235 });
236 candidates
237 .into_iter()
238 .map(|(_deg, label, subject)| crate::meta::LabelEntry { label, subject })
239 .collect()
240}
241
242pub(crate) fn compute_text_index(dict: &Dictionary, triples: &[(u32, u32, u32)]) -> Vec<u8> {
248 use crate::terms::{is_literal, literal_lexical};
249 let mut b = crate::text_index::TextIndexBuilder::new();
250 for &(s, _p, o) in triples {
251 let Some(term) = dict.object_term(o) else {
252 continue;
253 };
254 if !is_literal(&term) {
255 continue;
256 }
257 if let Some(lit) = literal_lexical(&term) {
258 b.add_text(&lit, s);
259 }
260 }
261 if b.is_empty() {
262 Vec::new()
263 } else {
264 b.build(writer_codec())
265 }
266}
267
268fn compute_char_sets(triples: &[(u32, u32, u32)]) -> Vec<crate::meta::CharSet> {
273 use std::collections::{BTreeSet, HashMap};
274 const MAX_CHAR_SETS: usize = 128;
275 let mut by_subject: HashMap<u32, BTreeSet<u32>> = HashMap::new();
276 for &(s, p, _o) in triples {
277 by_subject.entry(s).or_default().insert(p);
278 }
279 let mut shapes: HashMap<Vec<u32>, u64> = HashMap::new();
280 for set in by_subject.into_values() {
281 *shapes.entry(set.into_iter().collect()).or_insert(0) += 1;
282 }
283 let mut v: Vec<crate::meta::CharSet> = shapes
284 .into_iter()
285 .map(|(predicates, subjects)| crate::meta::CharSet {
286 predicates,
287 subjects,
288 })
289 .collect();
290 v.sort_by(|a, b| {
291 b.subjects
292 .cmp(&a.subjects)
293 .then_with(|| a.predicates.cmp(&b.predicates))
294 });
295 v.truncate(MAX_CHAR_SETS);
296 v
297}
298
299fn compute_predicate_stats(triples: &[(u32, u32, u32)]) -> Vec<crate::meta::PredStat> {
305 use std::collections::HashMap;
306 #[allow(clippy::type_complexity)]
307 let mut acc: HashMap<u32, (HashMap<u32, u32>, HashMap<u32, u32>, u64)> = HashMap::new();
308 for &(s, p, o) in triples {
309 let e = acc.entry(p).or_default();
310 *e.0.entry(s).or_insert(0) += 1;
311 *e.1.entry(o).or_insert(0) += 1;
312 e.2 += 1;
313 }
314 let mut stats: Vec<crate::meta::PredStat> = acc
315 .into_iter()
316 .map(|(predicate, (subj, obj, count))| crate::meta::PredStat {
317 predicate,
318 count,
319 distinct_subjects: subj.len() as u64,
320 distinct_objects: obj.len() as u64,
321 max_objects_per_subject: subj.values().copied().max().unwrap_or(0),
322 max_subjects_per_object: obj.values().copied().max().unwrap_or(0),
323 })
324 .collect();
325 stats.sort_by_key(|p| p.predicate);
326 stats
327}
328
329pub const CODEC_NONE: u8 = 0;
331pub const CODEC_ZSTD: u8 = 1;
333#[cfg(feature = "compression")]
335const ZSTD_LEVEL: i32 = 9;
336
337#[derive(Debug, thiserror::Error)]
338#[non_exhaustive]
339pub enum FileError {
340 #[error("header: {0}")]
341 Header(#[from] crate::header::HeaderError),
342 #[error("malformed container: {0}")]
343 Container(&'static str),
344 #[error("unknown codec: {0}")]
345 UnknownCodec(u8),
346 #[error("decompression failed: {0}")]
347 Decompress(std::io::Error),
348 #[error("io: {0}")]
349 Io(#[from] std::io::Error),
350}
351
352pub(crate) fn writer_codec() -> u8 {
355 if cfg!(feature = "compression") {
356 CODEC_ZSTD
357 } else {
358 CODEC_NONE
359 }
360}
361
362fn intersect_sorted(a: &[u32], b: &[u32]) -> Vec<u32> {
365 let mut out = Vec::with_capacity(a.len().min(b.len()));
366 let (mut i, mut j) = (0, 0);
367 while i < a.len() && j < b.len() {
368 match a[i].cmp(&b[j]) {
369 std::cmp::Ordering::Less => i += 1,
370 std::cmp::Ordering::Greater => j += 1,
371 std::cmp::Ordering::Equal => {
372 out.push(a[i]);
373 i += 1;
374 j += 1;
375 }
376 }
377 }
378 out
379}
380
381pub(crate) fn compress(codec: u8, bytes: &[u8]) -> Vec<u8> {
382 match codec {
383 #[cfg(feature = "compression")]
384 CODEC_ZSTD => {
385 zstd::encode_all(bytes, ZSTD_LEVEL).expect("zstd encode is infallible in-memory")
386 }
387 _ => bytes.to_vec(),
388 }
389}
390
391pub(crate) fn decompress(codec: u8, bytes: &[u8]) -> Result<Vec<u8>, FileError> {
392 match codec {
393 CODEC_NONE => Ok(bytes.to_vec()),
394 #[cfg(feature = "compression")]
403 CODEC_ZSTD => zstd::decode_all(bytes).map_err(FileError::Decompress),
404 #[cfg(not(feature = "compression"))]
409 CODEC_ZSTD => {
410 use std::io::Read;
411 let mut dec = ruzstd::StreamingDecoder::new(bytes)
412 .map_err(|e| FileError::Decompress(std::io::Error::other(e.to_string())))?;
413 let mut out = Vec::new();
414 dec.read_to_end(&mut out).map_err(FileError::Decompress)?;
415 Ok(out)
416 }
417 other => Err(FileError::UnknownCodec(other)),
418 }
419}
420
421const TILE_COALESCE_GAP: u64 = 4096;
425
426const DICT_COALESCE_GAP: u64 = 64 * 1024;
432
433pub const DEFAULT_EXPORT_BUDGET_MB: u64 = 4096;
440
441const MIN_DICT_CACHE_CAP: u64 = 64 << 20; #[derive(Debug, Clone, Copy)]
457pub struct MemoryBudget {
458 pub block_cache: u64,
460 pub dict_cache: u64,
462 pub tile_cache: u64,
464}
465
466pub fn split_memory_budget(total: u64) -> MemoryBudget {
477 if total == u64::MAX {
478 return MemoryBudget {
479 block_cache: u64::MAX,
480 dict_cache: u64::MAX,
481 tile_cache: u64::MAX,
482 };
483 }
484 let block_cache = (total / 8).min(crate::block_cache::DEFAULT_CACHE_CAP);
487 let tile_cache = (total / 16).min(256 << 20);
490 let dict_cache = total
493 .saturating_sub(block_cache)
494 .saturating_sub(tile_cache)
495 .max(MIN_DICT_CACHE_CAP);
496 MemoryBudget {
497 block_cache,
498 dict_cache,
499 tile_cache,
500 }
501}
502
503fn fmt_cap(cap: u64) -> String {
505 if cap == u64::MAX {
506 "unlimited".to_string()
507 } else {
508 format!("{}MiB", cap >> 20)
509 }
510}
511
512fn read_coalesced<R: RangeReader + ?Sized>(
518 reader: &R,
519 ranges: &[ByteRange],
520 gap: u64,
521) -> Option<Vec<Vec<u8>>> {
522 let mut spans: Vec<(u64, u64)> = Vec::new();
525 let mut span_of: Vec<usize> = Vec::with_capacity(ranges.len());
526 let mut i = 0;
527 while i < ranges.len() {
528 let start = ranges[i].offset;
529 let mut end = ranges[i].offset.checked_add(ranges[i].len)?;
530 let mut j = i + 1;
531 while j < ranges.len() {
532 let r = &ranges[j];
533 if r.offset < end || r.offset - end > gap {
534 break;
535 }
536 end = r.offset.checked_add(r.len)?;
537 j += 1;
538 }
539 let si = spans.len();
540 spans.push((start, end - start));
541 for _ in i..j {
542 span_of.push(si);
543 }
544 i = j;
545 }
546 let blobs = reader.read_many(&spans).ok()?;
547 if blobs.len() != spans.len() {
548 return None;
549 }
550 let mut out = Vec::with_capacity(ranges.len());
551 for (k, r) in ranges.iter().enumerate() {
552 let (span_start, _) = spans[span_of[k]];
553 let blob = &blobs[span_of[k]];
554 let lo = (r.offset - span_start) as usize;
555 let hi = lo.checked_add(r.len as usize)?;
556 out.push(blob.get(lo..hi)?.to_vec());
557 }
558 Some(out)
559}
560
561fn content_hash(parts: &[&[u8]]) -> [u8; 16] {
564 let mut h = blake3::Hasher::new();
565 for p in parts {
566 h.update(p);
567 }
568 let mut out = [0u8; 16];
569 out.copy_from_slice(&h.finalize().as_bytes()[..16]);
570 out
571}
572
573pub type TermTriple = (String, String, String);
575
576#[derive(Debug, Clone, Copy, PartialEq, Eq)]
579pub struct DumpPlan {
580 pub scan: Option<crate::index::ScanPlan>,
584 pub dictionary_bytes: u64,
593}
594
595#[derive(Debug, Clone)]
599pub struct LayoutSegment {
600 pub kind: &'static str,
601 pub label: String,
602 pub offset: u64,
603 pub len: u64,
604}
605
606#[derive(Debug, Clone, Copy, PartialEq, Eq)]
608pub struct ByteRange {
609 pub offset: u64,
610 pub len: u64,
611}
612
613impl ByteRange {
614 pub fn end(self) -> u64 {
615 self.offset + self.len
616 }
617}
618
619#[derive(Debug, Clone, PartialEq, Eq)]
621pub struct TripleProvenance {
622 pub terms: TermTriple,
624 pub ids: Triple,
626 pub graph: Option<String>,
628 pub matched_pattern: Pattern,
630 pub index_permutation: IndexPermutation,
632 pub dictionary_range: ByteRange,
634 pub index_range: ByteRange,
636 pub index_section_range: ByteRange,
639 pub pyramid_range: Option<ByteRange>,
641 pub tile: Option<String>,
645 pub tile_range: Option<ByteRange>,
648}
649
650fn encode_container(sections: &[&[u8]], codec: u8) -> Vec<u8> {
654 let mut out = Vec::new();
655 write_uvarint(&mut out, sections.len() as u64);
656 for s in sections {
657 let payload = compress(codec, s);
658 write_uvarint(&mut out, payload.len() as u64);
659 out.extend_from_slice(&payload);
660 }
661 out
662}
663
664fn decode_container(bytes: &[u8], codec: u8) -> Result<Vec<Vec<u8>>, FileError> {
666 let (n, mut pos) = read_uvarint(bytes).ok_or(FileError::Container("truncated count"))?;
667 let mut out = Vec::with_capacity((n as usize).min(bytes.len()));
671 for _ in 0..n {
672 let (len, used) =
673 read_uvarint(&bytes[pos..]).ok_or(FileError::Container("truncated length"))?;
674 pos += used;
675 let end = pos + len as usize;
676 if end > bytes.len() {
677 return Err(FileError::Container("section overruns buffer"));
678 }
679 out.push(decompress(codec, &bytes[pos..end])?);
680 pos = end;
681 }
682 Ok(out)
683}
684
685fn checked_end(off: u64, len: u64) -> Result<u64, FileError> {
686 off.checked_add(len)
687 .ok_or(FileError::Container("section range overflows"))
688}
689
690const DICT_CHUNK_BUDGET: usize = 64 * 1024;
693
694fn encode_chunked_dict_section(raw: &[u8], codec: u8) -> Vec<u8> {
711 let meta = crate::dict::parse_meta(raw).unwrap_or(crate::dict::SectionMeta {
712 term_count: 0,
713 restart_interval: 1,
714 restart_offsets: Vec::new(),
715 });
716 let body_start = meta
717 .restart_offsets
718 .first()
719 .copied()
720 .unwrap_or(raw.len() as u64);
721 let header = &raw[..(body_start.min(raw.len() as u64)) as usize];
722
723 let n_runs = meta.restart_offsets.len();
725 let mut bounds: Vec<(usize, u64, u64)> = Vec::new(); let mut r = 0;
727 while r < n_runs {
728 let start = meta.restart_offsets[r];
729 let mut r2 = r + 1;
730 while r2 < n_runs && meta.restart_offsets[r2] - start < DICT_CHUNK_BUDGET as u64 {
731 r2 += 1;
732 }
733 let end = if r2 < n_runs {
734 meta.restart_offsets[r2]
735 } else {
736 raw.len() as u64
737 };
738 bounds.push((r, start, end));
739 r = r2;
740 }
741
742 let compressed: Vec<Vec<u8>> = bounds
743 .iter()
744 .map(|&(_, s, e)| compress(codec, &raw[s as usize..e as usize]))
745 .collect();
746 let mut out = Vec::new();
747 write_uvarint(&mut out, header.len() as u64);
748 out.extend_from_slice(header);
749 write_uvarint(&mut out, bounds.len() as u64);
750 let mut prev_run = 0usize;
751 let mut prev_last: Option<Vec<u8>> = None;
755 for (i, (&(first_run, start, end), comp)) in bounds.iter().zip(&compressed).enumerate() {
756 let key = if i == 0 {
757 Vec::new()
759 } else {
760 let first = crate::dict::run_first_term(raw, start as usize).unwrap_or_default();
761 match &prev_last {
762 None => first,
765 Some(pl) => crate::dict::shortest_separator(pl, &first),
766 }
767 };
768 if i + 1 < bounds.len() {
769 let last_run_off = meta.restart_offsets[bounds[i + 1].0 - 1] as usize;
771 prev_last = crate::dict::run_last_term(raw, last_run_off, end as usize);
772 }
773 write_uvarint(&mut out, (first_run - prev_run) as u64);
774 write_uvarint(&mut out, key.len() as u64);
775 out.extend_from_slice(&key);
776 write_uvarint(&mut out, comp.len() as u64);
777 prev_run = first_run;
778 }
779 for comp in &compressed {
780 out.extend_from_slice(comp);
781 }
782 out
783}
784
785struct DictChunkEntry {
789 first_run: usize,
790 key: Vec<u8>,
791 body_start: u64,
792 start: u64,
793 end: u64,
794}
795
796fn parse_chunked_dict_dir(
800 bytes: &[u8],
801 total_len: u64,
802) -> Result<(crate::dict::SectionMeta, Vec<DictChunkEntry>), FileError> {
803 let mut pos = 0usize;
804 let take = |pos: &mut usize| -> Result<u64, FileError> {
805 let (v, n) = read_uvarint(bytes.get(*pos..).unwrap_or(&[]))
806 .ok_or(FileError::Container("truncated dict chunk directory"))?;
807 *pos += n;
808 Ok(v)
809 };
810 let header_len = take(&mut pos)? as usize;
811 let header = bytes
812 .get(pos..pos.saturating_add(header_len))
813 .ok_or(FileError::Container("truncated dict header"))?;
814 let meta = crate::dict::parse_meta(header)
815 .map_err(|_| FileError::Container("malformed dict header"))?;
816 pos += header_len;
817
818 let num_chunks = take(&mut pos)? as usize;
819 let mut entries = Vec::with_capacity(num_chunks.min(bytes.len()));
820 let mut lens = Vec::with_capacity(num_chunks.min(bytes.len()));
821 let mut prev_run = 0usize;
822 for _ in 0..num_chunks {
823 let drun = take(&mut pos)? as usize;
824 let klen = take(&mut pos)? as usize;
825 let key = bytes
826 .get(pos..pos.saturating_add(klen))
827 .ok_or(FileError::Container("truncated dict chunk key"))?
828 .to_vec();
829 pos += klen;
830 let clen = take(&mut pos)?;
831 let first_run = prev_run + drun;
832 let body_start = meta
833 .restart_offsets
834 .get(first_run)
835 .copied()
836 .ok_or(FileError::Container("dict chunk run out of range"))?;
837 entries.push(DictChunkEntry {
838 first_run,
839 key,
840 body_start,
841 start: 0,
842 end: 0,
843 });
844 lens.push(clen);
845 prev_run = first_run;
846 }
847 let mut start = pos as u64;
848 for (e, len) in entries.iter_mut().zip(lens) {
849 let end = start
850 .checked_add(len)
851 .filter(|&e| e <= total_len)
852 .ok_or(FileError::Container("dict chunk overruns section"))?;
853 e.start = start;
854 e.end = end;
855 start = end;
856 }
857 Ok((meta, entries))
858}
859
860fn read_dict_dir_ranged<R: RangeReader>(
872 reader: &R,
873 section: ByteRange,
874) -> Result<(crate::dict::SectionMeta, Vec<DictChunkEntry>), FileError> {
875 let total = section.len;
876 let init = 8192.min(total); let head = reader.read_at(section.offset, init)?;
884 let (header_len, n0) =
885 read_uvarint(&head).ok_or(FileError::Container("truncated dict header len"))?;
886 let hbase = n0; let (term_count, n1) = read_uvarint(head.get(hbase..).unwrap_or(&[]))
888 .ok_or(FileError::Container("truncated dict term_count"))?;
889 let (restart_interval, _n2) = read_uvarint(head.get(hbase + n1..).unwrap_or(&[]))
890 .ok_or(FileError::Container("truncated dict interval"))?;
891 if restart_interval == 0 {
892 return Err(FileError::Container("zero restart interval"));
893 }
894 let dir_start = (hbase as u64)
897 .checked_add(header_len)
898 .filter(|&d| d <= total)
899 .ok_or(FileError::Container("dict header overruns section"))?;
900 let dir_total = total - dir_start;
901 let meta = crate::dict::SectionMeta {
902 term_count: term_count as u32,
903 restart_interval: restart_interval as u32,
904 restart_offsets: Vec::new(),
905 };
906 let finish = |mut entries: Vec<DictChunkEntry>| {
907 for e in &mut entries {
908 e.start += dir_start; e.end += dir_start;
910 }
911 (meta.clone(), entries)
912 };
913 let mut have: Vec<u8> = Vec::new();
917 if dir_start < head.len() as u64 {
918 match parse_chunk_dir_only(&head[dir_start as usize..], dir_total)? {
919 ChunkDirParse::Done(entries) => return Ok(finish(entries)),
920 ChunkDirParse::Truncated { .. } => have = head[dir_start as usize..].to_vec(),
921 }
922 }
923 let mut want = 4096u64.min(dir_total).max(1);
936 loop {
937 let held = have.len() as u64;
938 if want > held {
939 let extra = reader.read_at(section.offset + dir_start + held, want - held)?;
940 if extra.is_empty() {
941 return Err(FileError::Container("truncated dict chunk directory"));
942 }
943 have.extend_from_slice(&extra);
944 }
945 let held = have.len() as u64;
946 match parse_chunk_dir_only(&have, dir_total)? {
947 ChunkDirParse::Done(entries) => return Ok(finish(entries)),
948 ChunkDirParse::Truncated { .. } if held >= dir_total => {
949 return Err(FileError::Container("truncated dict chunk directory"));
950 }
951 ChunkDirParse::Truncated {
952 parsed,
953 used,
954 total,
955 } => {
956 let est = if parsed >= 16 && used > 0 {
961 let whole = ((used as u64) / (parsed as u64)).saturating_mul(total as u64);
962 whole.saturating_add(whole / 8).saturating_add(64)
963 } else {
964 held.saturating_mul(2)
966 };
967 want = est
968 .max(held.saturating_add(1))
969 .min(held.saturating_mul(4))
970 .min(dir_total);
971 }
972 }
973 }
974}
975
976enum ChunkDirParse {
981 Done(Vec<DictChunkEntry>),
982 Truncated {
983 parsed: usize,
985 used: usize,
987 total: usize,
989 },
990}
991
992fn parse_chunk_dir_only(dir: &[u8], dir_total: u64) -> Result<ChunkDirParse, FileError> {
1002 let mut pos = 0usize;
1003 let take = |pos: &mut usize| -> Option<u64> {
1004 let (v, n) = read_uvarint(dir.get(*pos..).unwrap_or(&[]))?;
1005 *pos += n;
1006 Some(v)
1007 };
1008 let Some(num_chunks) = take(&mut pos) else {
1009 return Ok(ChunkDirParse::Truncated {
1010 parsed: 0,
1011 used: 0,
1012 total: 0,
1013 });
1014 };
1015 let num_chunks = num_chunks as usize;
1016 let mut entries = Vec::with_capacity(num_chunks.min(dir.len()));
1017 let mut lens: Vec<u64> = Vec::with_capacity(num_chunks.min(dir.len()));
1018 let mut prev_run = 0usize;
1019 for _ in 0..num_chunks {
1020 let entry_start = pos;
1021 let short = ChunkDirParse::Truncated {
1022 parsed: entries.len(),
1023 used: entry_start,
1024 total: num_chunks,
1025 };
1026 let (Some(drun), Some(klen)) = (take(&mut pos), take(&mut pos)) else {
1027 return Ok(short);
1028 };
1029 let (drun, klen) = (drun as usize, klen as usize);
1030 let Some(key) = dir.get(pos..pos.saturating_add(klen)) else {
1031 return Ok(short);
1032 };
1033 let key = key.to_vec();
1034 pos += klen;
1035 let Some(clen) = take(&mut pos) else {
1036 return Ok(short);
1037 };
1038 let first_run = prev_run + drun;
1039 entries.push(DictChunkEntry {
1040 first_run,
1041 key,
1042 body_start: 0,
1043 start: 0,
1044 end: 0,
1045 });
1046 lens.push(clen);
1047 prev_run = first_run;
1048 }
1049 let mut start = pos as u64;
1050 for (e, len) in entries.iter_mut().zip(lens) {
1051 let end = start
1052 .checked_add(len)
1053 .filter(|&e| e <= dir_total)
1054 .ok_or(FileError::Container("dict chunk overruns section"))?;
1055 e.start = start;
1056 e.end = end;
1057 start = end;
1058 }
1059 Ok(ChunkDirParse::Done(entries))
1060}
1061
1062fn decode_chunked_dict_section(
1066 payload: &[u8],
1067 codec: u8,
1068 cache: std::sync::Arc<crate::chunk_cache::ChunkCache>,
1069 section_index: u8,
1070) -> Result<crate::dict::ChunkedSection, FileError> {
1071 let (meta, entries) = parse_chunked_dict_dir(payload, payload.len() as u64)?;
1072 let chunks = entries
1073 .into_iter()
1074 .map(|e| {
1075 let body = decompress(codec, &payload[e.start as usize..e.end as usize])?;
1076 Ok((
1077 crate::dict::SectionChunk::new(e.first_run, e.key, e.body_start),
1078 body,
1079 ))
1080 })
1081 .collect::<Result<Vec<_>, FileError>>()?;
1082 Ok(crate::dict::ChunkedSection::resident(
1083 meta,
1084 chunks,
1085 cache,
1086 section_index,
1087 ))
1088}
1089
1090fn decode_dictionary_container(bytes: &[u8], codec: u8) -> Result<Dictionary, FileError> {
1091 let dsecs = decode_container(bytes, CODEC_NONE)?;
1092 if dsecs.len() != 4 {
1093 return Err(FileError::Container("expected 4 dictionary sections"));
1094 }
1095 let cache = crate::chunk_cache::ChunkCache::unlimited_arc();
1098 let mut sections = Vec::with_capacity(4);
1099 for (si, sec) in dsecs.iter().enumerate() {
1100 sections.push(decode_chunked_dict_section(
1101 sec,
1102 codec,
1103 cache.clone(),
1104 si as u8,
1105 )?);
1106 }
1107 let arr: [crate::dict::ChunkedSection; 4] = sections
1108 .try_into()
1109 .map_err(|_| FileError::Container("expected 4 dictionary sections"))?;
1110 Ok(Dictionary::from_chunked_sections(arr))
1111}
1112
1113fn encode_tiled_section(index: &GraphIndex, si: usize, codec: u8) -> Vec<u8> {
1119 let tiles = index.tile_sections()[si];
1120 let n = tiles.len();
1121 #[cfg(feature = "parallel")]
1128 let compressed: Vec<Vec<u8>> = {
1129 use rayon::prelude::*;
1130 (0..n)
1131 .into_par_iter()
1132 .map(|ti| compress(codec, &index.tile_body(si, ti)))
1133 .collect()
1134 };
1135 #[cfg(not(feature = "parallel"))]
1136 let compressed: Vec<Vec<u8>> = (0..n)
1137 .map(|ti| compress(codec, &index.tile_body(si, ti)))
1138 .collect();
1139 let mut out = Vec::new();
1140 write_uvarint(&mut out, tiles.len() as u64);
1141 let mut prev_min = 0u32;
1142 for (tile, comp) in tiles.iter().zip(&compressed) {
1143 let (min_a, max_a) = tile.leading_range();
1144 write_uvarint(&mut out, (min_a - prev_min) as u64);
1145 write_uvarint(&mut out, (max_a - min_a) as u64);
1146 write_uvarint(&mut out, comp.len() as u64);
1147 prev_min = min_a;
1148 }
1149 for comp in &compressed {
1150 out.extend_from_slice(comp);
1151 }
1152 for ti in 0..n {
1160 let body = index.tile_body(si, ti);
1161 let (min_b, max_b, min_c, max_c) = match crate::triples::TripleBlock::parse(&body) {
1162 Ok(b) => {
1163 let z = b.zone();
1164 (z.min_b, z.max_b, z.min_c, z.max_c)
1165 }
1166 Err(_) => (0, u32::MAX, 0, u32::MAX),
1167 };
1168 write_uvarint(&mut out, min_b as u64);
1169 write_uvarint(&mut out, (max_b - min_b) as u64);
1170 write_uvarint(&mut out, min_c as u64);
1171 write_uvarint(&mut out, (max_c - min_c) as u64);
1172 }
1173 out
1174}
1175
1176struct TileDirEntry {
1179 min_a: u32,
1180 max_a: u32,
1181 start: u64,
1182 end: u64,
1183}
1184
1185type TileSynopsis = (u32, u32, u32, u32);
1187
1188fn parse_tile_synopsis(
1196 payload: &[u8],
1197 trailer_start: usize,
1198 num_tiles: usize,
1199) -> Option<Vec<TileSynopsis>> {
1200 let mut pos = trailer_start;
1201 let take = |pos: &mut usize| -> Option<u32> {
1202 let (v, n) = read_uvarint(payload.get(*pos..)?)?;
1203 *pos += n;
1204 u32::try_from(v).ok()
1205 };
1206 let mut out = Vec::with_capacity(num_tiles.min(payload.len()));
1207 for _ in 0..num_tiles {
1208 let min_b = take(&mut pos)?;
1209 let max_b = min_b.checked_add(take(&mut pos)?)?;
1210 let min_c = take(&mut pos)?;
1211 let max_c = min_c.checked_add(take(&mut pos)?)?;
1212 out.push((min_b, max_b, min_c, max_c));
1213 }
1214 Some(out)
1215}
1216
1217fn parse_tile_directory(bytes: &[u8], total_len: u64) -> Result<Vec<TileDirEntry>, FileError> {
1222 let mut pos = 0usize;
1223 let take = |pos: &mut usize| -> Result<u64, FileError> {
1224 let (v, n) = read_uvarint(bytes.get(*pos..).unwrap_or(&[]))
1225 .ok_or(FileError::Container("truncated tile directory"))?;
1226 *pos += n;
1227 Ok(v)
1228 };
1229 let num_tiles = take(&mut pos)? as usize;
1230 let mut entries = Vec::with_capacity(num_tiles.min(bytes.len()));
1231 let mut prev_min = 0u32;
1232 let mut lens = Vec::with_capacity(num_tiles.min(bytes.len()));
1233 for _ in 0..num_tiles {
1234 let dmin = take(&mut pos)? as u32;
1235 let span = take(&mut pos)? as u32;
1236 let len = take(&mut pos)?;
1237 let min_a = prev_min.wrapping_add(dmin);
1238 entries.push(TileDirEntry {
1239 min_a,
1240 max_a: min_a.wrapping_add(span),
1241 start: 0,
1242 end: 0,
1243 });
1244 lens.push(len);
1245 prev_min = min_a;
1246 }
1247 let mut start = pos as u64;
1248 for (e, len) in entries.iter_mut().zip(lens) {
1249 let end = start
1250 .checked_add(len)
1251 .filter(|&e| e <= total_len)
1252 .ok_or(FileError::Container("tile overruns section"))?;
1253 e.start = start;
1254 e.end = end;
1255 start = end;
1256 }
1257 Ok(entries)
1258}
1259
1260fn read_tile_directory_ranged<R: RangeReader>(
1264 reader: &R,
1265 section: ByteRange,
1266) -> Result<Vec<TileDirEntry>, FileError> {
1267 let total = section.len;
1268 let mut prefetch = 4096u64.min(total);
1269 loop {
1270 let prefix = reader.read_at(section.offset, prefetch)?;
1271 match parse_tile_directory(&prefix, total) {
1272 Ok(dir) => return Ok(dir),
1273 Err(_) if prefetch < total => prefetch = prefetch.saturating_mul(2).min(total),
1274 Err(e) => return Err(e),
1275 }
1276 }
1277}
1278
1279fn read_tile_synopsis_ranged<R: RangeReader>(
1286 reader: &R,
1287 section: ByteRange,
1288 dir: &[TileDirEntry],
1289) -> Vec<Option<TileSynopsis>> {
1290 let n = dir.len();
1291 let none = vec![None; n];
1292 let trailer_start = dir.iter().map(|e| e.end).max().unwrap_or(0);
1293 let total = section.len;
1294 if n == 0 || trailer_start >= total {
1295 return none; }
1297 let trailer_len = total - trailer_start;
1298 let Ok(bytes) = reader.read_at(section.offset + trailer_start, trailer_len) else {
1299 return none;
1300 };
1301 match parse_tile_synopsis(&bytes, 0, n) {
1302 Some(v) => v.into_iter().map(Some).collect(),
1303 None => none,
1304 }
1305}
1306
1307fn tile_file_ranges(
1311 index_bytes: &[u8],
1312 container_offset: u64,
1313 section_ranges: &[ByteRange; NUM_PERMS],
1314) -> [Vec<(u32, u32, ByteRange)>; NUM_PERMS] {
1315 let mut out: [Vec<(u32, u32, ByteRange)>; NUM_PERMS] = Default::default();
1316 for (section, range) in out.iter_mut().zip(section_ranges) {
1317 if range.len == 0 || range.offset < container_offset {
1321 continue;
1322 }
1323 let start = (range.offset - container_offset) as usize;
1324 let Some(payload) = index_bytes.get(start..start + range.len as usize) else {
1325 continue;
1326 };
1327 if let Ok(dir) = parse_tile_directory(payload, payload.len() as u64) {
1328 *section = dir
1329 .into_iter()
1330 .map(|e| {
1331 (
1332 e.min_a,
1333 e.max_a,
1334 ByteRange {
1335 offset: range.offset + e.start,
1336 len: (e.end - e.start),
1337 },
1338 )
1339 })
1340 .collect();
1341 }
1342 }
1343 out
1344}
1345
1346fn decode_tiled_section(payload: &[u8], codec: u8) -> Result<Vec<(u32, u32, Vec<u8>)>, FileError> {
1349 parse_tile_directory(payload, payload.len() as u64)?
1350 .into_iter()
1351 .map(|e| {
1352 Ok((
1353 e.min_a,
1354 e.max_a,
1355 decompress(codec, &payload[e.start as usize..e.end as usize])?,
1356 ))
1357 })
1358 .collect()
1359}
1360
1361fn decode_index_container(
1371 bytes: &[u8],
1372 codec: u8,
1373 perms: PermSet,
1374) -> Result<GraphIndex, FileError> {
1375 let mut isecs = decode_container(bytes, CODEC_NONE)?;
1376 if isecs.len() != perms.len() {
1377 return Err(FileError::Container(
1378 "index container section count does not match the header permutation mask",
1379 ));
1380 }
1381 let mut sections: [Vec<(u32, u32, Vec<u8>)>; NUM_PERMS] = Default::default();
1382 for (perm, sec) in perms.iter().zip(isecs.iter_mut()) {
1383 sections[perm.section_index()] = decode_tiled_section(sec, codec)?;
1384 }
1385 Ok(GraphIndex::from_tiles(sections, perms))
1386}
1387
1388fn container_section_payload_ranges(
1389 bytes: &[u8],
1390 container_offset: u64,
1391 expected_sections: usize,
1392) -> Result<Vec<ByteRange>, FileError> {
1393 let (section_count, mut pos) =
1394 read_uvarint(bytes).ok_or(FileError::Container("truncated count"))?;
1395 let section_count = usize::try_from(section_count)
1396 .map_err(|_| FileError::Container("section count too large"))?;
1397 if section_count != expected_sections {
1398 return Err(FileError::Container("unexpected section count"));
1399 }
1400
1401 let mut ranges = Vec::with_capacity(section_count);
1402 for _ in 0..section_count {
1403 let remaining = bytes
1404 .get(pos..)
1405 .ok_or(FileError::Container("truncated length"))?;
1406 let (payload_len, used) =
1407 read_uvarint(remaining).ok_or(FileError::Container("truncated length"))?;
1408 pos = pos
1409 .checked_add(used)
1410 .ok_or(FileError::Container("section range overflows"))?;
1411 let payload_len_usize = usize::try_from(payload_len)
1412 .map_err(|_| FileError::Container("section length too large"))?;
1413 let payload_end = pos
1414 .checked_add(payload_len_usize)
1415 .ok_or(FileError::Container("section range overflows"))?;
1416 if payload_end > bytes.len() {
1417 return Err(FileError::Container("section overruns buffer"));
1418 }
1419 ranges.push(ByteRange {
1420 offset: checked_end(container_offset, pos as u64)?,
1421 len: payload_len,
1422 });
1423 pos = payload_end;
1424 }
1425
1426 Ok(ranges)
1427}
1428
1429fn decode_index_section_ranges(
1430 bytes: &[u8],
1431 container_offset: u64,
1432 perms: PermSet,
1433) -> Result<[ByteRange; NUM_PERMS], FileError> {
1434 let ranges = container_section_payload_ranges(bytes, container_offset, perms.len())?;
1435 let mut out = [ByteRange { offset: 0, len: 0 }; NUM_PERMS];
1436 for (perm, range) in perms.iter().zip(ranges) {
1437 out[perm.section_index()] = range;
1438 }
1439 Ok(out)
1440}
1441
1442fn read_uvarint_at<R: RangeReader>(
1443 reader: &R,
1444 absolute_offset: u64,
1445 container_end: u64,
1446) -> Result<(u64, u64), FileError> {
1447 if absolute_offset >= container_end {
1448 return Err(FileError::Container("truncated container varint"));
1449 }
1450 let remaining = container_end - absolute_offset;
1451 let probe_len = remaining.min(10);
1452 let bytes = reader.read_at(absolute_offset, probe_len)?;
1453 read_uvarint(&bytes)
1454 .map(|(value, used)| (value, used as u64))
1455 .ok_or(FileError::Container("truncated container varint"))
1456}
1457
1458fn locate_container_section_ranged<R: RangeReader>(
1461 reader: &R,
1462 container_offset: u64,
1463 container_len: u64,
1464 section_index: usize,
1465 expected_sections: u64,
1466) -> Result<ByteRange, FileError> {
1467 let container_end = checked_end(container_offset, container_len)?;
1468 let (section_count, used) = read_uvarint_at(reader, container_offset, container_end)?;
1469 if section_count != expected_sections {
1470 return Err(FileError::Container("unexpected container section count"));
1471 }
1472 if section_index >= section_count as usize {
1473 return Err(FileError::Container(
1474 "container section index out of bounds",
1475 ));
1476 }
1477
1478 let mut pos = checked_end(container_offset, used)?;
1479 for i in 0..section_count as usize {
1480 let (payload_len, len_used) = read_uvarint_at(reader, pos, container_end)?;
1481 pos = checked_end(pos, len_used)?;
1482 let payload_end = checked_end(pos, payload_len)?;
1483 if payload_end > container_end {
1484 return Err(FileError::Container("section overruns buffer"));
1485 }
1486 if i == section_index {
1487 return Ok(ByteRange {
1488 offset: pos,
1489 len: payload_len,
1490 });
1491 }
1492 pos = payload_end;
1493 }
1494 Err(FileError::Container("container section not found"))
1495}
1496
1497#[allow(clippy::type_complexity)]
1507fn open_index_container_lazy(
1508 reader: &std::sync::Arc<dyn RangeReader + Send + Sync>,
1509 container: ByteRange,
1510 block_codec: u8,
1511 has_synopsis: bool,
1512 read_concurrency: usize,
1513 perms: PermSet,
1514) -> Result<
1515 (
1516 GraphIndex,
1517 [ByteRange; NUM_PERMS],
1518 [Vec<(u32, u32, ByteRange)>; NUM_PERMS],
1519 ),
1520 FileError,
1521> {
1522 let mut index_section_ranges = [ByteRange { offset: 0, len: 0 }; NUM_PERMS];
1523 let mut tile_ranges: [Vec<(u32, u32, ByteRange)>; NUM_PERMS] = Default::default();
1524 #[allow(clippy::type_complexity)]
1525 let mut directories: [Vec<(u32, u32, Option<TileSynopsis>)>; NUM_PERMS] = Default::default();
1526 for (pos, perm) in perms.iter().enumerate() {
1527 let si = perm.section_index();
1528 let section = locate_container_section_ranged(
1529 reader,
1530 container.offset,
1531 container.len,
1532 pos,
1533 perms.len() as u64,
1534 )?;
1535 index_section_ranges[si] = section;
1536 let dir = read_tile_directory_ranged(reader, section)?;
1537 let syn = if has_synopsis {
1540 read_tile_synopsis_ranged(reader, section, &dir)
1541 } else {
1542 vec![None; dir.len()]
1543 };
1544 directories[si] = dir
1545 .iter()
1546 .zip(syn)
1547 .map(|(e, s)| (e.min_a, e.max_a, s))
1548 .collect();
1549 tile_ranges[si] = dir
1550 .into_iter()
1551 .map(|e| {
1552 (
1553 e.min_a,
1554 e.max_a,
1555 ByteRange {
1556 offset: section.offset + e.start,
1557 len: (e.end - e.start),
1558 },
1559 )
1560 })
1561 .collect();
1562 }
1563
1564 let codec = block_codec;
1569 let loader_ranges = tile_ranges.clone();
1570 let loader_reader = reader.clone();
1571 let loader: crate::index::TileLoader = Box::new(move |si, ti| {
1572 let (_, _, range) = loader_ranges.get(si)?.get(ti)?;
1573 let bytes = loader_reader.read_at(range.offset, range.len).ok()?;
1574 decompress(codec, &bytes).ok()
1575 });
1576 let bulk_ranges = tile_ranges.clone();
1577 let bulk_reader = reader.clone();
1578 let bulk: crate::index::TileBulkLoader = Box::new(move |si, tis| {
1579 let section = bulk_ranges.get(si)?;
1580 let want: Option<Vec<ByteRange>> = tis
1581 .iter()
1582 .map(|&ti| section.get(ti).map(|&(_, _, r)| r))
1583 .collect();
1584 let blobs = read_coalesced(bulk_reader.as_ref(), &want?, TILE_COALESCE_GAP)?;
1585 blobs.iter().map(|b| decompress(codec, b).ok()).collect()
1586 });
1587 let mut index =
1588 GraphIndex::from_remote_directories(directories, perms, loader).with_bulk_loader(bulk);
1589 index.set_tile_lens(std::array::from_fn(|si| {
1592 tile_ranges[si]
1593 .iter()
1594 .map(|&(_, _, r)| r.len.min(u32::MAX as u64) as u32)
1595 .collect()
1596 }));
1597 index.set_read_concurrency(read_concurrency);
1601 Ok((index, index_section_ranges, tile_ranges))
1602}
1603
1604pub fn write_file(
1608 dict: &Dictionary,
1609 index: &GraphIndex,
1610 has_quads: bool,
1611 pyramid_meta: &[u8],
1612 pyramid_levels: u16,
1613) -> Vec<u8> {
1614 write_dataset(dict, index, &[], has_quads, pyramid_meta, pyramid_levels)
1615}
1616
1617pub(crate) fn encode_index_container(index: &GraphIndex, codec: u8) -> Vec<u8> {
1623 let payloads: Vec<Vec<u8>> = index
1624 .perms()
1625 .iter()
1626 .map(|perm| encode_tiled_section(index, perm.section_index(), codec))
1627 .collect();
1628 let refs: Vec<&[u8]> = payloads.iter().map(|p| p.as_slice()).collect();
1629 encode_container(&refs, CODEC_NONE)
1630}
1631
1632fn encode_named_graphs(named: &[(String, GraphIndex)], codec: u8) -> Vec<u8> {
1634 let mut out = Vec::new();
1635 write_uvarint(&mut out, named.len() as u64);
1636 for (iri, index) in named {
1637 write_uvarint(&mut out, iri.len() as u64);
1638 out.extend_from_slice(iri.as_bytes());
1639 let container = encode_index_container(index, codec);
1640 write_uvarint(&mut out, container.len() as u64);
1641 out.extend_from_slice(&container);
1642 }
1643 out
1644}
1645
1646fn decode_named_graphs(
1647 bytes: &[u8],
1648 codec: u8,
1649 perms: PermSet,
1650) -> Result<Vec<(String, GraphIndex)>, FileError> {
1651 let (n, mut pos) = read_uvarint(bytes).ok_or(FileError::Container("truncated graph count"))?;
1652 let bound = |start: usize, len: u64| -> Result<usize, FileError> {
1655 start
1656 .checked_add(len as usize)
1657 .filter(|&e| e <= bytes.len())
1658 .ok_or(FileError::Container("named-graph field overruns buffer"))
1659 };
1660 let mut out = Vec::with_capacity((n as usize).min(bytes.len()));
1661 for _ in 0..n {
1662 let (ilen, u1) = read_uvarint(bytes.get(pos..).unwrap_or(&[]))
1663 .ok_or(FileError::Container("truncated iri len"))?;
1664 pos += u1;
1665 let iend = bound(pos, ilen)?;
1666 let iri = String::from_utf8_lossy(&bytes[pos..iend]).into_owned();
1667 pos = iend;
1668 let (clen, u2) = read_uvarint(bytes.get(pos..).unwrap_or(&[]))
1669 .ok_or(FileError::Container("truncated container len"))?;
1670 pos += u2;
1671 let cend = bound(pos, clen)?;
1672 let index = decode_index_container(&bytes[pos..cend], codec, perms)?;
1673 out.push((iri, index));
1674 pos = cend;
1675 }
1676 Ok(out)
1677}
1678
1679pub fn write_dataset(
1682 dict: &Dictionary,
1683 default_index: &GraphIndex,
1684 named: &[(String, GraphIndex)],
1685 has_quads: bool,
1686 pyramid_meta: &[u8],
1687 pyramid_levels: u16,
1688) -> Vec<u8> {
1689 write_dataset_with_metadata(
1690 dict,
1691 default_index,
1692 named,
1693 has_quads,
1694 pyramid_meta,
1695 pyramid_levels,
1696 &[],
1697 &[],
1698 )
1699}
1700
1701pub(crate) fn encode_dict_container(dict: &Dictionary, codec: u8) -> Vec<u8> {
1706 let raw_sections = dict.sections();
1707 let dict_payloads: Vec<Vec<u8>> = raw_sections
1708 .iter()
1709 .map(|raw| encode_chunked_dict_section(raw, codec))
1710 .collect();
1711 encode_container(
1712 &[
1713 dict_payloads[0].as_slice(),
1714 dict_payloads[1].as_slice(),
1715 dict_payloads[2].as_slice(),
1716 dict_payloads[3].as_slice(),
1717 ],
1718 CODEC_NONE,
1719 )
1720}
1721
1722#[allow(clippy::too_many_arguments)]
1733pub fn write_dataset_with_metadata(
1734 dict: &Dictionary,
1735 default_index: &GraphIndex,
1736 named: &[(String, GraphIndex)],
1737 has_quads: bool,
1738 pyramid_meta: &[u8],
1739 pyramid_levels: u16,
1740 metadata: &[u8],
1741 text_index: &[u8],
1742) -> Vec<u8> {
1743 let codec = writer_codec();
1744 let dict_container = encode_dict_container(dict, codec);
1745 write_dataset_from_parts(
1746 &dict_container,
1747 dict.term_count() as u64,
1748 default_index,
1749 named,
1750 has_quads,
1751 dict.has_quoted_triples(),
1752 pyramid_meta,
1753 pyramid_levels,
1754 metadata,
1755 text_index,
1756 codec,
1757 )
1758}
1759
1760#[allow(clippy::too_many_arguments)]
1765pub(crate) fn write_dataset_from_parts(
1766 dict_container: &[u8],
1767 term_count: u64,
1768 default_index: &GraphIndex,
1769 named: &[(String, GraphIndex)],
1770 has_quads: bool,
1771 has_quoted_triples: bool,
1772 pyramid_meta: &[u8],
1773 pyramid_levels: u16,
1774 metadata: &[u8],
1775 text_index: &[u8],
1776 codec: u8,
1777) -> Vec<u8> {
1778 let index_container = encode_index_container(default_index, codec);
1779 let named_section = encode_named_graphs(named, codec);
1780
1781 let meta_section_len = metadata.len() as u64;
1784 let dict_offset = HEADER_LEN as u64 + meta_section_len;
1785 let dict_len = dict_container.len() as u64;
1786 let index_offset = dict_offset + dict_len;
1787 let index_len = index_container.len() as u64;
1788 let pyr_offset = index_offset + index_len;
1789 let pyr_len = pyramid_meta.len() as u64;
1790 let text_offset = pyr_offset + pyr_len;
1792 let text_len = text_index.len() as u64;
1793 let named_offset = text_offset + text_len;
1794 let named_len = if named.is_empty() {
1795 0
1796 } else {
1797 named_section.len() as u64
1798 };
1799
1800 let mut parts: Vec<&[u8]> = Vec::with_capacity(5);
1806 if meta_section_len > 0 {
1807 parts.push(metadata);
1808 }
1809 parts.push(dict_container);
1810 parts.push(&index_container);
1811 parts.push(pyramid_meta);
1812 if text_len > 0 {
1813 parts.push(text_index);
1814 }
1815 if named_len > 0 {
1816 parts.push(&named_section);
1817 }
1818
1819 let schema_meta_len = crate::meta::schema_block_len(pyramid_meta);
1822
1823 let header = Header {
1824 version: crate::header::CURRENT_FORMAT_VERSION,
1825 flags: FLAG_TILE_SYNOPSIS
1826 | if has_quads { FLAG_HAS_QUADS } else { 0 }
1827 | if has_quoted_triples {
1828 FLAG_HAS_QUOTED_TRIPLES
1829 } else {
1830 0
1831 },
1832 metadata_offset: HEADER_LEN as u64,
1833 metadata_len: meta_section_len,
1834 dictionary_offset: dict_offset,
1835 dictionary_len: dict_len,
1836 root_dir_offset: index_offset,
1837 root_dir_len: index_len,
1838 pyramid_meta_offset: if pyr_len > 0 { pyr_offset } else { 0 },
1839 pyramid_meta_len: pyr_len,
1840 dict_codec: codec,
1841 block_codec: codec,
1842 pyramid_levels,
1843 perms: default_index.perms(),
1847 quad_count: default_index.triple_count() as u64
1848 + named
1849 .iter()
1850 .map(|(_, idx)| idx.triple_count() as u64)
1851 .sum::<u64>(),
1852 term_count,
1853 content_hash: content_hash(&parts),
1854 named_graphs_offset: if named_len > 0 { named_offset } else { 0 },
1855 named_graphs_len: named_len,
1856 schema_meta_len,
1857 text_index_offset: if text_len > 0 { text_offset } else { 0 },
1858 text_index_len: text_len,
1859 build_info_offset: 0,
1860 build_info_len: 0,
1861 extra_sections: Vec::new(),
1862 };
1863
1864 let mut out = Vec::with_capacity(
1865 HEADER_LEN
1866 + metadata.len()
1867 + dict_container.len()
1868 + index_container.len()
1869 + pyramid_meta.len()
1870 + text_index.len()
1871 + named_section.len()
1872 + MAGIC.len(),
1873 );
1874 out.extend_from_slice(&header.to_bytes());
1875 if meta_section_len > 0 {
1876 out.extend_from_slice(metadata);
1877 }
1878 out.extend_from_slice(dict_container);
1879 out.extend_from_slice(&index_container);
1880 out.extend_from_slice(pyramid_meta);
1881 if text_len > 0 {
1882 out.extend_from_slice(text_index);
1883 }
1884 if named_len > 0 {
1885 out.extend_from_slice(&named_section);
1886 }
1887 out.extend_from_slice(&MAGIC); out
1889}
1890
1891pub const RDF_TYPE: &str = "<http://www.w3.org/1999/02/22-rdf-syntax-ns#type>";
1893
1894pub fn schema_summary(rete: &Rete) -> Vec<(String, String, String, u32)> {
1901 use std::collections::{BTreeMap, HashMap};
1902 let triples = rete.dump(None);
1903
1904 let mut class_of: HashMap<&str, &str> = HashMap::new();
1905 for (s, p, o) in &triples {
1906 if p == RDF_TYPE {
1907 class_of.insert(s.as_str(), o.as_str());
1908 }
1909 }
1910 let classify = |t: &str| -> String {
1911 if let Some(c) = class_of.get(t) {
1912 (*c).to_string()
1913 } else if t.starts_with('"') {
1914 "(literal)".to_string()
1915 } else {
1916 "(untyped)".to_string()
1917 }
1918 };
1919
1920 let mut counts: BTreeMap<(String, String, String), u32> = BTreeMap::new();
1921 for (s, p, o) in &triples {
1922 if p == RDF_TYPE {
1923 continue; }
1925 *counts
1926 .entry((classify(s), p.clone(), classify(o)))
1927 .or_default() += 1;
1928 }
1929 counts
1930 .into_iter()
1931 .map(|((a, p, b), c)| (a, p, b, c))
1932 .collect()
1933}
1934
1935pub fn schema_classes(rete: &Rete) -> Vec<(String, u32)> {
1939 use std::collections::BTreeMap;
1940 let mut counts: BTreeMap<String, u32> = BTreeMap::new();
1941 for (_s, p, o) in rete.dump(None) {
1942 if p == RDF_TYPE {
1943 *counts.entry(o).or_default() += 1;
1944 }
1945 }
1946 let mut out: Vec<(String, u32)> = counts.into_iter().collect();
1947 out.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
1948 out
1949}
1950
1951pub fn read_metadata_ranged<R: RangeReader>(reader: &R) -> Result<Option<Vec<u8>>, FileError> {
1961 let head = reader.read_at(0, HEADER_LEN as u64)?;
1962 let header = Header::from_bytes(&head)?;
1963 if header.metadata_len == 0 {
1964 return Ok(None);
1965 }
1966 let bytes = reader.read_at(header.metadata_offset, header.metadata_len)?;
1967 Ok(Some(bytes))
1968}
1969
1970#[allow(clippy::type_complexity)]
1980pub fn read_card_and_build_info_ranged<R: RangeReader>(
1981 reader: &R,
1982) -> Result<(Option<Vec<u8>>, Option<Vec<u8>>), FileError> {
1983 read_card_and_build_info_with_header(reader).map(|(_, m, b)| (m, b))
1984}
1985
1986#[allow(clippy::type_complexity)]
1997pub fn read_card_and_build_info_with_header<R: RangeReader>(
1998 reader: &R,
1999) -> Result<(Header, Option<Vec<u8>>, Option<Vec<u8>>), FileError> {
2000 let head = reader.read_at(0, HEADER_LEN as u64)?;
2001 let header = Header::from_bytes(&head)?;
2002 let meta = (header.metadata_offset, header.metadata_len);
2003 let build = (header.build_info_offset, header.build_info_len);
2004 if meta.1 > 0 && build.1 > 0 && build.0 == meta.0 + meta.1 {
2005 let both = reader.read_at(meta.0, meta.1 + build.1)?;
2007 let (m, b) = both.split_at(meta.1 as usize);
2008 return Ok((header, Some(m.to_vec()), Some(b.to_vec())));
2009 }
2010 let fetch = |off: u64, len: u64| -> Result<Option<Vec<u8>>, FileError> {
2011 if len == 0 {
2012 return Ok(None);
2013 }
2014 Ok(Some(reader.read_at(off, len)?))
2015 };
2016 Ok((header, fetch(meta.0, meta.1)?, fetch(build.0, build.1)?))
2017}
2018
2019pub fn read_build_info(bytes: &[u8]) -> Result<Option<Vec<u8>>, FileError> {
2024 let header = Header::from_bytes(bytes)?;
2025 if header.build_info_len == 0 {
2026 return Ok(None);
2027 }
2028 let start = header.build_info_offset as usize;
2029 let end = start
2030 .checked_add(header.build_info_len as usize)
2031 .filter(|&e| e <= bytes.len())
2032 .ok_or(FileError::Container("build-info section overruns buffer"))?;
2033 Ok(Some(bytes[start..end].to_vec()))
2034}
2035
2036#[derive(Debug, Clone, PartialEq, Eq)]
2043pub struct BuildInfoPlan {
2044 pub header: [u8; HEADER_LEN],
2046 pub insert: u64,
2049 pub tail_start: u64,
2053 pub new_len: u64,
2055}
2056
2057pub fn plan_build_info(
2079 head: &[u8],
2080 file_len: u64,
2081 info_len: u64,
2082) -> Result<BuildInfoPlan, FileError> {
2083 let mut header = Header::from_bytes(head)?;
2084 let insert = HEADER_LEN as u64 + header.metadata_len;
2087 let old_len = header.build_info_len;
2088 if old_len > 0 && header.build_info_offset != insert {
2089 return Err(FileError::Container(
2090 "existing build-info section is not adjacent to the metadata",
2091 ));
2092 }
2093 let tail_start = insert.saturating_add(old_len);
2094 if tail_start > file_len || insert > file_len {
2095 return Err(FileError::Container("build-info splice out of bounds"));
2096 }
2097
2098 let shift = |off: &mut u64, len: u64| {
2100 if len > 0 && *off >= tail_start {
2101 *off = *off - old_len + info_len;
2102 }
2103 };
2104 shift(&mut header.dictionary_offset, header.dictionary_len);
2105 shift(&mut header.root_dir_offset, header.root_dir_len);
2106 shift(&mut header.pyramid_meta_offset, header.pyramid_meta_len);
2107 shift(&mut header.named_graphs_offset, header.named_graphs_len);
2108 shift(&mut header.text_index_offset, header.text_index_len);
2109 for s in &mut header.extra_sections {
2110 if s.length > 0 && s.offset >= tail_start {
2111 s.offset = s.offset - old_len + info_len;
2112 }
2113 }
2114 if info_len == 0 {
2115 header.build_info_offset = 0;
2116 header.build_info_len = 0;
2117 } else {
2118 header.build_info_offset = insert;
2119 header.build_info_len = info_len;
2120 }
2121
2122 Ok(BuildInfoPlan {
2123 header: header.to_bytes(),
2124 insert,
2125 tail_start,
2126 new_len: file_len - old_len + info_len,
2127 })
2128}
2129
2130pub fn attach_build_info(image: &[u8], info: &[u8]) -> Result<Vec<u8>, FileError> {
2136 let plan = plan_build_info(image, image.len() as u64, info.len() as u64)?;
2137 let mut out = Vec::with_capacity(plan.new_len as usize);
2138 out.extend_from_slice(&plan.header);
2139 out.extend_from_slice(&image[HEADER_LEN..plan.insert as usize]);
2140 out.extend_from_slice(info);
2141 out.extend_from_slice(&image[plan.tail_start as usize..]);
2142 debug_assert_eq!(out.len() as u64, plan.new_len);
2143 Ok(out)
2144}
2145
2146pub fn replace_metadata(image: &[u8], metadata: &[u8]) -> Result<Vec<u8>, FileError> {
2166 let mut header = Header::from_bytes(image)?;
2167 let insert = HEADER_LEN as u64;
2170 if header.metadata_len > 0 && header.metadata_offset != insert {
2171 return Err(FileError::Container(
2172 "metadata section does not sit immediately after the header",
2173 ));
2174 }
2175 let old_len = header.metadata_len;
2176 let tail_start = (insert + old_len) as usize;
2177 if tail_start > image.len() {
2178 return Err(FileError::Container("metadata splice out of bounds"));
2179 }
2180 let new_len = metadata.len() as u64;
2181
2182 let shift = |off: &mut u64, len: u64| {
2183 if len > 0 && *off >= insert + old_len {
2184 *off = *off - old_len + new_len;
2185 }
2186 };
2187 shift(&mut header.dictionary_offset, header.dictionary_len);
2188 shift(&mut header.root_dir_offset, header.root_dir_len);
2189 shift(&mut header.pyramid_meta_offset, header.pyramid_meta_len);
2190 shift(&mut header.named_graphs_offset, header.named_graphs_len);
2191 shift(&mut header.text_index_offset, header.text_index_len);
2192 shift(&mut header.build_info_offset, header.build_info_len);
2193 for s in &mut header.extra_sections {
2194 if s.length > 0 && s.offset >= insert + old_len {
2195 s.offset = s.offset - old_len + new_len;
2196 }
2197 }
2198 header.metadata_offset = insert;
2199 header.metadata_len = new_len;
2200
2201 let mut out = Vec::with_capacity(image.len() - old_len as usize + metadata.len());
2202 out.extend_from_slice(&header.to_bytes());
2203 out.extend_from_slice(metadata);
2204 out.extend_from_slice(&image[tail_start..]);
2205 header.content_hash = hash_of_sections(&header, &out)?;
2209 out[..HEADER_LEN].copy_from_slice(&header.to_bytes());
2210 Ok(out)
2211}
2212
2213fn hash_of_sections(header: &Header, bytes: &[u8]) -> Result<[u8; 16], FileError> {
2219 let slice = |off: u64, len: u64| -> Result<&[u8], FileError> {
2220 let end = off.saturating_add(len);
2227 bytes
2228 .get(off as usize..end as usize)
2229 .ok_or(FileError::Container("section overruns buffer"))
2230 };
2231 let d = slice(header.dictionary_offset, header.dictionary_len)?;
2232 let i = slice(header.root_dir_offset, header.root_dir_len)?;
2233 let m = if header.pyramid_meta_len > 0 {
2234 slice(header.pyramid_meta_offset, header.pyramid_meta_len)?
2235 } else {
2236 &[]
2237 };
2238 let mut parts: Vec<&[u8]> = Vec::with_capacity(6);
2239 if header.metadata_len > 0 {
2240 parts.push(slice(header.metadata_offset, header.metadata_len)?);
2241 }
2242 parts.push(d);
2243 parts.push(i);
2244 parts.push(m);
2245 if header.text_index_len > 0 {
2246 parts.push(slice(header.text_index_offset, header.text_index_len)?);
2247 }
2248 if header.named_graphs_len > 0 {
2249 parts.push(slice(header.named_graphs_offset, header.named_graphs_len)?);
2250 }
2251 Ok(content_hash(&parts))
2252}
2253
2254pub fn verify(bytes: &[u8]) -> Result<bool, FileError> {
2257 let header = Header::from_bytes(bytes)?;
2258 Ok(hash_of_sections(&header, bytes)? == header.content_hash)
2259}
2260
2261type PyramidLoader = Box<dyn Fn() -> Option<PyramidMeta> + Send + Sync>;
2263
2264enum PyramidSlot {
2270 Resident(Option<PyramidMeta>),
2271 Lazy {
2272 loader: PyramidLoader,
2273 cell: std::sync::OnceLock<Option<PyramidMeta>>,
2274 },
2275}
2276
2277type TextIndexLoader = Box<dyn Fn() -> Option<crate::text_index::TextIndex> + Send + Sync>;
2280
2281type TokenTableProbe = Box<dyn Fn() -> Option<u64> + Send + Sync>;
2284
2285enum TextIndexSlot {
2295 Resident {
2296 index: Option<crate::text_index::TextIndex>,
2297 token_table_len: Option<u64>,
2300 },
2301 Lazy {
2302 loader: TextIndexLoader,
2303 cell: std::sync::OnceLock<Option<crate::text_index::TextIndex>>,
2304 token_table: TokenTableProbe,
2306 token_table_cell: std::sync::OnceLock<Option<u64>>,
2307 },
2308}
2309
2310impl TextIndexSlot {
2311 fn index(&self) -> Option<&crate::text_index::TextIndex> {
2313 match self {
2314 TextIndexSlot::Resident { index, .. } => index.as_ref(),
2315 TextIndexSlot::Lazy { loader, cell, .. } => cell.get_or_init(loader).as_ref(),
2316 }
2317 }
2318
2319 fn token_table_len(&self) -> Option<u64> {
2322 match self {
2323 TextIndexSlot::Resident {
2324 token_table_len, ..
2325 } => *token_table_len,
2326 TextIndexSlot::Lazy {
2327 token_table,
2328 token_table_cell,
2329 ..
2330 } => *token_table_cell.get_or_init(token_table),
2331 }
2332 }
2333}
2334
2335fn resident_text_index_slot(section: Option<&[u8]>, codec: u8) -> Result<TextIndexSlot, FileError> {
2338 let Some(section) = section else {
2339 return Ok(TextIndexSlot::Resident {
2340 index: None,
2341 token_table_len: None,
2342 });
2343 };
2344 Ok(TextIndexSlot::Resident {
2345 index: Some(
2346 crate::text_index::TextIndex::from_section(section, codec)
2347 .map_err(|_| FileError::Container("malformed text index"))?,
2348 ),
2349 token_table_len: crate::text_index::TextIndex::postings_base(section).map(|b| b as u64),
2350 })
2351}
2352
2353pub fn read_text_index_token_table_len_ranged<R: RangeReader + ?Sized>(
2366 reader: &R,
2367 header: &Header,
2368) -> Option<u64> {
2369 if header.text_index_len == 0 {
2370 return None;
2371 }
2372 read_token_table_len(reader, header.text_index_offset, header.text_index_len)
2373}
2374
2375fn read_token_table_len<R: RangeReader + ?Sized>(reader: &R, off: u64, len: u64) -> Option<u64> {
2379 let head = reader.read_at(off, 10u64.min(len)).ok()?;
2380 let (ttlen, n) = crate::varint::read_uvarint(&head)?;
2381 Some((n as u64 + ttlen).min(len))
2382}
2383
2384const NAMED_GRAPH_RESIDENT_MAX: u64 = 1 << 20; const NAMED_SLAB: usize = 1024;
2395
2396const NAMED_WALK_CHUNK: u64 = 64 * 1024;
2400
2401const NAMED_WALK_CHUNK_MAX: u64 = 8 * 1024 * 1024;
2414
2415#[derive(Default)]
2417struct NamedEntry {
2418 meta: std::sync::OnceLock<(String, ByteRange)>,
2421 index: std::sync::OnceLock<Box<GraphIndex>>,
2424}
2425
2426#[derive(Default)]
2430struct NamedWalk {
2431 next: usize,
2433 pos: u64,
2437 buf: Vec<u8>,
2444 buf_off: u64,
2445 chunk: u64,
2452 big_seen: bool,
2459}
2460
2461struct LazyNamedGraphs {
2475 reader: std::sync::Arc<dyn RangeReader + Send + Sync>,
2476 section: ByteRange,
2478 codec: u8,
2479 has_synopsis: bool,
2480 read_concurrency: usize,
2481 perms: PermSet,
2484 #[allow(clippy::type_complexity)]
2488 dir: std::sync::OnceLock<(usize, Box<[std::sync::OnceLock<Box<[NamedEntry]>>]>)>,
2489 walk: std::sync::Mutex<NamedWalk>,
2490 failed: std::sync::atomic::AtomicBool,
2491 tile_cap: std::sync::atomic::AtomicU64,
2497}
2498
2499impl LazyNamedGraphs {
2500 fn new(
2501 reader: std::sync::Arc<dyn RangeReader + Send + Sync>,
2502 section: ByteRange,
2503 codec: u8,
2504 has_synopsis: bool,
2505 read_concurrency: usize,
2506 perms: PermSet,
2507 ) -> Self {
2508 LazyNamedGraphs {
2509 reader,
2510 section,
2511 codec,
2512 has_synopsis,
2513 read_concurrency,
2514 perms,
2515 dir: std::sync::OnceLock::new(),
2516 walk: std::sync::Mutex::new(NamedWalk::default()),
2517 failed: std::sync::atomic::AtomicBool::new(false),
2518 tile_cap: std::sync::atomic::AtomicU64::new(u64::MAX),
2519 }
2520 }
2521
2522 fn set_tile_cap(&self, cap: u64) {
2526 self.tile_cap
2527 .store(cap, std::sync::atomic::Ordering::Relaxed);
2528 self.for_each_opened(|g| {
2529 if g.is_lazy() {
2530 g.set_cache_cap(cap);
2531 }
2532 });
2533 }
2534
2535 fn fail(&self) {
2537 self.failed
2538 .store(true, std::sync::atomic::Ordering::Relaxed);
2539 }
2540
2541 #[allow(clippy::type_complexity)]
2543 fn directory(&self) -> Option<&(usize, Box<[std::sync::OnceLock<Box<[NamedEntry]>>]>)> {
2544 if let Some(d) = self.dir.get() {
2545 return Some(d);
2546 }
2547 let end = self.section.offset.checked_add(self.section.len)?;
2548 let (n, used) = match read_uvarint_at(&self.reader, self.section.offset, end) {
2549 Ok(v) => v,
2550 Err(_) => {
2551 self.fail();
2552 return None;
2553 }
2554 };
2555 if n > self.section.len / 2 {
2558 self.fail();
2559 return None;
2560 }
2561 let n = n as usize;
2562 let slabs = n.div_ceil(NAMED_SLAB);
2563 let table: Box<[std::sync::OnceLock<Box<[NamedEntry]>>]> =
2564 (0..slabs).map(|_| std::sync::OnceLock::new()).collect();
2565 let _ = self.dir.set((n, table));
2566 {
2567 let mut w = self.walk.lock().unwrap();
2570 if w.pos == 0 {
2571 w.pos = self.section.offset + used;
2572 }
2573 }
2574 self.dir.get()
2575 }
2576
2577 fn count(&self) -> usize {
2578 self.directory().map(|(n, _)| *n).unwrap_or(0)
2579 }
2580
2581 fn entry(&self, i: usize) -> Option<&NamedEntry> {
2583 let (n, table) = self.directory()?;
2584 if i >= *n {
2585 return None;
2586 }
2587 let slab = table[i / NAMED_SLAB].get_or_init(|| {
2588 let len = NAMED_SLAB.min(n - (i / NAMED_SLAB) * NAMED_SLAB);
2589 (0..len).map(|_| NamedEntry::default()).collect()
2590 });
2591 slab.get(i % NAMED_SLAB)
2592 }
2593
2594 fn next_read_len(w: &mut NamedWalk, pos: u64, end: u64, need: u64) -> u64 {
2602 let cur = w.chunk.max(NAMED_WALK_CHUNK);
2603 w.chunk = (cur * 2).min(NAMED_WALK_CHUNK_MAX);
2604 cur.max(need).min(end - pos)
2605 }
2606
2607 fn ensure_meta(&self, upto: usize, exhaustive: bool) -> Option<()> {
2622 let target = self.entry(upto)?;
2623 if target.meta.get().is_some() {
2624 return Some(());
2625 }
2626 let end = self.section.offset.checked_add(self.section.len)?;
2627 let mut w = self.walk.lock().unwrap();
2628 if exhaustive && !w.big_seen {
2629 w.chunk = NAMED_WALK_CHUNK_MAX;
2630 }
2631 while w.next <= upto {
2632 let pos = w.pos;
2633 if pos >= end {
2634 self.fail();
2636 return None;
2637 }
2638 let have =
2641 |b: &[u8], off: u64, need: u64| pos >= off && pos + need <= off + b.len() as u64;
2642 if !have(&w.buf, w.buf_off, 20.min(end - pos)) {
2643 let len = Self::next_read_len(&mut w, pos, end, 20.min(end - pos));
2644 w.buf = match self.reader.read_at(pos, len) {
2645 Ok(b) => b,
2646 Err(_) => {
2647 self.fail();
2648 return None;
2649 }
2650 };
2651 w.buf_off = pos;
2652 }
2653 let rel = (pos - w.buf_off) as usize;
2654 let Some((ilen, u1)) = read_uvarint(&w.buf[rel..]) else {
2655 self.fail();
2656 return None;
2657 };
2658 let header_need = u1 as u64 + ilen + 10; if pos + u1 as u64 + ilen > end {
2660 self.fail(); return None;
2662 }
2663 if !have(&w.buf, w.buf_off, header_need.min(end - pos)) {
2664 let len = Self::next_read_len(&mut w, pos, end, header_need);
2665 w.buf = match self.reader.read_at(pos, len) {
2666 Ok(b) => b,
2667 Err(_) => {
2668 self.fail();
2669 return None;
2670 }
2671 };
2672 w.buf_off = pos;
2673 }
2674 let rel = (pos - w.buf_off) as usize;
2675 let istart = rel + u1;
2676 let iend = istart + ilen as usize;
2677 let iri = String::from_utf8_lossy(&w.buf[istart..iend]).into_owned();
2678 let Some((clen, u2)) = read_uvarint(&w.buf[iend..]) else {
2679 self.fail();
2680 return None;
2681 };
2682 let cstart = pos + u1 as u64 + ilen + u2 as u64;
2683 let cend = match cstart.checked_add(clen) {
2684 Some(e) if e <= end => e,
2685 _ => {
2686 self.fail(); return None;
2688 }
2689 };
2690 if clen > NAMED_GRAPH_RESIDENT_MAX {
2691 w.big_seen = true;
2696 w.chunk = NAMED_WALK_CHUNK;
2697 }
2698 let range = ByteRange {
2699 offset: cstart,
2700 len: clen,
2701 };
2702 if let Some(e) = self.entry(w.next) {
2703 let _ = e.meta.set((iri, range));
2704 }
2705 w.pos = cend;
2706 w.next += 1;
2707 }
2708 Some(())
2709 }
2710
2711 fn name_at(&self, i: usize, exhaustive: bool) -> Option<&str> {
2714 self.ensure_meta(i, exhaustive)?;
2715 self.entry(i)?.meta.get().map(|(iri, _)| iri.as_str())
2716 }
2717
2718 fn graph_at(&self, i: usize, exhaustive: bool) -> Option<(&str, &GraphIndex)> {
2721 self.ensure_meta(i, exhaustive)?;
2722 let e = self.entry(i)?;
2723 let (iri, range) = e.meta.get()?;
2724 if e.index.get().is_none() {
2725 let opened = self.open_graph(*range)?;
2726 let _ = e.index.set(Box::new(opened));
2727 }
2728 Some((iri.as_str(), e.index.get()?.as_ref()))
2729 }
2730
2731 fn container_bytes(&self, range: ByteRange) -> Option<Vec<u8>> {
2735 {
2736 let w = self.walk.lock().unwrap();
2737 let buf_end = w.buf_off + w.buf.len() as u64;
2738 if range.offset >= w.buf_off && range.offset + range.len <= buf_end {
2739 let a = (range.offset - w.buf_off) as usize;
2740 return Some(w.buf[a..a + range.len as usize].to_vec());
2741 }
2742 }
2743 match self.reader.read_at(range.offset, range.len) {
2744 Ok(b) => Some(b),
2745 Err(_) => {
2746 self.fail();
2747 None
2748 }
2749 }
2750 }
2751
2752 fn open_graph(&self, range: ByteRange) -> Option<GraphIndex> {
2757 if range.len <= NAMED_GRAPH_RESIDENT_MAX {
2758 let bytes = self.container_bytes(range)?;
2759 match decode_index_container(&bytes, self.codec, self.perms) {
2760 Ok(g) => Some(g),
2761 Err(_) => {
2762 self.fail();
2763 None
2764 }
2765 }
2766 } else {
2767 match open_index_container_lazy(
2768 &self.reader,
2769 range,
2770 self.codec,
2771 self.has_synopsis,
2772 self.read_concurrency,
2773 self.perms,
2774 ) {
2775 Ok((g, _, _)) => {
2776 let cap = self.tile_cap.load(std::sync::atomic::Ordering::Relaxed);
2780 if cap != u64::MAX {
2781 g.set_cache_cap(cap);
2782 }
2783 Some(g)
2784 }
2785 Err(_) => {
2786 self.fail();
2787 None
2788 }
2789 }
2790 }
2791 }
2792
2793 fn release(&mut self, iri: &str) {
2799 let Some((_, table)) = self.dir.get_mut() else {
2800 return;
2801 };
2802 for slab in table.iter_mut() {
2803 let Some(entries) = slab.get_mut() else {
2804 continue;
2805 };
2806 for e in entries.iter_mut() {
2807 let matches = e
2808 .meta
2809 .get()
2810 .map(|(name, _)| name.as_str() == iri)
2811 .unwrap_or(false);
2812 if matches {
2813 e.index.take();
2814 return;
2815 }
2816 }
2817 }
2818 }
2819
2820 fn for_each_opened(&self, mut f: impl FnMut(&GraphIndex)) {
2822 if let Some((_, table)) = self.dir.get() {
2823 for slab in table.iter().filter_map(|s| s.get()) {
2824 for e in slab.iter() {
2825 if let Some(g) = e.index.get() {
2826 f(g);
2827 }
2828 }
2829 }
2830 }
2831 }
2832}
2833
2834enum NamedGraphsSlot {
2839 Resident(Vec<(String, GraphIndex)>),
2840 Lazy(LazyNamedGraphs),
2841}
2842
2843impl NamedGraphsSlot {
2844 fn count(&self) -> usize {
2845 match self {
2846 NamedGraphsSlot::Resident(v) => v.len(),
2847 NamedGraphsSlot::Lazy(l) => l.count(),
2848 }
2849 }
2850
2851 fn name_at(&self, i: usize, exhaustive: bool) -> Option<&str> {
2854 match self {
2855 NamedGraphsSlot::Resident(v) => v.get(i).map(|(iri, _)| iri.as_str()),
2856 NamedGraphsSlot::Lazy(l) => l.name_at(i, exhaustive),
2857 }
2858 }
2859
2860 fn graph_at(&self, i: usize, exhaustive: bool) -> Option<(&str, &GraphIndex)> {
2861 match self {
2862 NamedGraphsSlot::Resident(v) => v.get(i).map(|(iri, g)| (iri.as_str(), g)),
2863 NamedGraphsSlot::Lazy(l) => l.graph_at(i, exhaustive),
2864 }
2865 }
2866
2867 fn find(&self, iri: &str) -> Option<&GraphIndex> {
2868 match self {
2869 NamedGraphsSlot::Resident(v) => v.iter().find(|(name, _)| name == iri).map(|(_, g)| g),
2870 NamedGraphsSlot::Lazy(l) => {
2871 for i in 0..l.count() {
2876 if l.name_at(i, false)? == iri {
2877 return l.graph_at(i, false).map(|(_, g)| g);
2878 }
2879 }
2880 None
2881 }
2882 }
2883 }
2884
2885 fn load_incomplete(&self) -> bool {
2886 match self {
2887 NamedGraphsSlot::Resident(v) => v.iter().any(|(_, g)| g.load_incomplete()),
2888 NamedGraphsSlot::Lazy(l) => {
2889 if l.failed.load(std::sync::atomic::Ordering::Relaxed) {
2890 return true;
2891 }
2892 let mut bad = false;
2893 l.for_each_opened(|g| bad |= g.load_incomplete());
2894 bad
2895 }
2896 }
2897 }
2898
2899 fn reset_load_failures(&self) {
2900 match self {
2901 NamedGraphsSlot::Resident(v) => {
2902 for (_, g) in v {
2903 g.reset_load_failure();
2904 }
2905 }
2906 NamedGraphsSlot::Lazy(l) => {
2907 l.failed.store(false, std::sync::atomic::Ordering::Relaxed);
2908 l.for_each_opened(|g| g.reset_load_failure());
2909 }
2910 }
2911 }
2912}
2913
2914pub struct Rete {
2916 header: Header,
2917 dict: Dictionary,
2918 index: GraphIndex,
2919 index_section_ranges: [ByteRange; NUM_PERMS],
2920 tile_ranges: [Vec<(u32, u32, ByteRange)>; NUM_PERMS],
2924 pyramid: PyramidSlot,
2925 text_index: TextIndexSlot,
2926 named_graphs: NamedGraphsSlot,
2927 metadata: Vec<u8>,
2932 service_client: Option<Box<dyn crate::service::ServiceClient>>,
2937 service_error: std::sync::Mutex<Option<String>>,
2941}
2942
2943impl Rete {
2944 pub fn open(bytes: &[u8]) -> Result<Self, FileError> {
2947 let header = Header::from_bytes(bytes)?;
2948
2949 let region = |off: u64, len: u64| -> Result<&[u8], FileError> {
2953 let start = off as usize;
2954 let end = start
2955 .checked_add(len as usize)
2956 .filter(|&e| e <= bytes.len())
2957 .ok_or(FileError::Container("section range out of bounds"))?;
2958 Ok(&bytes[start..end])
2959 };
2960
2961 let dict = decode_dictionary_container(
2962 region(header.dictionary_offset, header.dictionary_len)?,
2963 header.dict_codec,
2964 )?;
2965
2966 let index_bytes = region(header.root_dir_offset, header.root_dir_len)?;
2967 let index = decode_index_container(index_bytes, header.block_codec, header.perms)?;
2968 let index_section_ranges =
2969 decode_index_section_ranges(index_bytes, header.root_dir_offset, header.perms)?;
2970
2971 let pyramid = PyramidSlot::Resident(if header.pyramid_meta_len > 0 {
2972 Some(
2973 PyramidMeta::decode(region(header.pyramid_meta_offset, header.pyramid_meta_len)?)
2974 .map_err(|_| FileError::Container("malformed pyramid meta"))?,
2975 )
2976 } else {
2977 None
2978 });
2979
2980 let text_index = resident_text_index_slot(
2983 if header.text_index_len > 0 {
2984 Some(region(header.text_index_offset, header.text_index_len)?)
2985 } else {
2986 None
2987 },
2988 header.block_codec,
2989 )?;
2990
2991 let named_graphs = NamedGraphsSlot::Resident(if header.named_graphs_len > 0 {
2992 decode_named_graphs(
2993 region(header.named_graphs_offset, header.named_graphs_len)?,
2994 header.block_codec,
2995 header.perms,
2996 )?
2997 } else {
2998 Vec::new()
2999 });
3000
3001 let metadata = if header.metadata_len > 0 {
3002 region(header.metadata_offset, header.metadata_len)?.to_vec()
3003 } else {
3004 Vec::new()
3005 };
3006
3007 let tile_ranges =
3008 tile_file_ranges(index_bytes, header.root_dir_offset, &index_section_ranges);
3009 Ok(Self {
3010 header,
3011 dict,
3012 index,
3013 index_section_ranges,
3014 tile_ranges,
3015 pyramid,
3016 text_index,
3017 named_graphs,
3018 metadata,
3019 service_client: None,
3020 service_error: std::sync::Mutex::new(None),
3021 })
3022 }
3023
3024 pub fn set_service_client(&mut self, client: Box<dyn crate::service::ServiceClient>) {
3029 self.service_client = Some(client);
3030 }
3031
3032 pub(crate) fn service_client(&self) -> Option<&dyn crate::service::ServiceClient> {
3033 self.service_client.as_deref()
3034 }
3035
3036 pub(crate) fn record_service_error(&self, msg: &str) {
3038 let mut e = self.service_error.lock().unwrap();
3039 if e.is_none() {
3040 *e = Some(msg.to_string());
3041 }
3042 }
3043
3044 pub(crate) fn take_service_error(&self) -> Option<String> {
3047 self.service_error.lock().unwrap().take()
3048 }
3049
3050 pub fn header(&self) -> &Header {
3051 &self.header
3052 }
3053
3054 pub fn file_layout(&self) -> Vec<LayoutSegment> {
3060 let h = &self.header;
3061 let seg = |kind: &'static str, label: String, offset: u64, len: u64| LayoutSegment {
3062 kind,
3063 label,
3064 offset,
3065 len,
3066 };
3067 let mut out = vec![seg(
3073 "header",
3074 format!("header (fixed {} bytes)", crate::header::HEADER_LEN),
3075 0,
3076 crate::header::HEADER_LEN as u64,
3077 )];
3078 if h.metadata_len > 0 {
3079 out.push(seg(
3080 "metadata",
3081 "metadata (dataset card)".into(),
3082 h.metadata_offset,
3083 h.metadata_len,
3084 ));
3085 }
3086 out.push(seg(
3087 "dictionary",
3088 "dictionary (4 front-coded term sections)".into(),
3089 h.dictionary_offset,
3090 h.dictionary_len,
3091 ));
3092 for (si, perm) in crate::index::ALL_PERMS.into_iter().enumerate() {
3093 let sec = self.index_section_ranges[si];
3094 if sec.len == 0 {
3095 continue;
3096 }
3097 let first_tile = self.tile_ranges[si]
3098 .first()
3099 .map(|&(_, _, r)| r.offset)
3100 .unwrap_or(sec.offset + sec.len);
3101 if first_tile > sec.offset {
3102 out.push(seg(
3103 "directory",
3104 format!("{} tile directory", perm.name()),
3105 sec.offset,
3106 first_tile - sec.offset,
3107 ));
3108 }
3109 for (ti, &(min_a, max_a, r)) in self.tile_ranges[si].iter().enumerate() {
3110 out.push(seg(
3111 "tile",
3112 format!("{} tile {ti} (leading ids {min_a}..{max_a})", perm.name()),
3113 r.offset,
3114 r.len,
3115 ));
3116 }
3117 }
3118 if h.pyramid_meta_len > 0 {
3119 out.push(seg(
3120 "pyramid",
3121 "pyramid summary (communities + superedges)".into(),
3122 h.pyramid_meta_offset,
3123 h.pyramid_meta_len,
3124 ));
3125 }
3126 if h.named_graphs_len > 0 {
3127 out.push(seg(
3128 "named-graphs",
3129 format!("named graphs ({})", self.named_graphs.count()),
3130 h.named_graphs_offset,
3131 h.named_graphs_len,
3132 ));
3133 }
3134 out.sort_by_key(|s| s.offset);
3135 out
3136 }
3137
3138 pub fn metadata(&self) -> Option<&[u8]> {
3143 if self.metadata.is_empty() {
3144 None
3145 } else {
3146 Some(&self.metadata)
3147 }
3148 }
3149
3150 pub fn dictionary(&self) -> &Dictionary {
3151 &self.dict
3152 }
3153
3154 pub fn pyramid(&self) -> Option<&PyramidMeta> {
3156 match &self.pyramid {
3157 PyramidSlot::Resident(p) => p.as_ref(),
3158 PyramidSlot::Lazy { loader, cell } => cell.get_or_init(loader).as_ref(),
3160 }
3161 }
3162
3163 pub fn pyramid_if_loaded(&self) -> Option<&PyramidMeta> {
3168 match &self.pyramid {
3169 PyramidSlot::Resident(p) => p.as_ref(),
3170 PyramidSlot::Lazy { cell, .. } => cell.get().and_then(|o| o.as_ref()),
3171 }
3172 }
3173
3174 pub fn predicate_stats(&self) -> &[crate::meta::PredStat] {
3178 self.pyramid_if_loaded()
3179 .map(|p| p.predicate_stats.as_slice())
3180 .unwrap_or(&[])
3181 }
3182
3183 pub fn char_sets(&self) -> &[crate::meta::CharSet] {
3186 self.pyramid_if_loaded()
3187 .map(|p| p.char_sets.as_slice())
3188 .unwrap_or(&[])
3189 }
3190
3191 pub fn label_index(&self) -> &[crate::meta::LabelEntry] {
3194 self.pyramid_if_loaded()
3195 .map(|p| p.label_index.as_slice())
3196 .unwrap_or(&[])
3197 }
3198
3199 pub fn prefix_search(&self, prefix: &str, limit: usize) -> Vec<(String, String)> {
3205 let Some(pyr) = self.pyramid() else {
3206 return Vec::new();
3207 };
3208 pyr.prefix_search(prefix, limit)
3209 .into_iter()
3210 .filter_map(|e| {
3211 self.dict
3212 .subject_term(e.subject)
3213 .map(|iri| (e.label.clone(), iri))
3214 })
3215 .collect()
3216 }
3217
3218 pub(crate) fn text_index(&self) -> Option<&crate::text_index::TextIndex> {
3221 self.text_index.index()
3222 }
3223
3224 pub fn has_text_index(&self) -> bool {
3227 self.header.text_index_len > 0
3228 }
3229
3230 pub fn text_index_token_table_len(&self) -> Option<u64> {
3245 self.text_index.token_table_len()
3246 }
3247
3248 pub fn text_search(&self, words: &[&str], prefix: Option<&str>, limit: usize) -> Vec<String> {
3258 text_search_in(self.text_index(), &self.dict, words, prefix, limit)
3259 }
3260
3261 pub fn default_index(&self) -> &GraphIndex {
3263 &self.index
3264 }
3265
3266 pub fn dump(&self, graph: Option<&str>) -> Vec<TermTriple> {
3268 let index = match graph {
3269 None => &self.index,
3270 Some(g) => match self.graph_index(g) {
3271 Some(i) => i,
3272 None => return Vec::new(),
3273 },
3274 };
3275 index
3276 .match_pattern((None, None, None))
3277 .into_iter()
3278 .filter_map(|(s, p, o)| {
3279 Some((
3280 self.dict.subject_term(s)?,
3281 self.dict.predicate_term(p)?,
3282 self.dict.object_term(o)?,
3283 ))
3284 })
3285 .collect()
3286 }
3287
3288 pub fn dump_each<F: FnMut(&str, &str, &str)>(&self, graph: Option<&str>, f: F) {
3304 self.dump_filtered_each(graph, None, None, None, f)
3305 }
3306
3307 pub fn dump_filtered_each<F: FnMut(&str, &str, &str)>(
3378 &self,
3379 graph: Option<&str>,
3380 s: Option<&str>,
3381 p: Option<&str>,
3382 o: Option<&str>,
3383 mut f: F,
3384 ) {
3385 let index = match graph {
3386 None => Some(&self.index),
3387 Some(g) => self.graph_index(g),
3388 };
3389 let Some(index) = index else { return };
3390 let Some(pattern) = self.resolve_query_pattern(s, p, o) else {
3393 return;
3394 };
3395 let dict = &self.dict;
3396 const WINDOW: usize = 4096;
3441 let cap = dict.cache_cap();
3452 let prefetch_window =
3453 cap == u64::MAX || cap >= (WINDOW as u64) * 2 * (DICT_CHUNK_BUDGET as u64);
3454 let mut ids: Vec<(u32, u32, u32)> = Vec::with_capacity(WINDOW);
3455 let mut nodes: Vec<u32> = Vec::with_capacity(WINDOW * 2);
3456 let mut preds: Vec<u32> = Vec::with_capacity(WINDOW);
3457 let mut resolver = crate::dictionary::WindowResolver::default();
3458 let mut flush = |ids: &mut Vec<(u32, u32, u32)>,
3459 resolver: &mut crate::dictionary::WindowResolver,
3460 f: &mut F| {
3461 if ids.is_empty() {
3462 return;
3463 }
3464 if prefetch_window {
3468 nodes.clear();
3469 preds.clear();
3470 for &(s, p, o) in ids.iter() {
3471 nodes.push(dict.subject_node(s));
3472 nodes.push(dict.object_node(o));
3473 preds.push(p);
3474 }
3475 dict.prefetch_terms(&nodes, &preds);
3476 }
3477 dict.resolve_window(ids, resolver, |s, p, o| f(s, p, o));
3478 ids.clear();
3479 };
3480 for t in index.scan_iter(pattern) {
3481 ids.push(t);
3482 if ids.len() == WINDOW {
3483 flush(&mut ids, &mut resolver, &mut f);
3484 }
3485 }
3486 flush(&mut ids, &mut resolver, &mut f);
3487 }
3488
3489 pub fn dump_plan(
3497 &self,
3498 graph: Option<&str>,
3499 s: Option<&str>,
3500 p: Option<&str>,
3501 o: Option<&str>,
3502 ) -> DumpPlan {
3503 let index = match graph {
3504 None => Some(&self.index),
3505 Some(g) => self.graph_index(g),
3506 };
3507 let scan = index
3508 .zip(self.resolve_query_pattern(s, p, o))
3509 .map(|(ix, pattern)| ix.scan_plan(pattern));
3510 DumpPlan {
3511 scan,
3512 dictionary_bytes: self.header.dictionary_len,
3513 }
3514 }
3515
3516 pub fn dump_iter(&self, graph: Option<&str>) -> impl Iterator<Item = TermTriple> + '_ {
3534 self.query_iter(graph, None, None, None)
3535 }
3536
3537 pub fn query_iter(
3548 &self,
3549 graph: Option<&str>,
3550 s: Option<&str>,
3551 p: Option<&str>,
3552 o: Option<&str>,
3553 ) -> impl Iterator<Item = TermTriple> + '_ {
3554 let index = match graph {
3555 None => Some(&self.index),
3556 Some(g) => self.graph_index(g),
3557 };
3558 index
3559 .zip(self.resolve_query_pattern(s, p, o))
3560 .into_iter()
3561 .flat_map(|(ix, pattern)| ix.scan_iter(pattern))
3562 .filter_map(move |(s, p, o)| {
3563 Some((
3564 self.dict.subject_term(s)?,
3565 self.dict.predicate_term(p)?,
3566 self.dict.object_term(o)?,
3567 ))
3568 })
3569 }
3570
3571 pub fn query_batch(
3590 &self,
3591 graph: Option<&str>,
3592 s: Option<&str>,
3593 p: Option<&str>,
3594 o: Option<&str>,
3595 cursor: u64,
3596 max_quads: usize,
3597 ) -> (Vec<TermTriple>, u64, bool) {
3598 let index = match graph {
3599 None => &self.index,
3600 Some(g) => match self.graph_index(g) {
3603 Some(i) => i,
3604 None => return (Vec::new(), cursor, true),
3605 },
3606 };
3607 let Some(pattern) = self.resolve_query_pattern(s, p, o) else {
3610 return (Vec::new(), cursor, true);
3611 };
3612 let (ids, next, done) = index.scan_batch(pattern, cursor, max_quads);
3613 if !ids.is_empty() {
3617 let mut nodes = Vec::with_capacity(ids.len() * 2);
3618 let mut preds = Vec::with_capacity(ids.len());
3619 for &(s, p, o) in &ids {
3620 nodes.push(self.dict.subject_node(s));
3621 nodes.push(self.dict.object_node(o));
3622 preds.push(p);
3623 }
3624 self.dict.prefetch_terms(&nodes, &preds);
3625 }
3626 let triples = ids
3627 .into_iter()
3628 .filter_map(|(s, p, o)| {
3629 Some((
3630 self.dict.subject_term(s)?,
3631 self.dict.predicate_term(p)?,
3632 self.dict.object_term(o)?,
3633 ))
3634 })
3635 .collect();
3636 (triples, next, done)
3637 }
3638
3639 pub fn dump_batch(
3665 &self,
3666 graph: Option<&str>,
3667 cursor: u32,
3668 max_quads: usize,
3669 ) -> (Vec<TermTriple>, u32, bool) {
3670 fn next_tile_start(tiles: &[crate::index::Tile], sid: u32) -> Option<u32> {
3675 let i = tiles.partition_point(|t| t.leading_range().1 < sid);
3676 tiles.get(i).map(|t| t.leading_range().0.max(sid))
3677 }
3678
3679 let index = match graph {
3680 None => self.default_index(),
3681 Some(g) => match self.graph_index(g) {
3684 Some(i) => i,
3685 None => return (Vec::new(), cursor, true),
3686 },
3687 };
3688 let tiles = index.tile_sections()[IndexPermutation::Spo.section_index()];
3692 let (Some(first), Some(last)) = (
3693 tiles.first().map(|t| t.leading_range().0),
3694 tiles.last().map(|t| t.leading_range().1),
3695 ) else {
3696 return (Vec::new(), cursor, true); };
3698
3699 let mut sid = cursor.max(first);
3700 let mut probes = max_quads.saturating_mul(4).max(1 << 16);
3701 let mut ids: Vec<(u32, u32, u32)> = Vec::new();
3702 let mut exhausted = false;
3703 while sid <= last && probes > 0 {
3704 probes -= 1;
3705 let hits = index.match_pattern((Some(sid), None, None));
3706 if hits.is_empty() {
3707 match next_tile_start(tiles, sid) {
3711 None => {
3712 exhausted = true;
3713 break;
3714 }
3715 Some(next) if next > sid => {
3716 sid = next;
3717 continue;
3718 }
3719 Some(_) => {}
3720 }
3721 } else {
3722 ids.extend(hits);
3723 }
3724 if sid == u32::MAX {
3725 exhausted = true;
3726 break;
3727 }
3728 sid += 1;
3729 if ids.len() >= max_quads {
3730 break;
3731 }
3732 }
3733 let done = exhausted || sid > last;
3734
3735 let dict = self.dictionary();
3742 if !ids.is_empty() {
3743 let mut nodes = Vec::with_capacity(ids.len() * 2);
3744 let mut preds = Vec::with_capacity(ids.len());
3745 for &(s, p, o) in &ids {
3746 nodes.push(dict.subject_node(s));
3747 nodes.push(dict.object_node(o));
3748 preds.push(p);
3749 }
3750 dict.prefetch_terms(&nodes, &preds);
3751 }
3752 let triples = ids
3753 .into_iter()
3754 .filter_map(|(s, p, o)| {
3755 Some((
3756 dict.subject_term(s)?,
3757 dict.predicate_term(p)?,
3758 dict.object_term(o)?,
3759 ))
3760 })
3761 .collect();
3762 (triples, sid, done)
3763 }
3764
3765 pub fn named_graph_count(&self) -> usize {
3768 self.named_graphs.count()
3769 }
3770
3771 pub fn named_graph_name_at(&self, i: usize) -> Option<&str> {
3774 self.named_graphs.name_at(i, false)
3775 }
3776
3777 pub fn named_graph_at(&self, i: usize) -> Option<(&str, &GraphIndex)> {
3781 self.named_graphs.graph_at(i, false)
3782 }
3783
3784 pub(crate) fn named_graph_name_at_demand(&self, i: usize, exhaustive: bool) -> Option<&str> {
3790 self.named_graphs.name_at(i, exhaustive)
3791 }
3792
3793 pub(crate) fn named_graph_at_demand(
3796 &self,
3797 i: usize,
3798 exhaustive: bool,
3799 ) -> Option<(&str, &GraphIndex)> {
3800 self.named_graphs.graph_at(i, exhaustive)
3801 }
3802
3803 pub fn graph_names(&self) -> Vec<&str> {
3808 (0..self.named_graphs.count())
3809 .filter_map(|i| self.named_graphs.name_at(i, false))
3810 .collect()
3811 }
3812
3813 pub fn graph_index(&self, iri: &str) -> Option<&GraphIndex> {
3815 self.named_graphs.find(iri)
3816 }
3817
3818 pub fn match_ids(
3821 &self,
3822 pattern: (Option<u32>, Option<u32>, Option<u32>),
3823 ) -> Vec<(u32, u32, u32)> {
3824 self.index.match_pattern(pattern)
3825 }
3826
3827 pub fn predicate_pairs(&self, predicate: &str) -> Vec<(u32, u32)> {
3830 let pid = match self.dict.predicate_id(predicate) {
3831 Some(p) => p,
3832 None => return Vec::new(),
3833 };
3834 self.index
3835 .match_pattern((None, Some(pid), None))
3836 .into_iter()
3837 .map(|(s, _p, o)| (self.dict.subject_node(s), self.dict.object_node(o)))
3838 .collect()
3839 }
3840
3841 pub fn open_ranged<R: RangeReader>(reader: &R) -> Result<Self, FileError> {
3845 let head = reader.read_at(0, HEADER_LEN as u64)?;
3846 let header = Header::from_bytes(&head)?;
3847
3848 let dict_bytes = reader.read_at(header.dictionary_offset, header.dictionary_len)?;
3849 let dict = decode_dictionary_container(&dict_bytes, header.dict_codec)?;
3850
3851 let index_bytes = reader.read_at(header.root_dir_offset, header.root_dir_len)?;
3852 let index = decode_index_container(&index_bytes, header.block_codec, header.perms)?;
3853 let index_section_ranges =
3854 decode_index_section_ranges(&index_bytes, header.root_dir_offset, header.perms)?;
3855
3856 let pyramid = PyramidSlot::Resident(if header.pyramid_meta_len > 0 {
3857 let mb = reader.read_at(header.pyramid_meta_offset, header.pyramid_meta_len)?;
3858 Some(
3859 PyramidMeta::decode(&mb)
3860 .map_err(|_| FileError::Container("malformed pyramid meta"))?,
3861 )
3862 } else {
3863 None
3864 });
3865
3866 let text_index_bytes = if header.text_index_len > 0 {
3869 Some(reader.read_at(header.text_index_offset, header.text_index_len)?)
3870 } else {
3871 None
3872 };
3873 let text_index = resident_text_index_slot(text_index_bytes.as_deref(), header.block_codec)?;
3874
3875 let named_graphs = NamedGraphsSlot::Resident(if header.named_graphs_len > 0 {
3876 let nb = reader.read_at(header.named_graphs_offset, header.named_graphs_len)?;
3877 decode_named_graphs(&nb, header.block_codec, header.perms)?
3878 } else {
3879 Vec::new()
3880 });
3881
3882 let tile_ranges =
3886 tile_file_ranges(&index_bytes, header.root_dir_offset, &index_section_ranges);
3887 Ok(Self {
3888 header,
3889 dict,
3890 index,
3891 index_section_ranges,
3892 tile_ranges,
3893 pyramid,
3894 text_index,
3895 named_graphs,
3896 metadata: Vec::new(),
3897 service_client: None,
3898 service_error: std::sync::Mutex::new(None),
3899 })
3900 }
3901
3902 pub fn open_ranged_lazy<R: RangeReader + Send + Sync + 'static>(
3914 reader: R,
3915 ) -> Result<Self, FileError> {
3916 let head = reader.read_at(0, HEADER_LEN as u64)?;
3917 let header = Header::from_bytes(&head)?;
3918 let reader = std::sync::Arc::new(reader);
3919 let read_concurrency = reader.concurrency();
3922
3923 let dict = ranged_chunked_dictionary(&reader, &header, [true; 4])?;
3927
3928 let reader_dyn: std::sync::Arc<dyn RangeReader + Send + Sync> = reader.clone();
3932 let (index, index_section_ranges, tile_ranges) = open_index_container_lazy(
3933 &reader_dyn,
3934 ByteRange {
3935 offset: header.root_dir_offset,
3936 len: header.root_dir_len,
3937 },
3938 header.block_codec,
3939 header.has_tile_synopsis(),
3940 read_concurrency,
3941 header.perms,
3942 )?;
3943
3944 let pyramid = ranged_pyramid_slot(&reader, &header);
3948
3949 let text_index = ranged_text_index_slot(&reader, &header);
3954
3955 let named_graphs = if header.named_graphs_len > 0 {
3962 NamedGraphsSlot::Lazy(LazyNamedGraphs::new(
3963 reader_dyn,
3964 ByteRange {
3965 offset: header.named_graphs_offset,
3966 len: header.named_graphs_len,
3967 },
3968 header.block_codec,
3969 header.has_tile_synopsis(),
3970 read_concurrency,
3971 header.perms,
3972 ))
3973 } else {
3974 NamedGraphsSlot::Resident(Vec::new())
3975 };
3976
3977 Ok(Self {
3978 header,
3979 dict,
3980 index,
3981 index_section_ranges,
3982 tile_ranges,
3983 pyramid,
3984 text_index,
3985 named_graphs,
3986 metadata: Vec::new(),
3987 service_client: None,
3988 service_error: std::sync::Mutex::new(None),
3989 })
3990 }
3991
3992 pub fn index_incomplete(&self) -> bool {
3997 self.index.load_incomplete()
3998 || self.dict.load_incomplete()
3999 || self.named_graphs.load_incomplete()
4000 }
4001
4002 pub fn reset_load_failures(&self) {
4010 self.index.reset_load_failure();
4011 self.dict.reset_load_failure();
4012 self.named_graphs.reset_load_failures();
4013 }
4014
4015 pub fn set_memory_budget(&self, budget: Option<u64>) -> MemoryBudget {
4032 let total = budget.unwrap_or(u64::MAX);
4033 let split = split_memory_budget(total);
4034 self.dict.set_cache_cap(split.dict_cache);
4035 self.index.set_cache_cap(split.tile_cache);
4036 if let NamedGraphsSlot::Lazy(lazy) = &self.named_graphs {
4041 lazy.set_tile_cap(split.tile_cache);
4042 }
4043 if std::env::var("RETE_OPEN_DEBUG").is_ok() {
4044 eprintln!(
4045 "[memory-budget] total={} dict_cache={} tile_cache={} block_cache={}",
4046 if total == u64::MAX {
4047 "unlimited".to_string()
4048 } else {
4049 format!("{}MiB", total >> 20)
4050 },
4051 fmt_cap(split.dict_cache),
4052 fmt_cap(split.tile_cache),
4053 fmt_cap(split.block_cache),
4054 );
4055 }
4056 split
4057 }
4058
4059 pub fn dict_cache_stats(&self) -> crate::chunk_cache::CacheStats {
4063 self.dict.cache_stats()
4064 }
4065
4066 pub fn index_cache_stats(&self) -> crate::chunk_cache::CacheStats {
4068 self.index.cache_stats()
4069 }
4070
4071 pub fn dict_cache_cap(&self) -> u64 {
4073 self.dict.cache_cap()
4074 }
4075
4076 pub fn release_named_graph(&mut self, iri: &str) {
4089 if let NamedGraphsSlot::Lazy(lazy) = &mut self.named_graphs {
4090 lazy.release(iri);
4091 }
4092 }
4093
4094 fn resolve_query_pattern(
4095 &self,
4096 s: Option<&str>,
4097 p: Option<&str>,
4098 o: Option<&str>,
4099 ) -> Option<Pattern> {
4100 let sid = match s {
4101 Some(t) => match self.dict.subject_id(t) {
4102 Some(id) => Some(id),
4103 None => return None,
4104 },
4105 None => None,
4106 };
4107 let pid = match p {
4108 Some(t) => match self.dict.predicate_id(t) {
4109 Some(id) => Some(id),
4110 None => return None,
4111 },
4112 None => None,
4113 };
4114 let oid = match o {
4115 Some(t) => match self.dict.object_id(t) {
4116 Some(id) => Some(id),
4117 None => return None,
4118 },
4119 None => None,
4120 };
4121 Some((sid, pid, oid))
4122 }
4123
4124 pub fn query_with_provenance(
4128 &self,
4129 s: Option<&str>,
4130 p: Option<&str>,
4131 o: Option<&str>,
4132 ) -> Vec<TripleProvenance> {
4133 let pattern = match self.resolve_query_pattern(s, p, o) {
4134 Some(pattern) => pattern,
4135 None => return Vec::new(),
4136 };
4137
4138 let index_permutation = GraphIndex::best_permutation_in(self.header.perms, pattern);
4139 let dictionary_range = ByteRange {
4140 offset: self.header.dictionary_offset,
4141 len: self.header.dictionary_len,
4142 };
4143 let index_range = ByteRange {
4144 offset: self.header.root_dir_offset,
4145 len: self.header.root_dir_len,
4146 };
4147 let index_section_range = self.index_section_ranges[index_permutation.section_index()];
4148 let pyramid_range = (self.header.pyramid_meta_len > 0).then_some(ByteRange {
4149 offset: self.header.pyramid_meta_offset,
4150 len: self.header.pyramid_meta_len,
4151 });
4152
4153 let tiles = &self.tile_ranges[index_permutation.section_index()];
4154 self.index
4155 .match_pattern(pattern)
4156 .into_iter()
4157 .filter_map(|(s, p, o)| {
4158 let terms = (
4159 self.dict.subject_term(s)?,
4160 self.dict.predicate_term(p)?,
4161 self.dict.object_term(o)?,
4162 );
4163 let a = index_permutation.forward((s, p, o)).0;
4166 let ti = tiles.partition_point(|&(_, max_a, _)| max_a < a);
4167 let (tile, tile_range) = match tiles.get(ti) {
4168 Some(&(min_a, _, range)) if min_a <= a => (
4169 Some(format!("{}/{ti}", index_permutation.name())),
4170 Some(range),
4171 ),
4172 _ => (None, None),
4173 };
4174 Some(TripleProvenance {
4175 terms,
4176 ids: (s, p, o),
4177 graph: None,
4178 matched_pattern: pattern,
4179 index_permutation,
4180 dictionary_range,
4181 index_range,
4182 index_section_range,
4183 pyramid_range,
4184 tile,
4185 tile_range,
4186 })
4187 })
4188 .collect()
4189 }
4190
4191 pub fn query(&self, s: Option<&str>, p: Option<&str>, o: Option<&str>) -> Vec<TermTriple> {
4195 self.query_with_provenance(s, p, o)
4196 .into_iter()
4197 .map(|m| m.terms)
4198 .collect()
4199 }
4200
4201 pub fn query_in_graph(
4209 &self,
4210 graph: Option<&str>,
4211 s: Option<&str>,
4212 p: Option<&str>,
4213 o: Option<&str>,
4214 ) -> Vec<TermTriple> {
4215 let pattern = match self.resolve_query_pattern(s, p, o) {
4216 Some(pattern) => pattern,
4217 None => return Vec::new(),
4218 };
4219 let index = match graph {
4220 None => &self.index,
4221 Some(g) => match self.graph_index(g) {
4222 Some(i) => i,
4223 None => return Vec::new(),
4224 },
4225 };
4226 let ids = index.match_pattern(pattern);
4227 if !ids.is_empty() {
4235 let mut nodes = Vec::with_capacity(ids.len() * 2);
4236 let mut preds = Vec::with_capacity(ids.len());
4237 for &(s, p, o) in &ids {
4238 nodes.push(self.dict.subject_node(s));
4239 nodes.push(self.dict.object_node(o));
4240 preds.push(p);
4241 }
4242 self.dict.prefetch_terms(&nodes, &preds);
4243 }
4244 ids.into_iter()
4245 .filter_map(|(s, p, o)| {
4246 Some((
4247 self.dict.subject_term(s)?,
4248 self.dict.predicate_term(p)?,
4249 self.dict.object_term(o)?,
4250 ))
4251 })
4252 .collect()
4253 }
4254
4255 pub fn query_quads(
4260 &self,
4261 s: Option<&str>,
4262 p: Option<&str>,
4263 o: Option<&str>,
4264 ) -> Vec<(TermTriple, Option<String>)> {
4265 let mut out: Vec<(TermTriple, Option<String>)> = self
4266 .query_in_graph(None, s, p, o)
4267 .into_iter()
4268 .map(|t| (t, None))
4269 .collect();
4270 for i in 0..self.named_graphs.count() {
4271 let Some(iri) = self.named_graphs.name_at(i, true) else {
4274 continue;
4275 };
4276 for triple in self.query_in_graph(Some(iri), s, p, o) {
4277 out.push((triple, Some(iri.to_string())));
4278 }
4279 }
4280 out
4281 }
4282
4283 pub fn query_ranged<R: RangeReader>(
4290 reader: &R,
4291 s: Option<&str>,
4292 p: Option<&str>,
4293 o: Option<&str>,
4294 ) -> Result<Vec<TermTriple>, FileError> {
4295 let routed = match route_pattern(reader, s, p, o)? {
4296 Some(routed) => routed,
4297 None => return Ok(Vec::new()),
4298 };
4299 let matches = fetch_routed_matches(reader, &routed)?;
4300 Ok(matches
4301 .into_iter()
4302 .filter_map(|(s, p, o)| {
4303 Some((
4304 routed.dict.subject_term(s)?,
4305 routed.dict.predicate_term(p)?,
4306 routed.dict.object_term(o)?,
4307 ))
4308 })
4309 .collect())
4310 }
4311
4312 pub fn route_pattern_ranged<R: RangeReader>(
4316 reader: &R,
4317 s: Option<&str>,
4318 p: Option<&str>,
4319 o: Option<&str>,
4320 ) -> Result<bool, FileError> {
4321 Ok(route_pattern(reader, s, p, o)?.is_some())
4322 }
4323}
4324
4325struct RoutedPattern {
4328 dict: Dictionary,
4329 pattern: Pattern,
4330 permutation: IndexPermutation,
4331 header: Header,
4332 section: ByteRange,
4334}
4335
4336fn route_pattern<R: RangeReader>(
4339 reader: &R,
4340 s: Option<&str>,
4341 p: Option<&str>,
4342 o: Option<&str>,
4343) -> Result<Option<RoutedPattern>, FileError> {
4344 let head = reader.read_at(0, HEADER_LEN as u64)?;
4345 let header = Header::from_bytes(&head)?;
4346
4347 let dict_bytes = reader.read_at(header.dictionary_offset, header.dictionary_len)?;
4348 let dict = decode_dictionary_container(&dict_bytes, header.dict_codec)?;
4349
4350 let Some(pattern) = resolve_query_pattern(&dict, s, p, o) else {
4351 return Ok(None);
4352 };
4353 let permutation = GraphIndex::best_permutation_in(header.perms, pattern);
4354 let section = locate_container_section_ranged(
4355 reader,
4356 header.root_dir_offset,
4357 header.root_dir_len,
4358 header
4359 .perms
4360 .position(permutation)
4361 .ok_or(FileError::Container("routed to an absent permutation"))?,
4362 header.perms.len() as u64,
4363 )?;
4364 Ok(Some(RoutedPattern {
4365 dict,
4366 pattern,
4367 permutation,
4368 header,
4369 section,
4370 }))
4371}
4372
4373fn fetch_routed_matches<R: RangeReader>(
4378 reader: &R,
4379 routed: &RoutedPattern,
4380) -> Result<Vec<Triple>, FileError> {
4381 let dir = read_tile_directory_ranged(reader, routed.section)?;
4382 let [pa, _, _] = routed.permutation.order_pattern(routed.pattern);
4383 let codec = routed.header.block_codec;
4384 let mut out = Vec::new();
4385 match pa {
4386 Some(a) => {
4389 for e in dir.iter().filter(|e| e.min_a <= a && a <= e.max_a) {
4390 let bytes = reader.read_at(routed.section.offset + e.start, e.end - e.start)?;
4391 let tile = decompress(codec, &bytes)?;
4392 out.extend(GraphIndex::match_serialized_block(
4393 &tile,
4394 routed.permutation,
4395 routed.pattern,
4396 ));
4397 }
4398 }
4399 None => {
4402 if let (Some(first), Some(last)) = (dir.first(), dir.last()) {
4403 let base = first.start;
4404 let body = reader.read_at(routed.section.offset + base, last.end - base)?;
4405 for e in &dir {
4406 let tile = decompress(
4407 codec,
4408 &body[(e.start - base) as usize..(e.end - base) as usize],
4409 )?;
4410 out.extend(GraphIndex::match_serialized_block(
4411 &tile,
4412 routed.permutation,
4413 routed.pattern,
4414 ));
4415 }
4416 }
4417 }
4418 }
4419 out.sort_unstable();
4420 Ok(out)
4421}
4422
4423fn resolve_query_pattern(
4424 dict: &Dictionary,
4425 s: Option<&str>,
4426 p: Option<&str>,
4427 o: Option<&str>,
4428) -> Option<Pattern> {
4429 let sid = match s {
4430 Some(t) => Some(dict.subject_id(t)?),
4431 None => None,
4432 };
4433 let pid = match p {
4434 Some(t) => Some(dict.predicate_id(t)?),
4435 None => None,
4436 };
4437 let oid = match o {
4438 Some(t) => Some(dict.object_id(t)?),
4439 None => None,
4440 };
4441 Some((sid, pid, oid))
4442}
4443
4444fn ranged_chunked_dictionary<R: RangeReader + Send + Sync + 'static>(
4459 reader: &std::sync::Arc<R>,
4460 header: &Header,
4461 want: [bool; 4],
4462) -> Result<Dictionary, FileError> {
4463 let cache = crate::chunk_cache::ChunkCache::unlimited_arc();
4466 let mut dict_sections: Vec<crate::dict::ChunkedSection> = Vec::with_capacity(4);
4467 for si in 0..4 {
4468 if !want[si] {
4469 dict_sections.push(crate::dict::ChunkedSection::from_parts(
4470 crate::dict::SectionMeta {
4471 term_count: 0,
4472 restart_interval: 1,
4473 restart_offsets: Vec::new(),
4474 },
4475 Vec::new(),
4476 None,
4477 cache.clone(),
4478 si as u8,
4479 ));
4480 continue;
4481 }
4482 let section = locate_container_section_ranged(
4483 reader.as_ref(),
4484 header.dictionary_offset,
4485 header.dictionary_len,
4486 si,
4487 4,
4488 )?;
4489 let (meta, entries) = read_dict_dir_ranged(reader.as_ref(), section)?;
4490 let ranges: Vec<ByteRange> = entries
4491 .iter()
4492 .map(|e| ByteRange {
4493 offset: section.offset + e.start,
4494 len: (e.end - e.start),
4495 })
4496 .collect();
4497 let chunks: Vec<crate::dict::SectionChunk> = entries
4498 .into_iter()
4499 .map(|e| crate::dict::SectionChunk::new(e.first_run, e.key, e.body_start))
4500 .collect();
4501 let chunk_reader = reader.clone();
4502 let codec = header.dict_codec;
4503 let loader_ranges = ranges.clone();
4504 let loader: crate::dict::ChunkLoader = Box::new(move |ci| {
4505 let range = loader_ranges.get(ci)?;
4506 let bytes = chunk_reader.read_at(range.offset, range.len).ok()?;
4507 decompress(codec, &bytes).ok()
4508 });
4509 let bulk_reader = reader.clone();
4512 let bulk: crate::dict::ChunkBulkLoader = Box::new(move |cis| {
4513 let want: Option<Vec<ByteRange>> =
4514 cis.iter().map(|&ci| ranges.get(ci).copied()).collect();
4515 let blobs = read_coalesced(bulk_reader.as_ref(), &want?, DICT_COALESCE_GAP)?;
4516 blobs.iter().map(|b| decompress(codec, b).ok()).collect()
4517 });
4518 dict_sections.push(
4519 crate::dict::ChunkedSection::from_parts(
4520 meta,
4521 chunks,
4522 Some(loader),
4523 cache.clone(),
4524 si as u8,
4525 )
4526 .with_bulk_loader(bulk),
4527 );
4528 }
4529 let dict_arr: [crate::dict::ChunkedSection; 4] = dict_sections
4530 .try_into()
4531 .map_err(|_| FileError::Container("expected 4 dictionary sections"))?;
4532 Ok(Dictionary::from_chunked_sections(dict_arr))
4533}
4534
4535fn ranged_pyramid_slot<R: RangeReader + Send + Sync + 'static>(
4539 reader: &std::sync::Arc<R>,
4540 header: &Header,
4541) -> PyramidSlot {
4542 if header.pyramid_meta_len == 0 {
4543 return PyramidSlot::Resident(None);
4544 }
4545 let pyr_reader = reader.clone();
4546 let pyr_off = header.pyramid_meta_offset;
4547 let pyr_len = header.pyramid_meta_len;
4548 PyramidSlot::Lazy {
4549 loader: Box::new(move || {
4550 let mb = pyr_reader.read_at(pyr_off, pyr_len).ok()?;
4551 PyramidMeta::decode(&mb).ok()
4552 }),
4553 cell: std::sync::OnceLock::new(),
4554 }
4555}
4556
4557fn ranged_text_index_slot<R: RangeReader + Send + Sync + 'static>(
4562 reader: &std::sync::Arc<R>,
4563 header: &Header,
4564) -> TextIndexSlot {
4565 if header.text_index_len == 0 {
4566 return TextIndexSlot::Resident {
4567 index: None,
4568 token_table_len: None,
4569 };
4570 }
4571 let ti_reader = reader.clone();
4572 let probe_reader = reader.clone();
4573 let ti_off = header.text_index_offset;
4574 let ti_len = header.text_index_len;
4575 let codec = header.block_codec;
4576 TextIndexSlot::Lazy {
4577 loader: Box::new(move || {
4578 let prefix_len = read_token_table_len(&*ti_reader, ti_off, ti_len)?;
4582 let prefix = ti_reader.read_at(ti_off, prefix_len).ok()?;
4583 let postings_base = crate::text_index::TextIndex::postings_base(&prefix)? as u64;
4584 let postings_abs = ti_off + postings_base;
4585 let pr = ti_reader.clone();
4586 let posting_loader =
4587 Box::new(move |off: u64, len: u64| pr.read_at(postings_abs + off, len).ok());
4588 crate::text_index::TextIndex::from_token_table(&prefix, codec, posting_loader).ok()
4589 }),
4590 cell: std::sync::OnceLock::new(),
4591 token_table: Box::new(move || read_token_table_len(&*probe_reader, ti_off, ti_len)),
4594 token_table_cell: std::sync::OnceLock::new(),
4595 }
4596}
4597
4598fn text_search_in(
4601 ti: Option<&crate::text_index::TextIndex>,
4602 dict: &Dictionary,
4603 words: &[&str],
4604 prefix: Option<&str>,
4605 limit: usize,
4606) -> Vec<String> {
4607 let Some(ti) = ti else {
4608 return Vec::new();
4609 };
4610 let mut acc: Option<Vec<u32>> = None;
4614 if let Some(p) = prefix {
4615 acc = Some(ti.prefix(&p.to_lowercase()));
4616 }
4617 for w in words {
4618 for tok in crate::text_index::tokenize(w) {
4619 let posting = ti.lookup(&tok);
4620 acc = Some(match acc {
4621 Some(a) => intersect_sorted(&a, &posting),
4622 None => posting,
4623 });
4624 if acc.as_ref().is_some_and(|a| a.is_empty()) {
4625 return Vec::new();
4626 }
4627 }
4628 }
4629 let ids = acc.unwrap_or_default();
4630 let mut out = Vec::with_capacity(if limit > 0 {
4631 limit.min(ids.len())
4632 } else {
4633 ids.len()
4634 });
4635 for id in ids {
4636 if let Some(iri) = dict.subject_term(id) {
4637 out.push(iri);
4638 if limit > 0 && out.len() >= limit {
4639 break;
4640 }
4641 }
4642 }
4643 out
4644}
4645
4646pub struct SearchView {
4664 header: Header,
4665 dict: Dictionary,
4666 pyramid: PyramidSlot,
4667 text_index: TextIndexSlot,
4668}
4669
4670impl SearchView {
4671 pub fn open_ranged<R: RangeReader + Send + Sync + 'static>(
4675 reader: R,
4676 ) -> Result<Self, FileError> {
4677 let head = reader.read_at(0, HEADER_LEN as u64)?;
4678 let header = Header::from_bytes(&head)?;
4679 let reader = std::sync::Arc::new(reader);
4680 let dict = ranged_chunked_dictionary(&reader, &header, [true, true, false, false])?;
4684 let pyramid = ranged_pyramid_slot(&reader, &header);
4685 let text_index = ranged_text_index_slot(&reader, &header);
4686 Ok(Self {
4687 header,
4688 dict,
4689 pyramid,
4690 text_index,
4691 })
4692 }
4693
4694 pub fn header(&self) -> &Header {
4696 &self.header
4697 }
4698
4699 pub fn has_text_index(&self) -> bool {
4701 self.header.text_index_len > 0
4702 }
4703
4704 pub fn has_pyramid(&self) -> bool {
4707 self.header.pyramid_meta_len > 0
4708 }
4709
4710 pub fn text_search(&self, words: &[&str], prefix: Option<&str>, limit: usize) -> Vec<String> {
4713 text_search_in(self.text_index.index(), &self.dict, words, prefix, limit)
4714 }
4715
4716 pub fn text_index_token_table_len(&self) -> Option<u64> {
4720 self.text_index.token_table_len()
4721 }
4722
4723 pub fn prefix_search(&self, prefix: &str, limit: usize) -> Vec<(String, String)> {
4726 let pyr = match &self.pyramid {
4727 PyramidSlot::Resident(p) => p.as_ref(),
4728 PyramidSlot::Lazy { loader, cell } => cell.get_or_init(loader).as_ref(),
4729 };
4730 let Some(pyr) = pyr else {
4731 return Vec::new();
4732 };
4733 pyr.prefix_search(prefix, limit)
4734 .into_iter()
4735 .filter_map(|e| {
4736 self.dict
4737 .subject_term(e.subject)
4738 .map(|iri| (e.label.clone(), iri))
4739 })
4740 .collect()
4741 }
4742}
4743
4744#[must_use]
4749pub struct SummaryView {
4750 pub round: u32,
4751 pub summary: Vec<SuperEdge>,
4752 pub class_hierarchy: Vec<ClassNode>,
4754 pub level_rollups: Vec<LevelRollup>,
4756 pub level_links: Vec<LevelLinks>,
4758 pub descriptors: Vec<CommunityDescriptor>,
4760 pub subclass_cycles: Vec<Vec<String>>,
4762 pub disjoint_pairs: Vec<(String, String)>,
4764 pub equivalent_pairs: Vec<(String, String)>,
4766 dict: Dictionary,
4767}
4768
4769impl SummaryView {
4770 pub fn open_ranged<R: RangeReader>(reader: &R) -> Result<Option<Self>, FileError> {
4772 let head = reader.read_at(0, HEADER_LEN as u64)?;
4773 let header = Header::from_bytes(&head)?;
4774 if header.pyramid_meta_len == 0 {
4775 return Ok(None);
4776 }
4777
4778 let dict_bytes = reader.read_at(header.dictionary_offset, header.dictionary_len)?;
4779 let dict = decode_dictionary_container(&dict_bytes, header.dict_codec)?;
4780
4781 let mb = reader.read_at(header.pyramid_meta_offset, header.pyramid_meta_len)?;
4782 let meta =
4783 PyramidMeta::decode(&mb).map_err(|_| FileError::Container("malformed pyramid meta"))?;
4784
4785 Ok(Some(SummaryView {
4786 round: meta.round,
4787 summary: meta.summary,
4788 class_hierarchy: meta.class_hierarchy,
4789 level_rollups: meta.level_rollups,
4790 level_links: meta.level_links,
4791 descriptors: meta.descriptors,
4792 subclass_cycles: meta.subclass_cycles,
4793 disjoint_pairs: meta.disjoint_pairs,
4794 equivalent_pairs: meta.equivalent_pairs,
4795 dict,
4796 }))
4797 }
4798
4799 pub fn level_count(&self) -> usize {
4801 self.level_rollups.len()
4802 }
4803
4804 pub fn level_rollup(&self, k: usize) -> Option<&LevelRollup> {
4807 self.level_rollups.get(k)
4808 }
4809
4810 pub fn predicate_term(&self, id: u32) -> Option<String> {
4812 self.dict.predicate_term(id)
4813 }
4814
4815 pub fn predicate_total(&self, predicate: &str) -> u32 {
4818 match self.dict.predicate_id(predicate) {
4819 Some(pid) => self
4820 .summary
4821 .iter()
4822 .filter(|e| e.predicate == pid)
4823 .map(|e| e.count)
4824 .sum(),
4825 None => 0,
4826 }
4827 }
4828
4829 pub fn predicate_totals(&self) -> Vec<(String, u32)> {
4831 let mut by_pred: std::collections::BTreeMap<u32, u32> = std::collections::BTreeMap::new();
4832 for e in &self.summary {
4833 *by_pred.entry(e.predicate).or_default() += e.count;
4834 }
4835 let mut out: Vec<(String, u32)> = by_pred
4836 .into_iter()
4837 .filter_map(|(pid, c)| self.dict.predicate_term(pid).map(|t| (t, c)))
4838 .collect();
4839 out.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
4840 out
4841 }
4842
4843 pub fn community_count(&self) -> usize {
4845 let mut comms = std::collections::BTreeSet::new();
4846 for e in &self.summary {
4847 comms.insert(e.s_comm);
4848 comms.insert(e.o_comm);
4849 }
4850 comms.len()
4851 }
4852
4853 pub fn tbox_coherence(&self) -> Vec<crate::reason::Inconsistency> {
4868 schema_coherence(
4869 &self.class_hierarchy,
4870 &self.subclass_cycles,
4871 &self.disjoint_pairs,
4872 &self.equivalent_pairs,
4873 )
4874 }
4875
4876 pub fn tbox_is_coherent(&self) -> bool {
4879 self.tbox_coherence().is_empty()
4880 }
4881}
4882
4883pub fn schema_coherence(
4889 class_hierarchy: &[ClassNode],
4890 subclass_cycles: &[Vec<String>],
4891 disjoint_pairs: &[(String, String)],
4892 equivalent_pairs: &[(String, String)],
4893) -> Vec<crate::reason::Inconsistency> {
4894 use crate::reason::Inconsistency;
4895 use std::collections::{BTreeMap, BTreeSet, VecDeque};
4896 const MAX_REACH: usize = 100_000;
4897
4898 let mut out: Vec<Inconsistency> = Vec::new();
4899
4900 for cyc in subclass_cycles {
4901 let detail = if cyc.len() == 1 {
4902 format!("{} is rdfs:subClassOf itself (a cycle)", cyc[0])
4903 } else {
4904 format!(
4905 "classes {{{}}} are mutually rdfs:subClassOf (a cycle)",
4906 cyc.join(", ")
4907 )
4908 };
4909 out.push(Inconsistency {
4910 kind: "subclass-cycle",
4911 detail,
4912 });
4913 }
4914
4915 if !disjoint_pairs.is_empty() {
4916 let mut adj: BTreeMap<&str, Vec<&str>> = BTreeMap::new();
4918 for n in class_hierarchy {
4919 let e = adj.entry(n.class.as_str()).or_default();
4920 for p in &n.parents {
4921 e.push(p.as_str());
4922 }
4923 }
4924 for (a, b) in equivalent_pairs {
4925 adj.entry(a.as_str()).or_default().push(b.as_str());
4926 adj.entry(b.as_str()).or_default().push(a.as_str());
4927 }
4928
4929 let mut focuses: BTreeSet<&str> =
4931 class_hierarchy.iter().map(|n| n.class.as_str()).collect();
4932 for (a, b) in disjoint_pairs.iter().chain(equivalent_pairs) {
4933 focuses.insert(a.as_str());
4934 focuses.insert(b.as_str());
4935 }
4936
4937 let mut seen: BTreeSet<&str> = BTreeSet::new();
4938 for &c in &focuses {
4939 let mut reach: BTreeSet<&str> = BTreeSet::new();
4941 let mut q: VecDeque<&str> = VecDeque::new();
4942 reach.insert(c);
4943 q.push_back(c);
4944 while let Some(x) = q.pop_front() {
4945 if reach.len() > MAX_REACH {
4946 break;
4947 }
4948 if let Some(ns) = adj.get(x) {
4949 for &p in ns {
4950 if reach.insert(p) {
4951 q.push_back(p);
4952 }
4953 }
4954 }
4955 }
4956 for (x, y) in disjoint_pairs {
4957 if reach.contains(x.as_str()) && reach.contains(y.as_str()) && seen.insert(c) {
4958 out.push(Inconsistency {
4959 kind: "unsatisfiable-class",
4960 detail: format!(
4961 "{c} is a subclass of both {x} and {y}, which are \
4962 owl:disjointWith — no individual can be a {c}"
4963 ),
4964 });
4965 break;
4966 }
4967 }
4968 }
4969 }
4970
4971 out.sort_by(|a, b| (a.kind, &a.detail).cmp(&(b.kind, &b.detail)));
4972 out
4973}
4974
4975pub fn read_schema_coherence_ranged<R: RangeReader>(
4983 reader: &R,
4984) -> Result<Option<Vec<crate::reason::Inconsistency>>, FileError> {
4985 let head = reader.read_at(0, HEADER_LEN as u64)?;
4986 let header = Header::from_bytes(&head)?;
4987 if header.pyramid_meta_len == 0 {
4988 return Ok(None);
4989 }
4990 if header.schema_meta_len > 0 && (header.schema_meta_len as u64) <= header.pyramid_meta_len {
4994 let off =
4995 header.pyramid_meta_offset + header.pyramid_meta_len - header.schema_meta_len as u64;
4996 let block = reader.read_at(off, header.schema_meta_len as u64)?;
4997 let (hierarchy, cycles, disjoint, equivalent) = crate::meta::decode_schema_block(&block)
4998 .map_err(|_| FileError::Container("malformed schema block"))?;
4999 return Ok(Some(schema_coherence(
5000 &hierarchy,
5001 &cycles,
5002 &disjoint,
5003 &equivalent,
5004 )));
5005 }
5006 let mb = reader.read_at(header.pyramid_meta_offset, header.pyramid_meta_len)?;
5008 let meta =
5009 PyramidMeta::decode(&mb).map_err(|_| FileError::Container("malformed pyramid meta"))?;
5010 Ok(Some(schema_coherence(
5011 &meta.class_hierarchy,
5012 &meta.subclass_cycles,
5013 &meta.disjoint_pairs,
5014 &meta.equivalent_pairs,
5015 )))
5016}
5017
5018#[allow(clippy::type_complexity)]
5026pub fn read_schema_summary_ranged<R: RangeReader>(
5027 reader: &R,
5028) -> Result<Option<(Vec<(String, u64)>, Vec<(String, String, String, u64)>)>, FileError> {
5029 let head = reader.read_at(0, HEADER_LEN as u64)?;
5030 let header = Header::from_bytes(&head)?;
5031 if header.pyramid_meta_len == 0 {
5032 return Ok(None);
5033 }
5034 if header.schema_meta_len > 0 && (header.schema_meta_len as u64) <= header.pyramid_meta_len {
5035 let off =
5036 header.pyramid_meta_offset + header.pyramid_meta_len - header.schema_meta_len as u64;
5037 let block = reader.read_at(off, header.schema_meta_len as u64)?;
5038 let summary = crate::meta::decode_schema_block_summary(&block)
5039 .map_err(|_| FileError::Container("malformed schema block"))?;
5040 return Ok(Some(summary));
5041 }
5042 let mb = reader.read_at(header.pyramid_meta_offset, header.pyramid_meta_len)?;
5044 let meta =
5045 PyramidMeta::decode(&mb).map_err(|_| FileError::Container("malformed pyramid meta"))?;
5046 if meta.level_rollups.is_empty() && meta.level_links.is_empty() {
5047 return Ok(None);
5048 }
5049 let classes = meta
5050 .level_rollups
5051 .iter()
5052 .max_by_key(|r| r.depth)
5053 .map(|r| r.classes.clone())
5054 .unwrap_or_default();
5055 let relations = meta
5056 .level_links
5057 .iter()
5058 .max_by_key(|l| l.depth)
5059 .map(|l| {
5060 l.links
5061 .iter()
5062 .map(|c| {
5063 (
5064 c.s_class.clone(),
5065 c.predicate.clone(),
5066 c.o_class.clone(),
5067 c.count,
5068 )
5069 })
5070 .collect()
5071 })
5072 .unwrap_or_default();
5073 Ok(Some((classes, relations)))
5074}
5075
5076#[cfg(test)]
5081pub(crate) fn dict_chunk_keys_for_test(image: &[u8]) -> Vec<Vec<Vec<u8>>> {
5082 let header = Header::from_bytes(&image[..HEADER_LEN]).unwrap();
5083 let s = header.dictionary_offset as usize;
5084 let e = s + header.dictionary_len as usize;
5085 decode_container(&image[s..e], CODEC_NONE)
5086 .unwrap()
5087 .iter()
5088 .map(|payload| {
5089 let (_meta, entries) = parse_chunked_dict_dir(payload, payload.len() as u64).unwrap();
5090 entries.into_iter().map(|c| c.key).collect()
5091 })
5092 .collect()
5093}
5094
5095#[cfg(test)]
5096mod tests {
5097 use super::*;
5098 use crate::dictionary::DictionaryBuilder;
5099 use crate::index::GraphIndexBuilder;
5100
5101 #[test]
5107 fn three_permutation_file_round_trips_every_open_path() {
5108 use crate::index::PermSet;
5109 use crate::reader::SliceReader;
5110
5111 let build = |perms: PermSet| {
5112 let mut db = crate::DictionaryBuilder::new();
5113 let mut triples = Vec::new();
5114 for i in 0..60u32 {
5115 let (s, p, o) = (
5116 format!("<http://ex/s{}>", i % 11),
5117 format!("<http://ex/p{}>", i % 3),
5118 format!("<http://ex/o{}>", i % 7),
5119 );
5120 db.observe(&s, &p, &o);
5121 triples.push((s, p, o));
5122 }
5123 let dict = db.build();
5124 let ids: Vec<(u32, u32, u32)> = triples
5125 .iter()
5126 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
5127 .collect();
5128 let def = GraphIndexBuilder::from_triples(ids.clone())
5129 .with_tile_budget(96)
5130 .with_perms(perms)
5131 .build();
5132 let named = GraphIndexBuilder::from_triples(ids)
5133 .with_tile_budget(96)
5134 .with_perms(perms)
5135 .build();
5136 write_dataset(
5137 &dict,
5138 &def,
5139 &[("<http://ex/g>".to_string(), named)],
5140 true,
5141 &[],
5142 0,
5143 )
5144 };
5145
5146 let six = build(PermSet::ALL);
5147 let three = build(PermSet::CORE);
5148 assert!(three.len() < six.len());
5149 assert_eq!(Header::from_bytes(&three).unwrap().perms, PermSet::CORE);
5150 assert_eq!(Header::from_bytes(&six).unwrap().perms, PermSet::ALL);
5151
5152 let shapes: Vec<(Option<&str>, Option<&str>, Option<&str>)> = vec![
5153 (None, None, None),
5154 (Some("<http://ex/s3>"), None, None),
5155 (None, Some("<http://ex/p1>"), None),
5156 (None, None, Some("<http://ex/o5>")),
5157 (Some("<http://ex/s3>"), Some("<http://ex/p1>"), None),
5158 (Some("<http://ex/s3>"), None, Some("<http://ex/o5>")),
5159 (None, Some("<http://ex/p1>"), Some("<http://ex/o5>")),
5160 (
5161 Some("<http://ex/s3>"),
5162 Some("<http://ex/p1>"),
5163 Some("<http://ex/o5>"),
5164 ),
5165 ];
5166
5167 for (image, tag) in [(&six, "six"), (&three, "three")] {
5168 let resident = Rete::open(image).unwrap();
5169 let leaked: &'static [u8] = Box::leak(image.clone().into_boxed_slice());
5170 let ranged = Rete::open_ranged(&SliceReader::new(leaked)).unwrap();
5171 let lazy = Rete::open_ranged_lazy(SliceReader::new(leaked)).unwrap();
5172 for (s, p, o) in &shapes {
5173 let want: Vec<_> = {
5174 let mut v = Rete::open(&six).unwrap().query(*s, *p, *o);
5175 v.sort();
5176 v
5177 };
5178 for (got, path) in [
5179 (resident.query(*s, *p, *o), "resident"),
5180 (ranged.query(*s, *p, *o), "ranged"),
5181 (lazy.query(*s, *p, *o), "lazy"),
5182 ] {
5183 let mut got = got;
5184 got.sort();
5185 assert_eq!(got, want, "{tag}/{path} disagreed on {:?}", (s, p, o));
5186 }
5187 let mut g = resident.query_in_graph(Some("<http://ex/g>"), *s, *p, *o);
5189 g.sort();
5190 assert_eq!(g, want, "{tag}/named disagreed on {:?}", (s, p, o));
5191 }
5192 }
5193
5194 let leaked: &'static [u8] = Box::leak(three.clone().into_boxed_slice());
5198 let reader = SliceReader::new(leaked);
5199 let routed = Rete::query_ranged(&reader, None, None, Some("<http://ex/o5>")).unwrap();
5200 let mut want = Rete::open(&six)
5201 .unwrap()
5202 .query(None, None, Some("<http://ex/o5>"));
5203 want.sort();
5204 let mut got = routed;
5205 got.sort();
5206 assert_eq!(got, want, "routed single-pattern read on a lean file");
5207 }
5208
5209 #[test]
5210 fn read_coalesced_merges_within_gap_and_splits_beyond() {
5211 use crate::reader::{CountingReader, SliceReader};
5212 let bytes = vec![0u8; 4096];
5213 let ranges = [
5215 ByteRange { offset: 0, len: 16 },
5216 ByteRange {
5217 offset: 48,
5218 len: 16,
5219 },
5220 ByteRange {
5221 offset: 1088,
5222 len: 16,
5223 },
5224 ];
5225 let r = CountingReader::new(SliceReader::new(&bytes));
5227 let out = read_coalesced(&r, &ranges, 16).unwrap();
5228 assert_eq!(out.len(), 3);
5229 assert_eq!(r.requests(), 3);
5230 let r = CountingReader::new(SliceReader::new(&bytes));
5232 read_coalesced(&r, &ranges, 64).unwrap();
5233 assert_eq!(r.requests(), 2);
5234 let r = CountingReader::new(SliceReader::new(&bytes));
5236 read_coalesced(&r, &ranges, 4096).unwrap();
5237 assert_eq!(r.requests(), 1);
5238 }
5239
5240 fn build_image() -> Vec<u8> {
5241 let triples = [
5242 ("Alice", "knows", "Bob"),
5243 ("Bob", "knows", "Carol"),
5244 ("Alice", "age", "30"),
5245 ];
5246 let mut db = DictionaryBuilder::new();
5247 for (s, p, o) in triples {
5248 db.observe(s, p, o);
5249 }
5250 let dict = db.build();
5251
5252 let mut ib = GraphIndexBuilder::new();
5253 for (s, p, o) in triples {
5254 ib.push(dict.encode(s, p, o).unwrap());
5255 }
5256 let index = ib.build();
5257
5258 let (meta, levels) = build_pyramid_meta(&dict, &triples_ids(&dict), DEFAULT_TILE_BUDGET);
5259 write_file(&dict, &index, false, &meta, levels)
5260 }
5261
5262 fn triples_ids(dict: &Dictionary) -> Vec<(u32, u32, u32)> {
5263 [
5264 ("Alice", "knows", "Bob"),
5265 ("Bob", "knows", "Carol"),
5266 ("Alice", "age", "30"),
5267 ]
5268 .iter()
5269 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
5270 .collect()
5271 }
5272
5273 #[test]
5274 fn file_round_trips_header_and_counts() {
5275 let bytes = build_image();
5276 let rete = Rete::open(&bytes).unwrap();
5277 assert_eq!(rete.header().quad_count, 3);
5278 assert!(rete.header().term_count >= 5);
5279 let expected_codec = writer_codec();
5280 assert_eq!(rete.header().dict_codec, expected_codec);
5281 assert_eq!(rete.header().block_codec, expected_codec);
5282 assert_eq!(&bytes[bytes.len() - 4..], &MAGIC); }
5284
5285 #[test]
5290 fn multi_tile_file_round_trips_and_routes() {
5291 let triples: Vec<(String, String, String)> = (0..200)
5292 .map(|i| {
5293 (
5294 format!("<http://ex/s/{i}>"),
5295 format!("<http://ex/p/{}>", i % 5),
5296 format!("<http://ex/o/{}>", i % 23),
5297 )
5298 })
5299 .collect();
5300 let mut db = DictionaryBuilder::new();
5301 for (s, p, o) in &triples {
5302 db.observe(s, p, o);
5303 }
5304 let dict = db.build();
5305 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
5306 for (s, p, o) in &triples {
5307 ib.push(dict.encode(s, p, o).unwrap());
5308 }
5309 let index = ib.build();
5310 assert!(
5311 index.tile_sections()[0].len() > 3,
5312 "tiny budget must force many tiles"
5313 );
5314 let bytes = write_file(&dict, &index, false, &[], 0);
5315
5316 let rete = Rete::open(&bytes).unwrap();
5317 assert_eq!(rete.header().version, crate::header::CURRENT_FORMAT_VERSION);
5318 assert_eq!(rete.query(None, None, None).len(), 200);
5319 assert_eq!(rete.query(Some("<http://ex/s/7>"), None, None).len(), 1);
5320 assert_eq!(
5321 rete.query(None, Some("<http://ex/p/3>"), None).len(),
5322 40,
5323 "predicate extent spans tiles"
5324 );
5325 assert_eq!(
5326 rete.query(None, None, Some("<http://ex/o/22>")).len(),
5327 8 );
5329
5330 use crate::reader::SliceReader;
5332 let reader = SliceReader::new(&bytes);
5333 let routed = Rete::query_ranged(&reader, Some("<http://ex/s/7>"), None, None).unwrap();
5334 assert_eq!(routed.len(), 1);
5335 let routed = Rete::query_ranged(&reader, None, Some("<http://ex/p/3>"), None).unwrap();
5336 assert_eq!(routed.len(), 40);
5337 let routed = Rete::query_ranged(&reader, None, None, Some("<http://ex/o/22>")).unwrap();
5338 assert_eq!(routed.len(), 8);
5339 }
5340
5341 #[test]
5350 fn tile_directory_offsets_survive_past_4gib() {
5351 let mut dir = Vec::new();
5352 write_uvarint(&mut dir, 2); write_uvarint(&mut dir, 5); write_uvarint(&mut dir, 0); write_uvarint(&mut dir, 3 << 30); write_uvarint(&mut dir, 1); write_uvarint(&mut dir, 0);
5358 write_uvarint(&mut dir, 2 << 30); let total = dir.len() as u64 + (3u64 << 30) + (2u64 << 30) + 64;
5360 let entries = parse_tile_directory(&dir, total).unwrap();
5361 assert_eq!(entries.len(), 2);
5362 assert_eq!(entries[1].start, dir.len() as u64 + (3u64 << 30));
5363 assert!(
5364 entries[1].end > u32::MAX as u64,
5365 "tail tile sits past 4 GiB"
5366 );
5367 assert!(parse_tile_directory(&dir, 1 << 20).is_err());
5369 }
5370
5371 #[test]
5372 fn tile_synopsis_trailer_round_trips() {
5373 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
5374 for i in 0..200u32 {
5375 ib.push((i, i % 7, i % 13));
5376 }
5377 let index = ib.build();
5378 let tile_count = index.tile_sections()[0].len();
5379 assert!(tile_count > 3, "tiny budget forces many tiles");
5380
5381 let payload = encode_tiled_section(&index, 0, CODEC_NONE);
5382 let dir = parse_tile_directory(&payload, payload.len() as u64).unwrap();
5383 assert_eq!(dir.len(), tile_count);
5384 let trailer_start = dir.iter().map(|e| e.end).max().unwrap();
5386 assert!(
5387 trailer_start < payload.len() as u64,
5388 "a trailer follows the tiles"
5389 );
5390 for e in &dir {
5391 assert!(
5392 e.end <= payload.len() as u64,
5393 "tiles still located within the payload"
5394 );
5395 }
5396 let syn = parse_tile_synopsis(&payload, trailer_start as usize, dir.len()).unwrap();
5397 for (e, (min_b, max_b, min_c, max_c)) in dir.iter().zip(syn) {
5398 let block = decompress(CODEC_NONE, &payload[e.start as usize..e.end as usize]).unwrap();
5399 let z = *crate::triples::TripleBlock::parse(&block).unwrap().zone();
5400 assert_eq!(
5401 (min_b, max_b, min_c, max_c),
5402 (z.min_b, z.max_b, z.min_c, z.max_c),
5403 "synopsis equals the tile's own zone"
5404 );
5405 }
5406 }
5407
5408 #[test]
5412 fn tile_synopsis_lazy_matches_reference_every_shape() {
5413 use crate::reader::{CountingReader, SliceReader};
5414 let triples: Vec<(String, String, String)> = (0..200u32)
5415 .map(|i| {
5416 (
5417 format!("<http://ex/s/{i:04}>"),
5418 format!("<http://ex/p/{}>", i % 7),
5419 format!("<http://ex/o/{:04}>", i % 13),
5420 )
5421 })
5422 .collect();
5423 let mut db = DictionaryBuilder::new();
5424 for (s, p, o) in &triples {
5425 db.observe(s, p, o);
5426 }
5427 let dict = db.build();
5428 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
5429 for (s, p, o) in &triples {
5430 ib.push(dict.encode(s, p, o).unwrap());
5431 }
5432 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
5433
5434 let eager = Rete::open(&bytes).unwrap();
5435 assert!(
5436 eager.header().has_tile_synopsis(),
5437 "new files set the synopsis flag"
5438 );
5439
5440 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
5443 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
5444 let lazy = Rete::open_ranged_lazy(reader).unwrap();
5445
5446 let brute = |s: Option<&str>, p: Option<&str>, o: Option<&str>| {
5447 let mut v: Vec<(String, String, String)> = triples
5448 .iter()
5449 .filter(|(a, b, c)| {
5450 s.is_none_or(|x| x == a) && p.is_none_or(|x| x == b) && o.is_none_or(|x| x == c)
5451 })
5452 .cloned()
5453 .collect();
5454 v.sort();
5455 v
5456 };
5457 let sv = [
5459 None,
5460 Some("<http://ex/s/0007>"),
5461 Some("<http://ex/s/0130>"),
5462 Some("<http://ex/s/9999>"),
5463 ];
5464 let pv = [
5465 None,
5466 Some("<http://ex/p/3>"),
5467 Some("<http://ex/p/6>"),
5468 Some("<http://ex/p/999>"),
5469 ];
5470 let ov = [
5471 None,
5472 Some("<http://ex/o/0000>"),
5473 Some("<http://ex/o/0012>"),
5474 Some("<http://ex/o/9999>"),
5475 ];
5476 for &s in &sv {
5477 for &p in &pv {
5478 for &o in &ov {
5479 let mut e = eager.query(s, p, o);
5480 e.sort();
5481 let mut l = lazy.query(s, p, o);
5482 l.sort();
5483 let r = brute(s, p, o);
5484 assert_eq!(e, r, "eager {s:?} {p:?} {o:?}");
5485 assert_eq!(l, r, "lazy {s:?} {p:?} {o:?} — synopsis over-pruned");
5486 }
5487 }
5488 }
5489 assert!(!lazy.index_incomplete(), "no lazy fetch failed");
5490 }
5491
5492 #[test]
5496 fn polyglot_offset_reads_lazily() {
5497 use crate::reader::{
5498 detect_polyglot_base, CountingReader, OffsetReader, SliceReader, POLYGLOT_DIGITS,
5499 POLYGLOT_MARKER,
5500 };
5501
5502 let image = build_image();
5506 let mut shell = Vec::new();
5507 shell.extend_from_slice(b"<!DOCTYPE html><html><head><!--");
5508 shell.extend_from_slice(POLYGLOT_MARKER);
5509 let digits_at = shell.len();
5510 shell.extend_from_slice(&[b'0'; POLYGLOT_DIGITS]); shell.extend_from_slice(b"--></head><body>a web page</body></html>\n");
5512 shell.resize(50_000, b' '); let base = shell.len() as u64;
5514 let digits = format!("{base:0width$}", width = POLYGLOT_DIGITS);
5515 shell[digits_at..digits_at + POLYGLOT_DIGITS].copy_from_slice(digits.as_bytes());
5516
5517 let mut poly = shell;
5518 poly.extend_from_slice(&image);
5519
5520 assert_ne!(
5522 &poly[0..4],
5523 b"RETE",
5524 "polyglot must not start with the magic"
5525 );
5526 let detected = detect_polyglot_base(&poly[..crate::header::HEADER_LEN]).unwrap();
5528 assert_eq!(detected, base);
5529
5530 let leaked: &'static [u8] = Box::leak(poly.into_boxed_slice());
5532 let counting = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
5533 let lazy = Rete::open_ranged_lazy(OffsetReader::new(counting.clone(), detected)).unwrap();
5534
5535 let plain = Rete::open(&image).unwrap();
5537 let mut want = plain.query(None, None, None);
5538 want.sort();
5539 let mut got = lazy.query(None, None, None);
5540 got.sort();
5541 assert_eq!(
5542 got, want,
5543 "polyglot lazy query differs from the plain .rete"
5544 );
5545 assert!(!got.is_empty(), "expected some triples");
5546 assert!(!lazy.index_incomplete(), "a lazy fetch failed");
5547
5548 let read = counting.bytes_read();
5551 eprintln!(
5552 "polyglot lazy read: HTML shell {base} B + .rete image {} B; \
5553 the query touched only {read} B (never the shell).",
5554 image.len()
5555 );
5556 assert!(
5557 read < base,
5558 "read {read} bytes; the HTML shell alone is {base} — the reader \
5559 touched the prefix instead of only the embedded .rete"
5560 );
5561 }
5562
5563 #[cfg(test)]
5566 fn build_text_indexed(triples: &[(String, String, String)]) -> Vec<u8> {
5567 let mut db = DictionaryBuilder::new();
5568 for (s, p, o) in triples {
5569 db.observe(s, p, o);
5570 }
5571 let dict = db.build();
5572 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
5573 let mut id_triples: Vec<(u32, u32, u32)> = Vec::with_capacity(triples.len());
5574 for (s, p, o) in triples {
5575 let t = dict.encode(s, p, o).unwrap();
5576 ib.push(t);
5577 id_triples.push(t);
5578 }
5579 let index = ib.build();
5580 let text_index = compute_text_index(&dict, &id_triples);
5581 assert!(
5582 !text_index.is_empty(),
5583 "literals should produce a text index"
5584 );
5585 write_dataset_with_metadata(&dict, &index, &[], false, &[], 0, &[], &text_index)
5586 }
5587
5588 #[test]
5592 fn text_index_eager_matches_brute_force() {
5593 let triples: Vec<(String, String, String)> = vec![
5594 (
5595 "<http://ex/s0>",
5596 "<http://ex/label>",
5597 "\"alpha glucose phosphate\"",
5598 ),
5599 ("<http://ex/s1>", "<http://ex/label>", "\"beta Glucose\""),
5600 ("<http://ex/s2>", "<http://ex/label>", "\"gamma fructose\""),
5601 (
5602 "<http://ex/s3>",
5603 "<http://ex/note>",
5604 "\"einstein relativity\"",
5605 ),
5606 (
5607 "<http://ex/s4>",
5608 "<http://ex/ref>",
5609 "<http://ex/not-a-literal>",
5610 ),
5611 ]
5612 .into_iter()
5613 .map(|(s, p, o)| (s.to_string(), p.to_string(), o.to_string()))
5614 .collect();
5615 let bytes = build_text_indexed(&triples);
5616 let rete = Rete::open(&bytes).unwrap();
5617 assert!(rete.has_text_index());
5618
5619 let brute = |words: &[&str]| -> Vec<String> {
5621 let mut v: Vec<String> = triples
5622 .iter()
5623 .filter(|(_, _, o)| {
5624 crate::terms::is_literal(o)
5625 && words.iter().all(|w| {
5626 let wl = w.to_lowercase();
5627 crate::terms::literal_lexical(o)
5628 .unwrap()
5629 .split(|c: char| !c.is_alphanumeric())
5630 .any(|t| t.to_lowercase() == wl)
5631 })
5632 })
5633 .map(|(s, _, _)| s.clone())
5634 .collect();
5635 v.sort();
5636 v.dedup();
5637 v
5638 };
5639
5640 let mut got = rete.text_search(&["glucose"], None, 0);
5641 got.sort();
5642 assert_eq!(got, brute(&["glucose"]), "case-insensitive single word");
5643
5644 let mut got = rete.text_search(&["glucose", "phosphate"], None, 0);
5646 got.sort();
5647 assert_eq!(got, brute(&["glucose", "phosphate"]));
5648
5649 assert!(rete.text_search(&["zzznope"], None, 0).is_empty());
5651
5652 let got = rete.text_search(&[], Some("ein"), 0);
5654 assert_eq!(got, vec!["<http://ex/s3>".to_string()]);
5655
5656 let mut db = DictionaryBuilder::new();
5658 for (s, p, o) in &triples {
5659 db.observe(s, p, o);
5660 }
5661 let dict = db.build();
5662 let mut ib = GraphIndexBuilder::new();
5663 for (s, p, o) in &triples {
5664 ib.push(dict.encode(s, p, o).unwrap());
5665 }
5666 let plain = write_dataset(&dict, &ib.build(), &[], false, &[], 0);
5667 let plain_rete = Rete::open(&plain).unwrap();
5668 assert!(!plain_rete.has_text_index());
5669 assert!(plain_rete.text_search(&["glucose"], None, 0).is_empty());
5670 }
5671
5672 #[test]
5676 fn text_index_lazy_faults_only_queried_postings() {
5677 use crate::reader::{CountingReader, SliceReader};
5678 let mut triples: Vec<(String, String, String)> = (0..300u32)
5681 .map(|i| {
5682 (
5683 format!("<http://ex/s/{i:04}>"),
5684 "<http://ex/label>".to_string(),
5685 format!("\"common word number {i}\""),
5686 )
5687 })
5688 .collect();
5689 for i in [3u32, 77, 250] {
5690 triples.push((
5691 format!("<http://ex/s/{i:04}>"),
5692 "<http://ex/tag>".to_string(),
5693 "\"raretoken\"".to_string(),
5694 ));
5695 }
5696 let bytes = build_text_indexed(&triples);
5697 let eager = Rete::open(&bytes).unwrap();
5698 let mut want = eager.text_search(&["raretoken"], None, 0);
5699 want.sort();
5700 assert_eq!(want.len(), 3);
5701
5702 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
5703 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
5704 let lazy = Rete::open_ranged_lazy(reader.clone()).unwrap();
5705 let before = reader.bytes_read();
5708 let mut got = lazy.text_search(&["raretoken"], None, 0);
5709 got.sort();
5710 assert_eq!(got, want, "lazy search matches eager");
5711 let pulled = reader.bytes_read() - before;
5712 let ti_len = eager.header().text_index_len;
5715 assert!(
5716 pulled < ti_len,
5717 "search pulled {pulled} B but the section is {ti_len} B — faulted too much"
5718 );
5719 assert!(!lazy.index_incomplete());
5720 }
5721
5722 #[test]
5727 fn search_view_matches_full_open_for_far_fewer_bytes() {
5728 use crate::reader::{CountingReader, SliceReader};
5729 let filler = "lorem ipsum dolor sit amet consectetur adipiscing elit ".repeat(20);
5732 let mut triples: Vec<(String, String, String)> = (0..300u32)
5733 .map(|i| {
5734 (
5735 format!("<http://ex/s/{i:04}>"),
5736 "<http://ex/abstract>".to_string(),
5737 format!("\"{filler} number {i}\""),
5738 )
5739 })
5740 .collect();
5741 for i in [3u32, 77, 250] {
5742 triples.push((
5743 format!("<http://ex/s/{i:04}>"),
5744 "<http://ex/tag>".to_string(),
5745 "\"raretoken\"".to_string(),
5746 ));
5747 }
5748 let bytes = build_text_indexed(&triples);
5749 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
5750
5751 let full_reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
5752 let full = Rete::open_ranged_lazy(full_reader.clone()).unwrap();
5753 let mut want = full.text_search(&["raretoken"], None, 0);
5754 want.sort();
5755 assert_eq!(want.len(), 3);
5756
5757 let search_reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
5758 let view = SearchView::open_ranged(search_reader.clone()).unwrap();
5759 assert!(view.has_text_index());
5760 let mut got = view.text_search(&["raretoken"], None, 0);
5761 got.sort();
5762 assert_eq!(got, want, "search view matches the full open");
5763 assert!(
5764 search_reader.bytes_read() < full_reader.bytes_read(),
5765 "search read {} B, full open read {} B — the narrow open bought nothing",
5766 search_reader.bytes_read(),
5767 full_reader.bytes_read()
5768 );
5769
5770 let mut prefix_hits = view.prefix_search("", 5);
5773 prefix_hits.sort();
5774 assert!(!prefix_hits.is_empty() || !view.has_pyramid());
5775 }
5776
5777 #[test]
5787 fn dict_chunk_directory_probe_costs_about_one_directory() {
5788 use crate::reader::{CountingReader, SliceReader};
5789 let head = "x".repeat(3000);
5803 let triples: Vec<(String, String, String)> = (0..2000u32)
5804 .map(|i| {
5805 let tail = char::from(b'a' + (i % 26) as u8).to_string().repeat(1000);
5806 (
5807 format!("<http://ex/s/{i:04}>"),
5808 "<http://ex/abstract>".to_string(),
5809 format!("\"{head}{i:04}{tail}\""),
5810 )
5811 })
5812 .collect();
5813 let mut db = DictionaryBuilder::new();
5814 for (s, p, o) in &triples {
5815 db.observe(s, p, o);
5816 }
5817 let dict = db.build();
5818 let mut ib = GraphIndexBuilder::new();
5819 for (s, p, o) in &triples {
5820 ib.push(dict.encode(s, p, o).unwrap());
5821 }
5822 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
5823
5824 let reader = CountingReader::new(SliceReader::new(&bytes));
5825 let header = Header::from_bytes(&reader.read_at(0, HEADER_LEN as u64).unwrap()).unwrap();
5826 let section = locate_container_section_ranged(
5828 &reader,
5829 header.dictionary_offset,
5830 header.dictionary_len,
5831 2,
5832 4,
5833 )
5834 .unwrap();
5835
5836 let before = reader.bytes_read();
5837 let (_meta, entries) = read_dict_dir_ranged(&reader, section).unwrap();
5838 let spent = reader.bytes_read() - before;
5839
5840 let dir_end = entries[0].start;
5842 let head = reader.read_at(section.offset, 64).unwrap();
5843 let (header_len, n0) = read_uvarint(&head).unwrap();
5844 let dir_start = n0 as u64 + header_len;
5845 let dir_len = dir_end - dir_start;
5846 let dir_total = section.len - dir_start;
5847 assert!(
5848 dir_start + dir_len > 8192,
5849 "directory fits the header prefix ({dir_len} B) — the probe never runs"
5850 );
5851
5852 let mut old = 0u64;
5854 let mut p = 4096u64.min(dir_total).max(1);
5855 loop {
5856 let got = p.min(dir_total);
5857 old += got;
5858 if got >= dir_len || p >= dir_total {
5859 break;
5860 }
5861 p = p.saturating_mul(2).min(dir_total);
5862 }
5863 assert!(
5864 spent <= dir_len * 2,
5865 "probe read {spent} B for a {dir_len} B directory — past the append-only bound"
5866 );
5867 assert!(
5868 spent < old,
5869 "probe read {spent} B; the re-reading loop it replaced spent {old} B"
5870 );
5871 }
5872
5873 fn verbatim_keyed_twin(payload: &[u8], codec: u8) -> Vec<u8> {
5883 let (_meta, entries) = parse_chunked_dict_dir(payload, payload.len() as u64).unwrap();
5884 let (header_len, n0) = read_uvarint(payload).unwrap();
5885 let mut out = Vec::new();
5886 write_uvarint(&mut out, header_len);
5887 out.extend_from_slice(&payload[n0..n0 + header_len as usize]);
5888 write_uvarint(&mut out, entries.len() as u64);
5889 let mut prev_run = 0usize;
5890 for e in &entries {
5891 let body = decompress(codec, &payload[e.start as usize..e.end as usize]).unwrap();
5892 let first = crate::dict::run_first_term(&body, 0).unwrap();
5893 write_uvarint(&mut out, (e.first_run - prev_run) as u64);
5894 write_uvarint(&mut out, first.len() as u64);
5895 out.extend_from_slice(&first);
5896 write_uvarint(&mut out, e.end - e.start);
5897 prev_run = e.first_run;
5898 }
5899 for e in &entries {
5900 out.extend_from_slice(&payload[e.start as usize..e.end as usize]);
5901 }
5902 out
5903 }
5904
5905 fn chunk_bounds_of(payload: &[u8], codec: u8) -> Vec<(Vec<u8>, Vec<u8>)> {
5907 let (_meta, entries) = parse_chunked_dict_dir(payload, payload.len() as u64).unwrap();
5908 entries
5909 .iter()
5910 .map(|e| {
5911 let body = decompress(codec, &payload[e.start as usize..e.end as usize]).unwrap();
5912 (
5913 crate::dict::run_first_term(&body, 0).unwrap(),
5914 crate::dict::run_last_term(&body, 0, body.len()).unwrap(),
5915 )
5916 })
5917 .collect()
5918 }
5919
5920 fn adversarial_terms() -> Vec<String> {
5926 let mut terms: Vec<String> = Vec::new();
5927 for i in 0..6_000 {
5928 terms.push(format!("<http://example.org/work/{i:06}>"));
5929 }
5930 let shared = "X".repeat(900);
5931 for i in 0..4_000 {
5932 terms.push(format!("\"{shared}-{i:05}-tail\""));
5933 }
5934 let deep = "Y".repeat(3_000); for i in 0..120 {
5936 terms.push(format!("\"{deep}{i:03}\""));
5937 }
5938 for i in 0..6 {
5939 terms.push(format!("\"BIG-{i:02}-{}\"", "Z".repeat(90_000)));
5941 }
5942 for i in 0..800 {
5943 terms.push(format!(
5944 "\"h\u{e9}llo-{}-\u{1F600}{i:04}\"",
5945 "\u{e9}".repeat(40)
5946 ));
5947 }
5948 for i in 0..2_000 {
5949 terms.push(format!("\"{i:05}-{}\"", "T".repeat(1_800)));
5950 }
5951 terms.sort();
5952 terms.dedup();
5953 terms
5954 }
5955
5956 fn chunked_section_of(terms: &[String]) -> Vec<u8> {
5957 let mut b = crate::dict::DictSectionBuilder::new().with_restart_interval(16);
5958 for t in terms {
5959 b.push(t.clone());
5960 }
5961 b.build()
5962 }
5963
5964 #[test]
5974 fn chunk_directory_keys_are_separators_not_terms() {
5975 let terms = adversarial_terms();
5976 let raw = chunked_section_of(&terms);
5977 let payload = encode_chunked_dict_section(&raw, CODEC_NONE);
5978 let (_meta, entries) = parse_chunked_dict_dir(&payload, payload.len() as u64).unwrap();
5979 let bounds = chunk_bounds_of(&payload, CODEC_NONE);
5980 assert!(
5981 entries.len() > 20,
5982 "need a many-chunk section, got {}",
5983 entries.len()
5984 );
5985
5986 assert!(entries[0].key.is_empty(), "chunk 0 needs no separator");
5987 let mut key_bytes = 0usize;
5988 let mut term_bytes = 0usize;
5989 let mut differ = 0usize;
5990 for (i, e) in entries.iter().enumerate() {
5991 key_bytes += e.key.len();
5992 term_bytes += bounds[i].0.len();
5993 if e.key != bounds[i].0 {
5994 differ += 1;
5995 }
5996 assert!(
5997 e.key <= bounds[i].0,
5998 "chunk {i}: key must be <= its own first term"
5999 );
6000 if i > 0 {
6001 assert!(
6002 e.key.as_slice() > bounds[i - 1].1.as_slice(),
6003 "chunk {i}: key must be > the previous chunk's last term"
6004 );
6005 assert!(
6006 e.key > entries[i - 1].key,
6007 "chunk {i}: keys must ascend for partition_point"
6008 );
6009 assert_eq!(
6011 e.key,
6012 crate::dict::shortest_separator(&bounds[i - 1].1, &bounds[i].0),
6013 "chunk {i}: key is not the SHORTEST separator"
6014 );
6015 }
6016 }
6017 assert!(
6018 differ > 0,
6019 "every key equalled its first term — the writer stores terms again"
6020 );
6021 assert!(
6022 key_bytes * 4 < term_bytes,
6023 "separators bought almost nothing: {key_bytes} B vs {term_bytes} B of first terms"
6024 );
6025 eprintln!(
6026 "chunks={} first-term keys={term_bytes} B separator keys={key_bytes} B ({:.1}x)",
6027 entries.len(),
6028 term_bytes as f64 / key_bytes.max(1) as f64
6029 );
6030 }
6031
6032 #[test]
6043 fn separator_keys_route_identically_to_verbatim_first_terms() {
6044 let terms = adversarial_terms();
6045 let raw = chunked_section_of(&terms);
6046 let sep_payload = encode_chunked_dict_section(&raw, CODEC_NONE);
6047 let ver_payload = verbatim_keyed_twin(&sep_payload, CODEC_NONE);
6048
6049 let (_ms, es) = parse_chunked_dict_dir(&sep_payload, sep_payload.len() as u64).unwrap();
6051 let (_mv, ev) = parse_chunked_dict_dir(&ver_payload, ver_payload.len() as u64).unwrap();
6052 assert_eq!(es.len(), ev.len());
6053 for (a, b) in es.iter().zip(ev.iter()) {
6054 assert_eq!(a.first_run, b.first_run);
6055 assert_eq!(
6056 &sep_payload[a.start as usize..a.end as usize],
6057 &ver_payload[b.start as usize..b.end as usize],
6058 "chunk bodies must be identical"
6059 );
6060 }
6061 assert!(
6062 sep_payload.len() < ver_payload.len(),
6063 "separator payload is not smaller"
6064 );
6065
6066 let sec_s = decode_chunked_dict_section(
6067 &sep_payload,
6068 CODEC_NONE,
6069 crate::chunk_cache::ChunkCache::unlimited_arc(),
6070 0,
6071 )
6072 .unwrap();
6073 let sec_v = decode_chunked_dict_section(
6074 &ver_payload,
6075 CODEC_NONE,
6076 crate::chunk_cache::ChunkCache::unlimited_arc(),
6077 0,
6078 )
6079 .unwrap();
6080
6081 for (i, t) in terms.iter().enumerate() {
6083 let want = Some(i as u32 + 1);
6084 assert_eq!(sec_v.id(t), want, "verbatim id() lost a term");
6085 assert_eq!(sec_s.id(t), want, "separator id() lost a term");
6086 }
6087 for i in 0..terms.len() as u32 {
6089 assert_eq!(sec_s.term(i + 1), sec_v.term(i + 1));
6090 assert_eq!(
6091 sec_s.term(i + 1).as_deref(),
6092 Some(terms[i as usize].as_str())
6093 );
6094 }
6095 let bounds = chunk_bounds_of(&sep_payload, CODEC_NONE);
6097 let mut probes: Vec<String> = Vec::new();
6098 let push = |bytes: &[u8], probes: &mut Vec<String>| {
6099 if let Ok(s) = std::str::from_utf8(bytes) {
6100 probes.push(s.to_string());
6101 }
6102 };
6103 for (first, last) in &bounds {
6104 for base in [first, last] {
6105 push(base, &mut probes); let mut after = base.clone();
6107 after.push(0x01); push(&after, &mut probes);
6109 let mut far = base.clone();
6110 far.push(b'~');
6111 push(&far, &mut probes);
6112 if !base.is_empty() {
6113 push(&base[..base.len() - 1], &mut probes); let mut down = base.clone();
6115 *down.last_mut().unwrap() = down.last().unwrap().wrapping_sub(1);
6116 push(&down, &mut probes);
6117 let mut up = base.clone();
6118 *up.last_mut().unwrap() = up.last().unwrap().wrapping_add(1);
6119 push(&up, &mut probes);
6120 }
6121 }
6122 }
6123 for e in &es {
6125 push(&e.key, &mut probes);
6126 let mut plus = e.key.clone();
6127 plus.push(0);
6128 push(&plus, &mut probes);
6129 }
6130 probes.push(String::new()); probes.push("\u{1}".to_string()); probes.push("\u{10FFFF}\u{10FFFF}".to_string()); probes.push(terms.first().unwrap().clone());
6134 probes.push(terms.last().unwrap().clone());
6135
6136 let mut hits = 0usize;
6137 for p in &probes {
6138 let (a, b) = (sec_v.id(p), sec_s.id(p));
6139 assert_eq!(
6140 a,
6141 b,
6142 "probe diverged: {:?}",
6143 p.get(..p.len().min(48)).unwrap_or(p.as_str())
6144 );
6145 hits += usize::from(a.is_some());
6146 }
6147 assert!(
6148 probes.len() > 500 && hits > 50,
6149 "probe set is too thin: {} probes, {hits} of them present",
6150 probes.len()
6151 );
6152 eprintln!(
6153 "{} boundary probes, {hits} present, all identical",
6154 probes.len()
6155 );
6156 }
6157
6158 mod separator_props {
6164 use super::*;
6165 use proptest::prelude::*;
6166
6167 fn term_set() -> impl Strategy<Value = Vec<String>> {
6174 prop::collection::vec(("[a-c]{1,6}", 0u8..3, 300usize..1200), 60..250).prop_map(
6175 |specs| {
6176 let mut terms: Vec<String> = specs
6177 .into_iter()
6178 .enumerate()
6179 .map(|(i, (core, kind, pad))| {
6180 let block = "z".repeat(pad);
6181 match kind {
6182 0 => format!("<http://ex/{core}/{i:05}>"),
6183 1 => format!("\"{block}{core}{i:05}\""),
6184 _ => format!("\"{core}{i:05}{block}\""),
6185 }
6186 })
6187 .collect();
6188 terms.sort();
6189 terms.dedup();
6190 terms
6191 },
6192 )
6193 }
6194
6195 proptest! {
6196 #[test]
6197 fn prop_separator_keyed_section_answers_like_verbatim(
6198 terms in term_set(),
6199 extra in prop::collection::vec(any::<String>(), 0..8),
6200 ) {
6201 let raw = chunked_section_of(&terms);
6202 let sep = encode_chunked_dict_section(&raw, CODEC_NONE);
6203 let ver = verbatim_keyed_twin(&sep, CODEC_NONE);
6204 prop_assert!(sep.len() <= ver.len());
6205
6206 let (_m, es) = parse_chunked_dict_dir(&sep, sep.len() as u64).unwrap();
6208 let bounds = chunk_bounds_of(&sep, CODEC_NONE);
6209 prop_assert!(es[0].key.is_empty());
6210 for i in 1..es.len() {
6211 prop_assert!(es[i].key.as_slice() > bounds[i - 1].1.as_slice());
6212 prop_assert!(es[i].key <= bounds[i].0);
6213 prop_assert!(es[i].key > es[i - 1].key);
6214 }
6215
6216 let sec_s = decode_chunked_dict_section(
6217 &sep,
6218 CODEC_NONE,
6219 crate::chunk_cache::ChunkCache::unlimited_arc(),
6220 0,
6221 )
6222 .unwrap();
6223 let sec_v = decode_chunked_dict_section(
6224 &ver,
6225 CODEC_NONE,
6226 crate::chunk_cache::ChunkCache::unlimited_arc(),
6227 0,
6228 )
6229 .unwrap();
6230
6231 for (i, t) in terms.iter().enumerate() {
6232 prop_assert_eq!(sec_s.id(t), Some(i as u32 + 1));
6233 prop_assert_eq!(sec_v.id(t), Some(i as u32 + 1));
6234 prop_assert_eq!(sec_s.term(i as u32 + 1), sec_v.term(i as u32 + 1));
6235 }
6236
6237 let mut probes: Vec<String> = extra;
6239 probes.push(String::new());
6240 for (first, last) in &bounds {
6241 for base in [first, last] {
6242 if let Ok(s) = std::str::from_utf8(base) {
6243 probes.push(s.to_string());
6244 probes.push(format!("{s}\u{1}"));
6245 probes.push(s[..s.len() - s.chars().next_back()
6246 .map(char::len_utf8).unwrap_or(0)].to_string());
6247 }
6248 }
6249 }
6250 for e in &es {
6251 if let Ok(s) = std::str::from_utf8(&e.key) {
6252 probes.push(s.to_string());
6253 }
6254 }
6255 for p in &probes {
6256 prop_assert_eq!(sec_s.id(p), sec_v.id(p), "probe diverged: {:?}", p);
6257 }
6258 }
6259 }
6260 }
6261
6262 #[test]
6267 fn text_index_is_tamper_evident_and_verifies() {
6268 let triples: Vec<(String, String, String)> = vec![(
6269 "<http://ex/s0>".to_string(),
6270 "<http://ex/label>".to_string(),
6271 "\"alpha glucose phosphate\"".to_string(),
6272 )];
6273 let bytes = build_text_indexed(&triples);
6274 let header = Rete::open(&bytes).unwrap().header().clone();
6275 assert!(header.text_index_len > 0);
6276 assert!(verify(&bytes).unwrap(), "a text-indexed build must verify");
6277
6278 let mut tampered = bytes.clone();
6279 tampered[header.text_index_offset as usize] ^= 0xff;
6280 assert!(
6281 !verify(&tampered).unwrap(),
6282 "tampering with the text index must break verify()"
6283 );
6284 }
6285
6286 #[test]
6290 fn synopsis_cuts_remote_fetch_bytes() {
6291 use crate::header::FLAG_TILE_SYNOPSIS;
6292 use crate::reader::{CountingReader, SliceReader};
6293
6294 let triples: Vec<(String, String, String)> = (0..400u32)
6298 .map(|i| {
6299 (
6300 format!("<http://ex/s/{i:04}>"),
6301 "<http://ex/p>".to_string(),
6302 format!("<http://ex/o/{i:04}>"),
6303 )
6304 })
6305 .collect();
6306 let mut db = DictionaryBuilder::new();
6307 for (s, p, o) in &triples {
6308 db.observe(s, p, o);
6309 }
6310 let dict = db.build();
6311 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
6312 for (s, p, o) in &triples {
6313 ib.push(dict.encode(s, p, o).unwrap());
6314 }
6315 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
6316
6317 let q = (Some("<http://ex/s/0395>"), None, Some("<http://ex/o/0005>"));
6321 let query_bytes = |image: &[u8]| -> (u64, usize) {
6325 let leaked: &'static [u8] = Box::leak(image.to_vec().into_boxed_slice());
6326 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
6327 let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
6328 let before = reader.bytes_read(); let n = rete.query(q.0, q.1, q.2).len();
6330 assert!(!rete.index_incomplete());
6331 (reader.bytes_read() - before, n)
6332 };
6333
6334 let (on_bytes, on_n) = query_bytes(&bytes);
6335 let mut off = bytes.clone();
6337 off[5] &= !FLAG_TILE_SYNOPSIS;
6338 let (off_bytes, off_n) = query_bytes(&off);
6339
6340 assert_eq!(on_n, 0, "the pair never co-occurs");
6341 assert_eq!(off_n, 0, "same answer without the synopsis");
6342 assert!(
6346 on_bytes < off_bytes,
6347 "synopsis skips the routed tile fetch: {on_bytes} < {off_bytes}"
6348 );
6349 }
6350
6351 #[test]
6357 fn double_bound_object_join_eager_matches_lazy() {
6358 use crate::reader::SliceReader;
6359 let occ = "<http://ex/occ>";
6360 let phys = "<http://ex/physicist>";
6361 let phil = "<http://ex/philosopher>";
6362 let label = "<http://www.w3.org/2000/01/rdf-schema#label>";
6363 let mut triples: Vec<(String, String, String)> = Vec::new();
6365 for i in 0..20u32 {
6366 triples.push((format!("<http://ex/p/{i:02}>"), occ.into(), phys.into()));
6367 if i < 10 {
6368 triples.push((format!("<http://ex/p/{i:02}>"), occ.into(), phil.into()));
6369 }
6370 triples.push((
6371 format!("<http://ex/p/{i:02}>"),
6372 label.into(),
6373 format!("\"Name {i:02}\""),
6374 ));
6375 }
6376 let mut db = DictionaryBuilder::new();
6377 for (s, p, o) in &triples {
6378 db.observe(s, p, o);
6379 }
6380 let dict = db.build();
6381 let mut ib = GraphIndexBuilder::new().with_tile_budget(16);
6382 for (s, p, o) in &triples {
6383 ib.push(dict.encode(s, p, o).unwrap());
6384 }
6385 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
6386
6387 let q = "SELECT ?l WHERE { \
6388 ?p <http://ex/occ> <http://ex/physicist> ; \
6389 <http://ex/occ> <http://ex/philosopher> ; \
6390 <http://www.w3.org/2000/01/rdf-schema#label> ?l }";
6391 let run = |rete: &Rete| -> Vec<String> {
6392 let (_, sols) = crate::eval_sparql(rete, q).unwrap();
6393 let mut v: Vec<String> = sols.iter().filter_map(|b| b.get("l").cloned()).collect();
6394 v.sort();
6395 v
6396 };
6397
6398 let eager_rows = run(&Rete::open(&bytes).unwrap());
6399 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
6400 let lazy = Rete::open_ranged_lazy(std::sync::Arc::new(SliceReader::new(leaked))).unwrap();
6401 let lazy_rows = run(&lazy);
6402 assert!(!lazy.index_incomplete());
6403
6404 assert_eq!(eager_rows.len(), 10, "the 10 physicist∩philosopher labels");
6405 assert_eq!(eager_rows, lazy_rows, "eager and lazy must agree exactly");
6406 }
6407
6408 #[test]
6412 fn multi_chunk_dictionary_round_trips() {
6413 let mut db = DictionaryBuilder::new();
6414 let term = |i: u32| format!("<http://example.org/some/long/prefix/entity/{i:06}>");
6415 for i in 0..6000u32 {
6416 db.observe(&term(i), "<http://ex/p>", &term(i + 1));
6417 }
6418 let dict = db.build();
6419 let mut ib = GraphIndexBuilder::new();
6420 for i in 0..6000u32 {
6421 ib.push(
6422 dict.encode(&term(i), "<http://ex/p>", &term(i + 1))
6423 .unwrap(),
6424 );
6425 }
6426 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
6427 let rete = Rete::open(&bytes).unwrap();
6428 let d = rete.dictionary();
6429 assert_eq!(d.term_count(), dict.term_count());
6430 for i in (0..6000).step_by(97).chain([0, 1, 5999, 6000]) {
6431 let t = term(i);
6432 let sid = dict.subject_id(&t);
6433 assert_eq!(d.subject_id(&t), sid, "subject_id({t})");
6434 if let Some(id) = sid {
6435 assert_eq!(d.subject_term(id).as_deref(), Some(t.as_str()));
6436 }
6437 let oid = dict.object_id(&t);
6438 assert_eq!(d.object_id(&t), oid, "object_id({t})");
6439 }
6440 assert_eq!(d.subject_id("<http://example.org/absent>"), None);
6441 assert_eq!(d.predicate_id("<http://ex/p>"), Some(1));
6442 assert_eq!(d.predicate_term(1).as_deref(), Some("<http://ex/p>"));
6443 }
6444
6445 #[test]
6446 #[cfg(feature = "compression")]
6447 fn compression_shrinks_repetitive_data() {
6448 let mut db = DictionaryBuilder::new();
6452 let triples: Vec<(String, String, String)> = (0..500)
6453 .map(|i| {
6454 (
6455 format!("<http://example.org/entity/{i}>"),
6456 "<http://example.org/p/relatedTo>".to_string(),
6457 format!("<http://example.org/entity/{}>", (i + 1) % 500),
6458 )
6459 })
6460 .collect();
6461 for (s, p, o) in &triples {
6462 db.observe(s, p, o);
6463 }
6464 let dict = db.build();
6465 let mut ib = GraphIndexBuilder::new();
6466 for (s, p, o) in &triples {
6467 ib.push(dict.encode(s, p, o).unwrap());
6468 }
6469 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
6470
6471 let raw: usize = triples
6472 .iter()
6473 .map(|(s, p, o)| s.len() + p.len() + o.len())
6474 .sum();
6475 assert!(
6476 bytes.len() < raw / 2,
6477 "expected strong compression: file {} vs raw terms {raw}",
6478 bytes.len()
6479 );
6480
6481 let rete = Rete::open(&bytes).unwrap();
6483 let r = rete.query(Some("<http://example.org/entity/0>"), None, None);
6484 assert_eq!(r.len(), 1);
6485 assert_eq!(r[0].2, "<http://example.org/entity/1>");
6486 }
6487
6488 fn big_file_with_pyramid() -> Vec<u8> {
6489 let triples: Vec<(String, String, String)> = (0..300)
6491 .map(|i| {
6492 (
6493 format!("<http://ex/e{i}>"),
6494 "<http://ex/next>".to_string(),
6495 format!("<http://ex/e{}>", (i + 1) % 300),
6496 )
6497 })
6498 .collect();
6499 let mut db = DictionaryBuilder::new();
6500 for (s, p, o) in &triples {
6501 db.observe(s, p, o);
6502 }
6503 let dict = db.build();
6504 let ids: Vec<_> = triples
6505 .iter()
6506 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
6507 .collect();
6508 let mut ib = GraphIndexBuilder::new();
6509 for &t in &ids {
6510 ib.push(t);
6511 }
6512 let (meta, levels) = build_pyramid_meta(&dict, &ids, DEFAULT_TILE_BUDGET);
6513 write_file(&dict, &ib.build(), false, &meta, levels)
6514 }
6515
6516 #[test]
6517 fn ranged_open_is_minimal_and_correct() {
6518 use crate::reader::{CountingReader, SliceReader};
6519 let bytes = big_file_with_pyramid();
6520
6521 let full = CountingReader::new(SliceReader::new(&bytes));
6522 let rete = Rete::open_ranged(&full).unwrap();
6523 assert!(full.requests() <= 4, "requests = {}", full.requests());
6525 assert_eq!(
6526 rete.query(Some("<http://ex/e0>"), None, None)[0].2,
6527 "<http://ex/e1>"
6528 );
6529
6530 let summ_reader = CountingReader::new(SliceReader::new(&bytes));
6532 let view = SummaryView::open_ranged(&summ_reader).unwrap().unwrap();
6533 assert!(!view.summary.is_empty());
6534 assert!(
6535 summ_reader.bytes_read() < bytes.len() as u64,
6536 "summary read {} of {} bytes",
6537 summ_reader.bytes_read(),
6538 bytes.len()
6539 );
6540 assert!(summ_reader.bytes_read() < full.bytes_read());
6542 }
6543
6544 #[test]
6545 fn content_hash_is_set_and_verifies() {
6546 let bytes = build_image();
6547 let rete = Rete::open(&bytes).unwrap();
6548 assert_ne!(
6549 rete.header().content_hash,
6550 [0u8; 16],
6551 "hash must be populated"
6552 );
6553 assert!(verify(&bytes).unwrap(), "freshly built file verifies");
6554
6555 assert_eq!(
6557 Rete::open(&build_image()).unwrap().header().content_hash,
6558 rete.header().content_hash
6559 );
6560
6561 let mut tampered = bytes.clone();
6563 let last = tampered.len() - 5; tampered[last] ^= 0xff;
6565 assert!(!verify(&tampered).unwrap());
6566 }
6567
6568 fn build_with_metadata(meta: &[u8]) -> Vec<u8> {
6570 let triples = [
6571 ("Alice", "knows", "Bob"),
6572 ("Bob", "knows", "Carol"),
6573 ("Alice", "age", "30"),
6574 ];
6575 let mut db = DictionaryBuilder::new();
6576 for (s, p, o) in triples {
6577 db.observe(s, p, o);
6578 }
6579 let dict = db.build();
6580 let ids: Vec<_> = triples
6581 .iter()
6582 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
6583 .collect();
6584 let mut ib = GraphIndexBuilder::new();
6585 for &t in &ids {
6586 ib.push(t);
6587 }
6588 let (pmeta, levels) = build_pyramid_meta(&dict, &ids, DEFAULT_TILE_BUDGET);
6589 write_dataset_with_metadata(&dict, &ib.build(), &[], false, &pmeta, levels, meta, &[])
6590 }
6591
6592 #[test]
6593 fn metadata_round_trips_and_shifts_offsets() {
6594 let card = br#"{"title":"My Dataset"}"#;
6595 let bytes = build_with_metadata(card);
6596 let rete = Rete::open(&bytes).unwrap();
6597
6598 assert_eq!(rete.metadata(), Some(card.as_slice()));
6600 let h = rete.header();
6601 assert_eq!(h.metadata_offset, HEADER_LEN as u64);
6602 assert_eq!(h.metadata_len, card.len() as u64);
6603 assert_eq!(h.dictionary_offset, HEADER_LEN as u64 + card.len() as u64);
6605
6606 assert_eq!(
6608 rete.query(Some("Bob"), Some("knows"), Some("Carol")).len(),
6609 1
6610 );
6611 assert!(verify(&bytes).unwrap());
6613 }
6614
6615 #[test]
6616 fn empty_metadata_is_byte_identical_to_plain_writer() {
6617 assert_eq!(
6620 build_with_metadata(&[]),
6621 build_image(),
6622 "empty-metadata output must equal the plain writer byte-for-byte"
6623 );
6624 }
6625
6626 #[test]
6631 fn replace_metadata_matches_a_direct_write() {
6632 let first = br#"{"title":"My Dataset","queries":[1,2,3]}"#;
6633 let bytes = build_with_metadata(first);
6634 for target in [
6635 br#"{"title":"My Dataset","queries":[1,2]}"#.as_slice(), br#"{"title":"My Dataset","queries":[1,2,3,4,5,6,7]}"#.as_slice(), b"".as_slice(), first.as_slice(), ] {
6640 let spliced = replace_metadata(&bytes, target).unwrap();
6641 assert_eq!(
6642 spliced,
6643 build_with_metadata(target),
6644 "splicing {} bytes of metadata must equal a direct write",
6645 target.len()
6646 );
6647 assert!(
6648 verify(&spliced).unwrap(),
6649 "the new hash covers the new card"
6650 );
6651 let rete = Rete::open(&spliced).unwrap();
6652 assert_eq!(
6653 rete.metadata(),
6654 (!target.is_empty()).then_some(target),
6655 "the new payload reads back verbatim"
6656 );
6657 assert_eq!(
6659 rete.query(Some("Bob"), Some("knows"), Some("Carol")).len(),
6660 1
6661 );
6662 }
6663 }
6664
6665 #[test]
6668 fn replace_metadata_keeps_build_info_adjacent() {
6669 let with_info =
6670 attach_build_info(&build_with_metadata(br#"{"a":1}"#), br#"{"builder":"x"}"#).unwrap();
6671 let spliced = replace_metadata(&with_info, br#"{"a":1,"b":2}"#).unwrap();
6672 let h = Header::from_bytes(&spliced).unwrap();
6673 assert_eq!(h.metadata_offset, HEADER_LEN as u64);
6674 assert_eq!(h.build_info_offset, HEADER_LEN as u64 + h.metadata_len);
6675 assert_eq!(
6676 read_build_info(&spliced).unwrap().as_deref(),
6677 Some(br#"{"builder":"x"}"#.as_slice())
6678 );
6679 assert!(verify(&spliced).unwrap());
6680 assert_eq!(
6682 h.content_hash,
6683 Header::from_bytes(&build_with_metadata(br#"{"a":1,"b":2}"#))
6684 .unwrap()
6685 .content_hash
6686 );
6687 }
6688
6689 #[test]
6690 fn metadata_is_tamper_evident() {
6691 let card = br#"{"title":"x"}"#;
6692 let mut bytes = build_with_metadata(card);
6693 assert!(verify(&bytes).unwrap());
6694 bytes[HEADER_LEN + 2] ^= 0xff;
6696 assert!(
6697 !verify(&bytes).unwrap(),
6698 "tampering with the card must break verify()"
6699 );
6700 }
6701
6702 #[test]
6703 fn ranged_opens_do_not_fetch_metadata() {
6704 use crate::reader::{CountingReader, SliceReader};
6705 let card = vec![0xABu8; 512]; let bytes = build_with_metadata(&card);
6707 let total = bytes.len() as u64;
6708
6709 let r = CountingReader::new(SliceReader::new(&bytes));
6711 let rete = Rete::open_ranged(&r).unwrap();
6712 assert!(
6713 rete.metadata().is_none(),
6714 "open_ranged must not load the card"
6715 );
6716 assert!(r.requests() <= 4, "requests = {}", r.requests());
6717 assert!(
6718 r.bytes_read() <= total - card.len() as u64,
6719 "read {} of {} bytes; the {}-byte card must be skipped",
6720 r.bytes_read(),
6721 total,
6722 card.len()
6723 );
6724
6725 let rs = CountingReader::new(SliceReader::new(&bytes));
6727 let view = SummaryView::open_ranged(&rs).unwrap().unwrap();
6728 assert!(!view.summary.is_empty());
6729 assert!(rs.bytes_read() <= total - card.len() as u64);
6730 }
6731
6732 #[test]
6733 fn metadata_ranged_fetches_only_header_and_card() {
6734 use crate::reader::{CountingReader, SliceReader};
6735 let card = vec![0xCDu8; 384];
6738 let bytes = build_with_metadata(&card);
6739
6740 let r = CountingReader::new(SliceReader::new(&bytes));
6741 let got = read_metadata_ranged(&r).unwrap().unwrap();
6742 assert_eq!(got, card, "the card reads back verbatim");
6743 assert_eq!(r.requests(), 2, "exactly header + metadata ranges");
6744 assert_eq!(
6745 r.bytes_read(),
6746 HEADER_LEN as u64 + card.len() as u64,
6747 "no dictionary/index/pyramid bytes are touched"
6748 );
6749
6750 let plain = build_image();
6752 let rp = CountingReader::new(SliceReader::new(&plain));
6753 assert!(read_metadata_ranged(&rp).unwrap().is_none());
6754 assert_eq!(rp.requests(), 1, "header only for a cardless file");
6755 assert_eq!(rp.bytes_read(), HEADER_LEN as u64);
6756 }
6757
6758 #[test]
6759 fn build_info_attaches_outside_the_hash() {
6760 let card = br#"{"title":"My Dataset"}"#;
6761 let base = build_with_metadata(card);
6762 let base_hash = Header::from_bytes(&base).unwrap().content_hash;
6763
6764 let info = br#"{"built_at":"2026-08-04T00:00:00Z","builder":"rete-cli 0.3.2"}"#;
6765 let with = attach_build_info(&base, info).unwrap();
6766
6767 assert_eq!(read_build_info(&with).unwrap().unwrap(), info);
6769 let h = Header::from_bytes(&with).unwrap();
6770 assert_eq!(h.build_info_offset, HEADER_LEN as u64 + card.len() as u64);
6771 assert_eq!(h.build_info_len, info.len() as u64);
6772 assert_eq!(h.expected_file_len(), Some(with.len() as u64));
6773
6774 assert_eq!(h.content_hash, base_hash);
6778 assert!(verify(&with).unwrap());
6779 let other = attach_build_info(&base, br#"{"built_at":"1999-01-01T00:00:00Z"}"#).unwrap();
6780 assert_eq!(Header::from_bytes(&other).unwrap().content_hash, base_hash);
6781
6782 let rete = Rete::open(&with).unwrap();
6784 assert_eq!(
6785 rete.query(Some("Bob"), Some("knows"), Some("Carol")).len(),
6786 1
6787 );
6788 assert_eq!(rete.metadata(), Some(card.as_slice()));
6789
6790 let replaced = attach_build_info(&with, br#"{"builder":"other"}"#).unwrap();
6793 assert_eq!(
6794 read_build_info(&replaced).unwrap().unwrap(),
6795 br#"{"builder":"other"}"#
6796 );
6797 assert!(verify(&replaced).unwrap());
6798 let stripped = attach_build_info(&replaced, &[]).unwrap();
6799 assert_eq!(stripped, base, "stripping build-info restores the image");
6800 }
6801
6802 #[test]
6803 fn build_info_attaches_on_a_metadata_free_file() {
6804 let base = build_image();
6806 let info = b"{\"builder\":\"x\"}";
6807 let with = attach_build_info(&base, info).unwrap();
6808 let h = Header::from_bytes(&with).unwrap();
6809 assert_eq!(h.build_info_offset, HEADER_LEN as u64);
6810 assert_eq!(h.metadata_len, 0);
6811 assert_eq!(read_build_info(&with).unwrap().unwrap(), info);
6812 assert!(verify(&with).unwrap());
6813 assert_eq!(
6814 h.content_hash,
6815 Header::from_bytes(&base).unwrap().content_hash
6816 );
6817 Rete::open(&with).unwrap();
6818 }
6819
6820 #[test]
6827 fn a_plan_rebuilds_exactly_what_the_in_memory_splice_produces() {
6828 let card = br#"{"title":"My Dataset"}"#;
6829 let base = build_with_metadata(card);
6830 let base_hash = Header::from_bytes(&base).unwrap().content_hash;
6831 for start in [
6832 base.clone(),
6833 attach_build_info(&base, b"{\"a\":1}").unwrap(),
6834 ] {
6835 for info in [
6836 br#"{"builder":"rete-cli 0.3.2","query_costs":{}}"#.as_slice(), b"{}".as_slice(), b"".as_slice(), ] {
6840 let plan = plan_build_info(&start, start.len() as u64, info.len() as u64).unwrap();
6841 let mut streamed = Vec::new();
6842 streamed.extend_from_slice(&plan.header);
6843 streamed.extend_from_slice(&start[HEADER_LEN..plan.insert as usize]);
6844 streamed.extend_from_slice(info);
6845 streamed.extend_from_slice(&start[plan.tail_start as usize..]);
6846 assert_eq!(streamed.len() as u64, plan.new_len);
6847 assert_eq!(
6848 streamed,
6849 attach_build_info(&start, info).unwrap(),
6850 "a {}-byte section streamed must equal the same section spliced",
6851 info.len()
6852 );
6853 assert!(verify(&streamed).unwrap());
6855 let h = Header::from_bytes(&streamed).unwrap();
6856 assert_eq!(h.content_hash, base_hash);
6857 assert_eq!(h.expected_file_len(), Some(plan.new_len));
6858 let rete = Rete::open(&streamed).unwrap();
6859 assert_eq!(rete.metadata(), Some(card.as_slice()));
6860 assert_eq!(
6861 rete.query(Some("Bob"), Some("knows"), Some("Carol")).len(),
6862 1
6863 );
6864 }
6865 }
6866 }
6867
6868 #[test]
6871 fn a_plan_refuses_a_header_that_overruns_the_file() {
6872 let base = build_with_metadata(br#"{"title":"x"}"#);
6873 assert!(plan_build_info(&base, 8, 16).is_err());
6874 assert!(plan_build_info(&base[..HEADER_LEN], base.len() as u64, 16).is_ok());
6875 }
6876
6877 #[test]
6878 fn card_and_build_info_ranged_is_one_header_plus_one_range() {
6879 use crate::reader::{CountingReader, SliceReader};
6880 let card = vec![0xCDu8; 384];
6881 let info = vec![0xEFu8; 200];
6882 let bytes = attach_build_info(&build_with_metadata(&card), &info).unwrap();
6883
6884 let r = CountingReader::new(SliceReader::new(&bytes));
6887 let (m, b) = read_card_and_build_info_ranged(&r).unwrap();
6888 assert_eq!(m.unwrap(), card);
6889 assert_eq!(b.unwrap(), info);
6890 assert_eq!(r.requests(), 2, "header + one coalesced range");
6891 assert_eq!(
6892 r.bytes_read(),
6893 HEADER_LEN as u64 + card.len() as u64 + info.len() as u64
6894 );
6895
6896 let plain = build_with_metadata(&card);
6898 let rp = CountingReader::new(SliceReader::new(&plain));
6899 let (m, b) = read_card_and_build_info_ranged(&rp).unwrap();
6900 assert_eq!(m.unwrap(), card);
6901 assert!(b.is_none());
6902 assert_eq!(rp.requests(), 2);
6903 assert_eq!(
6908 r.requests(),
6909 rp.requests(),
6910 "build info must not cost an extra request",
6911 );
6912
6913 let none = build_image();
6915 let rn = CountingReader::new(SliceReader::new(&none));
6916 let (m, b) = read_card_and_build_info_ranged(&rn).unwrap();
6917 assert!(m.is_none() && b.is_none());
6918 assert_eq!(rn.requests(), 1);
6919 }
6920
6921 #[test]
6922 fn schema_summary_groups_by_type() {
6923 let rt = RDF_TYPE;
6924 let bytes = build_from(&[
6925 ("Alice", rt, "Person"),
6926 ("Bob", rt, "Person"),
6927 ("NYC", rt, "City"),
6928 ("Alice", "knows", "Bob"),
6929 ("Alice", "livesIn", "NYC"),
6930 ("Alice", "name", "\"Alice\""),
6931 ]);
6932 let rete = Rete::open(&bytes).unwrap();
6933 let summary = schema_summary(&rete);
6934 assert!(summary.contains(&("Person".into(), "knows".into(), "Person".into(), 1)));
6936 assert!(summary.contains(&("Person".into(), "livesIn".into(), "City".into(), 1)));
6937 assert!(summary.contains(&("Person".into(), "name".into(), "(literal)".into(), 1)));
6938 assert!(!summary.iter().any(|(_, p, _, _)| p == RDF_TYPE));
6940
6941 let classes = schema_classes(&rete);
6943 assert_eq!(
6944 classes,
6945 vec![("Person".into(), 2u32), ("City".into(), 1u32)]
6946 );
6947 }
6948
6949 fn build_from(triples: &[(&str, &str, &str)]) -> Vec<u8> {
6950 let mut db = DictionaryBuilder::new();
6951 for (s, p, o) in triples {
6952 db.observe(s, p, o);
6953 }
6954 let dict = db.build();
6955 let mut ib = GraphIndexBuilder::new();
6956 for (s, p, o) in triples {
6957 ib.push(dict.encode(s, p, o).unwrap());
6958 }
6959 write_file(&dict, &ib.build(), false, &[], 0)
6960 }
6961
6962 fn build_with_pyramid(triples: &[(&str, &str, &str)]) -> Vec<u8> {
6964 let mut db = DictionaryBuilder::new();
6965 for (s, p, o) in triples {
6966 db.observe(s, p, o);
6967 }
6968 let dict = db.build();
6969 let encoded: Vec<_> = triples
6970 .iter()
6971 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
6972 .collect();
6973 let mut ib = GraphIndexBuilder::new();
6974 for t in &encoded {
6975 ib.push(*t);
6976 }
6977 let (meta, levels) = build_pyramid_meta(&dict, &encoded, DEFAULT_TILE_BUDGET);
6978 write_dataset(&dict, &ib.build(), &[], false, &meta, levels)
6979 }
6980
6981 #[test]
6982 fn tbox_coherence_flags_unsatisfiable_class_index_free() {
6983 use crate::reader::{CountingReader, SliceReader};
6984 let rt = RDF_TYPE;
6985 let sub = "<http://www.w3.org/2000/01/rdf-schema#subClassOf>";
6986 let disj = "<http://www.w3.org/2002/07/owl#disjointWith>";
6987 let bytes = build_with_pyramid(&[
6990 ("<http://ex/C>", sub, "<http://ex/D>"),
6991 ("<http://ex/C>", sub, "<http://ex/E>"),
6992 ("<http://ex/D>", disj, "<http://ex/E>"),
6993 ("<http://ex/x>", rt, "<http://ex/C>"),
6994 ]);
6995
6996 let r = CountingReader::new(SliceReader::new(&bytes));
6997 let view = SummaryView::open_ranged(&r).unwrap().unwrap();
6998 let points = view.tbox_coherence();
6999 assert!(
7000 points
7001 .iter()
7002 .any(|i| i.kind == "unsatisfiable-class" && i.detail.contains("http://ex/C>")),
7003 "expected C unsatisfiable from the schema alone, got {points:?}"
7004 );
7005
7006 let header = Header::from_bytes(&bytes[..HEADER_LEN]).unwrap();
7009 assert!(
7010 r.bytes_read() <= bytes.len() as u64 - header.root_dir_len,
7011 "tbox_coherence must not read the triple index"
7012 );
7013 }
7014
7015 #[test]
7016 fn schema_coherence_reads_only_the_schema_block() {
7017 use crate::reader::{CountingReader, SliceReader};
7018 let rt = RDF_TYPE;
7019 let sub = "<http://www.w3.org/2000/01/rdf-schema#subClassOf>";
7020 let disj = "<http://www.w3.org/2002/07/owl#disjointWith>";
7021 let mut triples: Vec<(String, String, String)> = vec![
7025 ("<http://ex/C>".into(), sub.into(), "<http://ex/D>".into()),
7026 ("<http://ex/C>".into(), sub.into(), "<http://ex/E>".into()),
7027 ("<http://ex/D>".into(), disj.into(), "<http://ex/E>".into()),
7028 ];
7029 for i in 0..500 {
7030 let s = format!("<http://ex/x{i}>");
7031 triples.push((s.clone(), rt.into(), "<http://ex/C>".into()));
7032 triples.push((
7033 s,
7034 "<http://ex/label>".into(),
7035 format!("\"unique label {i}\""),
7036 ));
7037 }
7038 let trefs: Vec<(&str, &str, &str)> = triples
7039 .iter()
7040 .map(|(s, p, o)| (s.as_str(), p.as_str(), o.as_str()))
7041 .collect();
7042 let bytes = build_with_pyramid(&trefs);
7043
7044 let header = Header::from_bytes(&bytes[..HEADER_LEN]).unwrap();
7045 assert!(
7046 header.schema_meta_len > 0,
7047 "the writer recorded a schema-block length"
7048 );
7049 assert!(
7050 (header.schema_meta_len as u64) < header.pyramid_meta_len,
7051 "schema block ({}) should be far smaller than the whole pyramid-meta ({})",
7052 header.schema_meta_len,
7053 header.pyramid_meta_len
7054 );
7055
7056 let r = CountingReader::new(SliceReader::new(&bytes));
7057 let points = read_schema_coherence_ranged(&r).unwrap().unwrap();
7058 assert!(points.iter().any(|i| i.kind == "unsatisfiable-class"));
7059 assert!(
7061 r.bytes_read() <= HEADER_LEN as u64 + header.schema_meta_len as u64,
7062 "read {} bytes; expected <= header + schema block ({})",
7063 r.bytes_read(),
7064 HEADER_LEN as u64 + header.schema_meta_len as u64
7065 );
7066 }
7067
7068 #[test]
7069 fn tbox_coherence_clean_schema_is_coherent() {
7070 use crate::reader::SliceReader;
7071 let rt = RDF_TYPE;
7072 let sub = "<http://www.w3.org/2000/01/rdf-schema#subClassOf>";
7073 let bytes = build_with_pyramid(&[
7074 ("<http://ex/Dog>", sub, "<http://ex/Animal>"),
7075 ("<http://ex/x>", rt, "<http://ex/Dog>"),
7076 ]);
7077 let view = SummaryView::open_ranged(&SliceReader::new(&bytes))
7078 .unwrap()
7079 .unwrap();
7080 assert!(view.tbox_is_coherent(), "a plain hierarchy is coherent");
7081 }
7082
7083 #[test]
7084 fn named_graphs_round_trip() {
7085 let all = [
7087 ("Alice", "knows", "Bob"), ("Bob", "age", "30"), ];
7090 let mut db = DictionaryBuilder::new();
7091 for (s, p, o) in all {
7092 db.observe(s, p, o);
7093 }
7094 let dict = db.build();
7095
7096 let mut def = GraphIndexBuilder::new();
7097 def.push(dict.encode("Alice", "knows", "Bob").unwrap());
7098 let mut g1 = GraphIndexBuilder::new();
7099 g1.push(dict.encode("Bob", "age", "30").unwrap());
7100
7101 let named = vec![("http://ex/g1".to_string(), g1.build())];
7102 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
7103
7104 assert!(verify(&bytes).unwrap());
7105 let rete = Rete::open(&bytes).unwrap();
7106 assert_eq!(rete.graph_names(), vec!["http://ex/g1"]);
7107 let gi = rete.graph_index("http://ex/g1").unwrap();
7109 assert_eq!(gi.triple_count(), 1);
7110 assert!(rete.graph_index("http://ex/missing").is_none());
7111
7112 assert_eq!(rete.header().quad_count, 2);
7115
7116 assert_eq!(rete.query(Some("Alice"), None, None).len(), 1);
7118
7119 assert_eq!(
7121 rete.dump(None),
7122 vec![("Alice".into(), "knows".into(), "Bob".into())]
7123 );
7124 assert_eq!(
7125 rete.dump(Some("http://ex/g1")),
7126 vec![("Bob".into(), "age".into(), "30".into())]
7127 );
7128 }
7129
7130 #[test]
7136 fn named_graph_container_opens_tile_lazily_like_the_root() {
7137 let mut db = DictionaryBuilder::new();
7138 let node = |n: u32| format!("<http://ex/n{n}>");
7139 let p = "<http://ex/p>".to_string();
7140 db.observe(&node(0), &p, &node(1));
7141 for i in 0..500u32 {
7142 db.observe(&node(1000 + i), &p, &node(2000 + i));
7143 }
7144 let dict = db.build();
7145 let mut def = GraphIndexBuilder::new();
7146 def.push(dict.encode(&node(0), &p, &node(1)).unwrap());
7147 let mut g = GraphIndexBuilder::new().with_tile_budget(256); for i in 0..500u32 {
7149 g.push(dict.encode(&node(1000 + i), &p, &node(2000 + i)).unwrap());
7150 }
7151 let named = vec![("<http://ex/g>".to_string(), g.build())];
7152 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
7153 let rete = Rete::open(&bytes).unwrap();
7154 let header = rete.header();
7155
7156 let soff = header.named_graphs_offset as usize;
7159 let send = soff + header.named_graphs_len as usize;
7160 let sec = &bytes[soff..send];
7161 let (n, mut pos) = read_uvarint(sec).unwrap();
7162 assert_eq!(n, 1);
7163 let (ilen, u1) = read_uvarint(&sec[pos..]).unwrap();
7164 pos += u1 + ilen as usize;
7165 let (clen, u2) = read_uvarint(&sec[pos..]).unwrap();
7166 pos += u2;
7167 let container = ByteRange {
7168 offset: (soff + pos) as u64,
7169 len: clen,
7170 };
7171
7172 use crate::reader::SliceReader;
7173 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
7174 let reader: std::sync::Arc<dyn RangeReader + Send + Sync> =
7175 std::sync::Arc::new(SliceReader::new(leaked));
7176 let (lazy_idx, _, _) = open_index_container_lazy(
7177 &reader,
7178 container,
7179 header.block_codec,
7180 header.has_tile_synopsis(),
7181 1,
7182 header.perms,
7183 )
7184 .unwrap();
7185 let mut want = rete
7186 .graph_index("<http://ex/g>")
7187 .unwrap()
7188 .match_pattern((None, None, None));
7189 let mut got = lazy_idx.match_pattern((None, None, None));
7190 want.sort_unstable();
7191 got.sort_unstable();
7192 assert_eq!(got.len(), 500);
7193 assert_eq!(got, want);
7194 assert!(!lazy_idx.load_incomplete());
7195 }
7196
7197 #[test]
7203 fn dump_iter_matches_dump_and_stops_early() {
7204 use crate::reader::{CountingReader, SliceReader};
7205
7206 let triples: Vec<(String, String, String)> = (0..2000u32)
7208 .map(|i| {
7209 (
7210 format!("<http://ex/s/{i:04}>"),
7211 "<http://ex/p>".to_string(),
7212 format!("<http://ex/o/{i:04}>"),
7213 )
7214 })
7215 .collect();
7216 let mut db = DictionaryBuilder::new();
7217 for (s, p, o) in &triples {
7218 db.observe(s, p, o);
7219 }
7220 let dict = db.build();
7221 let mut ib = GraphIndexBuilder::new().with_tile_budget(256);
7222 for (s, p, o) in &triples {
7223 ib.push(dict.encode(s, p, o).unwrap());
7224 }
7225 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
7226
7227 let rete = Rete::open(&bytes).unwrap();
7229 assert_eq!(rete.dump_iter(None).collect::<Vec<_>>(), rete.dump(None));
7230 assert_eq!(rete.dump_iter(Some("http://ex/nope")).count(), 0);
7232
7233 let read_bytes = |take: usize| -> u64 {
7235 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
7236 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
7237 let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
7238 let before = reader.bytes_read();
7239 let n = rete.dump_iter(None).take(take).count();
7240 assert_eq!(n, take.min(triples.len()));
7241 assert!(!rete.index_incomplete());
7242 reader.bytes_read() - before
7243 };
7244 let first_only = read_bytes(1);
7245 let everything = read_bytes(triples.len());
7246 assert!(
7248 first_only < everything,
7249 "dump_iter is not lazy: one triple read {first_only} B, a full scan {everything} B"
7250 );
7251 #[cfg(feature = "compression")]
7258 assert!(
7259 first_only * 2 < everything,
7260 "dump_iter is not lazy: one triple read {first_only} B of the {everything} B a full scan reads"
7261 );
7262
7263 let fresh = Rete::open(&bytes).unwrap();
7268 let drain = |max_quads: usize| {
7273 let (mut out, mut cursor, mut calls) = (Vec::new(), 0u32, 0);
7274 loop {
7275 let (batch, next, done) = fresh.dump_batch(None, cursor, max_quads);
7276 out.extend(batch);
7277 cursor = next;
7278 calls += 1;
7279 assert!(calls < 10_000, "cursor failed to advance");
7280 if done {
7281 break out;
7282 }
7283 }
7284 };
7285 let eager = fresh.dump(None);
7286 for max in [1usize, 7, 128, 100_000] {
7287 assert_eq!(drain(max), eager, "batch size {max} changed the dump");
7288 }
7289 let (t, _, done) = fresh.dump_batch(Some("http://ex/nope"), 0, 16);
7292 assert!(t.is_empty() && done);
7293 }
7294
7295 #[test]
7303 fn query_batch_matches_query_in_graph_and_stops_early() {
7304 use crate::reader::{CountingReader, SliceReader};
7305
7306 let triples: Vec<(String, String, String)> = (0..2000u32)
7309 .map(|i| {
7310 (
7311 format!("<http://ex/s/{:04}>", i / 2),
7312 format!("<http://ex/p/{}>", i % 2),
7313 format!("<http://ex/o/{i:04}>"),
7314 )
7315 })
7316 .collect();
7317 let mut db = DictionaryBuilder::new();
7318 for (s, p, o) in &triples {
7319 db.observe(s, p, o);
7320 }
7321 let dict = db.build();
7322 let mut ib = GraphIndexBuilder::new().with_tile_budget(256);
7323 for (s, p, o) in &triples {
7324 ib.push(dict.encode(s, p, o).unwrap());
7325 }
7326 let bytes = write_file(&dict, &ib.build(), false, &[], 0);
7327 let rete = Rete::open(&bytes).unwrap();
7328
7329 let drain = |s: Option<&str>, p: Option<&str>, o: Option<&str>, max: usize| {
7330 let (mut out, mut cursor, mut calls) = (Vec::new(), 0u64, 0);
7331 loop {
7332 let (batch, next, done) = rete.query_batch(None, s, p, o, cursor, max);
7333 assert!(!batch.is_empty() || done, "empty non-final batch");
7336 out.extend(batch);
7337 cursor = next;
7338 calls += 1;
7339 assert!(calls < 20_000, "cursor failed to advance");
7340 if done {
7341 break out;
7342 }
7343 }
7344 };
7345 let shapes: [(Option<&str>, Option<&str>, Option<&str>); 6] = [
7346 (None, None, None),
7347 (Some("<http://ex/s/0007>"), None, None),
7348 (None, Some("<http://ex/p/1>"), None),
7349 (None, None, Some("<http://ex/o/0042>")),
7350 (Some("<http://ex/s/0007>"), Some("<http://ex/p/0>"), None),
7351 (
7352 Some("<http://ex/s/0007>"),
7353 Some("<http://ex/p/0>"),
7354 Some("<http://ex/o/0014>"),
7355 ),
7356 ];
7357 for (s, p, o) in shapes {
7358 let eager = rete.query_in_graph(None, s, p, o);
7359 assert!(
7360 !eager.is_empty(),
7361 "fixture shape {s:?} {p:?} {o:?} matches nothing"
7362 );
7363 assert_eq!(
7365 rete.query_iter(None, s, p, o).collect::<Vec<_>>(),
7366 eager,
7367 "query_iter disagrees for {s:?} {p:?} {o:?}"
7368 );
7369 for max in [1usize, 3, 64, 100_000] {
7370 let mut got = drain(s, p, o, max);
7371 let mut want = eager.clone();
7372 got.sort();
7375 want.sort();
7376 assert_eq!(got, want, "batch size {max} changed {s:?} {p:?} {o:?}");
7377 }
7378 if s.is_none() && p.is_none() && o.is_none() {
7380 assert_eq!(drain(s, p, o, 7), eager, "unbound batching reordered rows");
7381 }
7382 }
7383
7384 let (t, _, done) = rete.query_batch(None, Some("<http://ex/nope>"), None, None, 0, 16);
7387 assert!(t.is_empty() && done);
7388 let (t, _, done) = rete.query_batch(Some("http://ex/nope"), None, None, None, 0, 16);
7389 assert!(t.is_empty() && done);
7390 assert_eq!(
7391 rete.query_iter(None, Some("<http://ex/nope>"), None, None)
7392 .count(),
7393 0
7394 );
7395
7396 let read_bytes = |take: usize| -> u64 {
7399 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
7400 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
7401 let rete = Rete::open_ranged_lazy(reader.clone()).unwrap();
7402 let before = reader.bytes_read();
7403 let (mut n, mut cursor) = (0usize, 0u64);
7404 loop {
7405 let (batch, next, done) = rete.query_batch(None, None, None, None, cursor, 1);
7406 n += batch.len();
7407 cursor = next;
7408 if done || n >= take {
7409 break;
7410 }
7411 }
7412 assert!(n >= take.min(triples.len()));
7413 assert!(!rete.index_incomplete());
7414 reader.bytes_read() - before
7415 };
7416 let first_only = read_bytes(1);
7417 let everything = read_bytes(triples.len());
7418 assert!(
7419 first_only < everything,
7420 "query_batch is not lazy: one row read {first_only} B, a full scan {everything} B"
7421 );
7422 #[cfg(feature = "compression")]
7423 assert!(
7424 first_only * 2 < everything,
7425 "query_batch is not lazy: one row read {first_only} B of the {everything} B a full scan reads"
7426 );
7427 }
7428
7429 #[test]
7432 fn query_batch_is_graph_scoped() {
7433 let mut db = DictionaryBuilder::new();
7434 for (s, p, o) in [
7435 ("Alice", "knows", "Bob"),
7436 ("Alice", "knows", "Carol"),
7437 ("Alice", "knows", "Dave"),
7438 ("Erin", "knows", "Frank"),
7439 ("Alice", "knows", "Zoe"),
7440 ] {
7441 db.observe(s, p, o);
7442 }
7443 let dict = db.build();
7444 let mut def = GraphIndexBuilder::new();
7445 def.push(dict.encode("Alice", "knows", "Bob").unwrap());
7446 def.push(dict.encode("Alice", "knows", "Carol").unwrap());
7447 let mut g1 = GraphIndexBuilder::new();
7448 g1.push(dict.encode("Alice", "knows", "Dave").unwrap());
7449 g1.push(dict.encode("Erin", "knows", "Frank").unwrap());
7450 let mut g2 = GraphIndexBuilder::new();
7451 g2.push(dict.encode("Alice", "knows", "Zoe").unwrap());
7452 let named = vec![
7453 ("http://ex/g1".to_string(), g1.build()),
7454 ("http://ex/g2".to_string(), g2.build()),
7455 ];
7456 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
7457 let rete = Rete::open(&bytes).unwrap();
7458 for graph in [None, Some("http://ex/g1"), Some("http://ex/g2")] {
7459 let eager = rete.query_in_graph(graph, None, None, None);
7460 let (mut out, mut cursor) = (Vec::new(), 0u64);
7461 loop {
7462 let (batch, next, done) = rete.query_batch(graph, None, None, None, cursor, 1);
7463 out.extend(batch);
7464 cursor = next;
7465 if done {
7466 break;
7467 }
7468 }
7469 assert_eq!(out, eager, "graph {graph:?} batched differently");
7470 }
7471 }
7472
7473 #[test]
7474 fn query_in_graph_is_graph_scoped() {
7475 let mut db = DictionaryBuilder::new();
7478 for (s, p, o) in [
7479 ("Alice", "knows", "Bob"),
7480 ("Alice", "knows", "Carol"),
7481 ("Alice", "knows", "Dave"),
7482 ] {
7483 db.observe(s, p, o);
7484 }
7485 let dict = db.build();
7486
7487 let mut def = GraphIndexBuilder::new();
7488 def.push(dict.encode("Alice", "knows", "Bob").unwrap());
7489 def.push(dict.encode("Alice", "knows", "Carol").unwrap());
7490 let mut g1 = GraphIndexBuilder::new();
7491 g1.push(dict.encode("Alice", "knows", "Dave").unwrap());
7492
7493 let named = vec![("http://ex/g1".to_string(), g1.build())];
7494 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
7495 let rete = Rete::open(&bytes).unwrap();
7496
7497 let mut def_objs: Vec<String> = rete
7499 .query_in_graph(None, Some("Alice"), Some("knows"), None)
7500 .into_iter()
7501 .map(|(_, _, o)| o)
7502 .collect();
7503 def_objs.sort();
7504 assert_eq!(def_objs, vec!["Bob".to_string(), "Carol".to_string()]);
7505
7506 assert_eq!(
7508 rete.query_in_graph(Some("http://ex/g1"), Some("Alice"), None, None),
7509 vec![("Alice".into(), "knows".into(), "Dave".into())]
7510 );
7511
7512 assert_eq!(rete.query_in_graph(None, None, None, None).len(), 2);
7514 assert_eq!(
7515 rete.query_in_graph(Some("http://ex/g1"), None, None, None)
7516 .len(),
7517 1
7518 );
7519
7520 assert!(rete
7522 .query_in_graph(Some("http://ex/missing"), None, None, None)
7523 .is_empty());
7524 }
7525
7526 #[test]
7539 fn bounded_query_in_graph_does_not_fault_the_whole_dictionary() {
7540 use crate::reader::{CountingReader, SliceReader};
7541
7542 let payload = |i: u32| -> String {
7548 let mut s = String::with_capacity(800);
7549 let mut x = i as u64 ^ 0x9e37_79b9_7f4a_7c15;
7550 while s.len() < 800 {
7551 x = x
7552 .wrapping_mul(6364136223846793005)
7553 .wrapping_add(1442695040888963407);
7554 s.push_str(&format!("{:016x}", x ^ (x >> 29)));
7555 }
7556 s
7557 };
7558 let quads: Vec<(String, String, String, bool)> = (0..2000u32)
7565 .map(|i| {
7566 let named = i % 2 == 0;
7567 let ns = if named { 'a' } else { 'b' };
7568 (
7569 format!("<http://ex/{ns}/s/{i:04}>"),
7570 "<http://ex/abstract>".to_string(),
7571 format!("\"{ns}{} number {i}\"", payload(i)),
7572 named,
7573 )
7574 })
7575 .collect();
7576 let mut db = DictionaryBuilder::new();
7577 for (s, p, o, _) in &quads {
7578 db.observe(s, p, o);
7579 }
7580 let dict = db.build();
7581 let mut def = GraphIndexBuilder::new().with_tile_budget(256);
7582 let mut g1 = GraphIndexBuilder::new().with_tile_budget(256);
7583 for (s, p, o, named) in &quads {
7584 let t = dict.encode(s, p, o).unwrap();
7585 if *named {
7586 g1.push(t);
7587 } else {
7588 def.push(t);
7589 }
7590 }
7591 let named = vec![("http://ex/g1".to_string(), g1.build())];
7592 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
7593
7594 let eager = Rete::open(&bytes).unwrap();
7595 let dict_len = eager.header().dictionary_len;
7596
7597 let leaked: &'static [u8] = Box::leak(bytes.clone().into_boxed_slice());
7600 let probe = |f: &dyn Fn(&Rete) -> Vec<TermTriple>| -> (Vec<TermTriple>, u64) {
7601 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
7602 let lazy = Rete::open_ranged_lazy(reader.clone()).unwrap();
7603 let before = reader.bytes_read();
7604 let out = f(&lazy);
7605 assert!(!lazy.index_incomplete(), "a chunk fetch failed");
7606 (out, reader.bytes_read() - before)
7607 };
7608
7609 let want = eager.query_in_graph(
7611 Some("http://ex/g1"),
7612 Some("<http://ex/a/s/0100>"),
7613 None,
7614 None,
7615 );
7616 assert_eq!(want.len(), 1, "fixture: expected exactly one match");
7617 let (got, pulled) = probe(&|r| {
7618 r.query_in_graph(
7619 Some("http://ex/g1"),
7620 Some("<http://ex/a/s/0100>"),
7621 None,
7622 None,
7623 )
7624 });
7625 assert_eq!(got, want, "lazy query_in_graph disagrees with eager");
7626 assert!(
7627 pulled * 4 < dict_len,
7628 "a 1-row query_in_graph faulted {pulled} B of a {dict_len} B dictionary — \
7629 it is prefetching the whole thing again"
7630 );
7631
7632 let (empty, q_pulled) =
7640 probe(&|r| r.query_in_graph(Some("http://ex/nope"), None, None, None));
7641 assert!(empty.is_empty());
7642 let (_, d_pulled) = probe(&|r| {
7643 let mut n = 0usize;
7644 r.dump_each(Some("http://ex/nope"), |_, _, _| n += 1);
7645 assert_eq!(n, 0, "an absent graph yielded triples");
7646 Vec::new()
7647 });
7648 assert_eq!(
7649 q_pulled, d_pulled,
7650 "query_in_graph and dump_each disagree on what an absent graph costs"
7651 );
7652 assert!(
7653 q_pulled * 4 < dict_len,
7654 "an absent graph faulted {q_pulled} B of a {dict_len} B dictionary"
7655 );
7656
7657 let want: Vec<TermTriple> = eager.dump(Some("http://ex/g1"));
7660 assert_eq!(want.len(), 1000);
7661 let (got, pulled) = probe(&|r| {
7662 let mut out = Vec::new();
7663 r.dump_each(Some("http://ex/g1"), |s, p, o| {
7664 out.push((s.to_string(), p.to_string(), o.to_string()))
7665 });
7666 out
7667 });
7668 assert_eq!(got, want, "lazy dump_each disagrees with eager dump");
7669 assert!(
7670 pulled * 3 < dict_len * 2,
7671 "dumping half the graph faulted {pulled} B of a {dict_len} B dictionary — \
7672 the other graph's chunks are being pulled in too"
7673 );
7674 }
7675
7676 #[cfg(test)]
7681 fn build_multichunk_image() -> Vec<u8> {
7682 use crate::index::PermSet;
7683 let mut db = DictionaryBuilder::new();
7684 let mut triples: Vec<(String, String, String)> = (0..4000u32)
7685 .map(|i| {
7686 (
7687 format!("<http://ex/s/{i:07}>"),
7688 format!("<http://ex/p/{}>", i % 8),
7689 format!("\"object literal number {i:07} with padding text to widen the term\""),
7690 )
7691 })
7692 .collect();
7693 triples.push((
7694 "<http://ex/s/big>".to_string(),
7695 "<http://ex/p/big>".to_string(),
7696 format!("\"{}\"", "Z".repeat(8000)),
7697 ));
7698 for (s, p, o) in &triples {
7699 db.observe(s, p, o);
7700 }
7701 let dict = db.build();
7702 let mut idx = GraphIndexBuilder::new()
7703 .with_tile_budget(64)
7704 .with_perms(PermSet::ALL);
7705 for (s, p, o) in &triples {
7706 idx.push(dict.encode(s, p, o).unwrap());
7707 }
7708 write_dataset(&dict, &idx.build(), &[], true, &[], 0)
7709 }
7710
7711 #[test]
7720 fn tiny_dict_cache_evicts_yet_dump_is_byte_identical() {
7721 use crate::reader::SliceReader;
7722 let leaked: &'static [u8] = Box::leak(build_multichunk_image().into_boxed_slice());
7723 let dump = |cap: u64| -> (
7724 Vec<(String, String, String)>,
7725 crate::chunk_cache::CacheStats,
7726 ) {
7727 let rete = Rete::open_ranged_lazy(SliceReader::new(leaked)).unwrap();
7728 rete.dict.set_cache_cap(cap);
7729 let mut out = Vec::new();
7730 rete.dump_filtered_each(None, None, None, None, |s, p, o| {
7731 out.push((s.to_string(), p.to_string(), o.to_string()))
7732 });
7733 assert!(
7734 !rete.index_incomplete(),
7735 "no lazy fault should fail at cap {cap}"
7736 );
7737 let stats = rete.dict_cache_stats();
7738 (out, stats)
7739 };
7740
7741 let (unlimited, us) = dump(u64::MAX);
7742 assert_eq!(unlimited.len(), 4001, "every triple dumped");
7743 assert_eq!(us.evictions, 0, "unlimited cap evicts nothing");
7744
7745 let (tiny, ts) = dump(128 * 1024);
7748 assert_eq!(tiny, unlimited, "128 KiB-cap dump differs from unlimited");
7749 assert!(ts.evictions > 0, "a tiny cap must evict (stats {ts:?})");
7750
7751 let (starved, ss) = dump(1);
7753 assert_eq!(starved, unlimited, "1-byte-cap dump differs from unlimited");
7754 assert!(
7755 ss.evictions > 0,
7756 "1-byte cap evicts every insert (stats {ss:?})"
7757 );
7758 }
7759
7760 #[test]
7765 fn releasing_a_named_graph_drops_and_reopens_its_index() {
7766 use crate::index::PermSet;
7767 use crate::reader::SliceReader;
7768 let mut db = DictionaryBuilder::new();
7769 let all: Vec<(String, String, String)> = (0..40u32)
7770 .map(|i| {
7771 (
7772 format!("<http://ex/s{}>", i % 7),
7773 format!("<http://ex/p{}>", i % 3),
7774 format!("<http://ex/o{}>", i % 5),
7775 )
7776 })
7777 .collect();
7778 for (s, p, o) in &all {
7779 db.observe(s, p, o);
7780 }
7781 let dict = db.build();
7782 let mk = |take: fn(u32) -> bool| {
7783 let mut g = GraphIndexBuilder::new()
7784 .with_tile_budget(16)
7785 .with_perms(PermSet::ALL);
7786 for (i, (s, p, o)) in all.iter().enumerate() {
7787 if take(i as u32) {
7788 g.push(dict.encode(s, p, o).unwrap());
7789 }
7790 }
7791 g.build()
7792 };
7793 let def = mk(|_| true);
7794 let g1 = mk(|i| i % 2 == 0);
7795 let g2 = mk(|i| i % 3 == 0);
7796 let image = write_dataset(
7797 &dict,
7798 &def,
7799 &[
7800 ("<http://ex/g1>".to_string(), g1),
7801 ("<http://ex/g2>".to_string(), g2),
7802 ],
7803 true,
7804 &[],
7805 0,
7806 );
7807 let leaked: &'static [u8] = Box::leak(image.into_boxed_slice());
7808 let mut rete = Rete::open_ranged_lazy(SliceReader::new(leaked)).unwrap();
7809
7810 let opened = |r: &Rete| -> usize {
7811 match &r.named_graphs {
7812 NamedGraphsSlot::Lazy(l) => {
7813 let mut n = 0;
7814 l.for_each_opened(|_| n += 1);
7815 n
7816 }
7817 _ => 0,
7818 }
7819 };
7820
7821 assert_eq!(opened(&rete), 0, "nothing opened before first touch");
7822 assert!(rete.graph_index("<http://ex/g1>").is_some());
7824 assert_eq!(opened(&rete), 1, "g1 opened");
7825 rete.release_named_graph("<http://ex/g1>");
7827 assert_eq!(opened(&rete), 0, "g1 released");
7828 rete.release_named_graph("<http://ex/absent>");
7830 rete.release_named_graph("<http://ex/g1>");
7831 assert_eq!(opened(&rete), 0);
7832 assert!(rete.graph_index("<http://ex/g1>").is_some());
7834 assert_eq!(opened(&rete), 1, "g1 re-opened after release");
7835 }
7836
7837 #[test]
7855 fn filtered_dump_is_exactly_the_unfiltered_dump_filtered() {
7856 use crate::index::PermSet;
7857 use crate::reader::SliceReader;
7858
7859 let terms = |i: u32| {
7860 (
7861 format!("<http://ex/s{}>", i % 11),
7862 format!("<http://ex/p{}>", i % 3),
7863 if i.is_multiple_of(4) {
7867 format!("\"lit {}\"@en", i % 5)
7868 } else {
7869 format!("<http://ex/o{}>", i % 7)
7870 },
7871 )
7872 };
7873 let build = |perms: PermSet| {
7874 let mut db = DictionaryBuilder::new();
7875 let all: Vec<_> = (0..90u32).map(terms).collect();
7876 for (s, p, o) in &all {
7877 db.observe(s, p, o);
7878 }
7879 let dict = db.build();
7880 let mut def = GraphIndexBuilder::new()
7881 .with_tile_budget(16)
7882 .with_perms(perms);
7883 let mut g1 = GraphIndexBuilder::new()
7884 .with_tile_budget(16)
7885 .with_perms(perms);
7886 for (i, (s, p, o)) in all.iter().enumerate() {
7887 let t = dict.encode(s, p, o).unwrap();
7888 if i % 2 == 0 {
7892 g1.push(t);
7893 }
7894 def.push(t);
7895 }
7896 write_dataset(
7897 &dict,
7898 &def.build(),
7899 &[("<http://ex/g1>".to_string(), g1.build())],
7900 true,
7901 &[],
7902 0,
7903 )
7904 };
7905
7906 let shapes: Vec<(Option<&str>, Option<&str>, Option<&str>)> = vec![
7907 (None, None, None),
7908 (Some("<http://ex/s3>"), None, None),
7909 (None, Some("<http://ex/p1>"), None),
7910 (None, None, Some("<http://ex/o5>")),
7911 (None, None, Some("\"lit 0\"@en")),
7912 (Some("<http://ex/s3>"), Some("<http://ex/p0>"), None),
7913 (Some("<http://ex/s4>"), None, Some("<http://ex/o4>")),
7914 (None, Some("<http://ex/p1>"), Some("<http://ex/o5>")),
7915 (
7916 Some("<http://ex/s3>"),
7917 Some("<http://ex/p0>"),
7918 Some("<http://ex/o3>"),
7919 ),
7920 (Some("<http://ex/nope>"), None, None),
7924 (None, Some("<http://ex/nope>"), None),
7925 (None, None, Some("<http://ex/nope>")),
7926 ];
7927
7928 let mut any_pruned = false;
7929 for perms in [PermSet::ALL, PermSet::CORE] {
7930 let image = build(perms);
7931 let leaked: &'static [u8] = Box::leak(image.clone().into_boxed_slice());
7932 let opens: Vec<(&str, Rete)> = vec![
7933 ("resident", Rete::open(&image).unwrap()),
7934 (
7935 "ranged",
7936 Rete::open_ranged(&SliceReader::new(leaked)).unwrap(),
7937 ),
7938 (
7939 "lazy",
7940 Rete::open_ranged_lazy(SliceReader::new(leaked)).unwrap(),
7941 ),
7942 ];
7943 for graph in [None, Some("<http://ex/g1>"), Some("<http://ex/absent>")] {
7944 let full: Vec<TermTriple> = {
7946 let mut v = Vec::new();
7947 Rete::open(&image).unwrap().dump_each(graph, |s, p, o| {
7948 v.push((s.to_string(), p.to_string(), o.to_string()))
7949 });
7950 v.sort();
7951 v
7952 };
7953 for &(s, p, o) in &shapes {
7954 let mut want = full.clone();
7955 want.retain(|(ts, tp, to)| {
7956 s.is_none_or(|x| x == ts)
7957 && p.is_none_or(|x| x == tp)
7958 && o.is_none_or(|x| x == to)
7959 });
7960 for (path, rete) in &opens {
7961 let mut got = Vec::new();
7962 rete.dump_filtered_each(graph, s, p, o, |s, p, o| {
7963 got.push((s.to_string(), p.to_string(), o.to_string()))
7964 });
7965 got.sort();
7966 assert_eq!(
7967 got,
7968 want,
7969 "{path}/{perms:?} graph={graph:?} pattern={:?} — filtered dump \
7970 is not the unfiltered dump filtered",
7971 (s, p, o)
7972 );
7973 assert!(!rete.index_incomplete(), "{path}: a fetch failed");
7974 }
7975 let plan = opens[2].1.dump_plan(graph, s, p, o);
7979 if graph == Some("<http://ex/absent>")
7980 || [s, p, o].contains(&Some("<http://ex/nope>"))
7981 {
7982 assert!(
7983 plan.scan.is_none(),
7984 "an unmatchable dump still planned a scan: {:?}",
7985 (graph, s, p, o)
7986 );
7987 } else {
7988 let scan = plan.scan.expect("a matchable dump plans a scan");
7989 assert!(
7990 PermSet::CORE.contains(scan.permutation),
7991 "a dump routed to {:?}, outside PermSet::CORE — a \
7992 --permutations 3 file has no such section",
7993 scan.permutation
7994 );
7995 assert!(scan.tiles_admitted <= scan.tiles_routed);
7996 assert!(scan.tiles_routed <= scan.tiles_total);
7997 assert!(scan.tile_bytes <= scan.section_bytes);
7998 if scan.tiles_admitted < scan.tiles_total {
7999 any_pruned = true;
8000 }
8001 }
8002 }
8003 }
8004 }
8005 assert!(
8006 any_pruned,
8007 "no shape pruned a single tile — the fixture is too small to be \
8008 measuring anything"
8009 );
8010 }
8011
8012 #[test]
8022 fn a_predicate_scoped_dump_fetches_less_than_the_graph() {
8023 use crate::reader::{CountingReader, SliceReader};
8024
8025 let mut db = DictionaryBuilder::new();
8031 let mut all = Vec::new();
8032 for s in 0..200u32 {
8033 for p in 0..20u32 {
8034 let t = (
8035 format!("<http://ex/s/{s:04}>"),
8036 format!("<http://ex/p/{p:02}>"),
8037 format!("\"object {s:04}/{p:02}\""),
8038 );
8039 db.observe(&t.0, &t.1, &t.2);
8040 all.push(t);
8041 }
8042 }
8043 let dict = db.build();
8044 let mut ib = GraphIndexBuilder::new().with_tile_budget(64);
8045 for (s, p, o) in &all {
8046 ib.push(dict.encode(s, p, o).unwrap());
8047 }
8048 let image = write_dataset(&dict, &ib.build(), &[], false, &[], 0);
8049 let leaked: &'static [u8] = Box::leak(image.clone().into_boxed_slice());
8050
8051 let probe = |s: Option<&str>, p: Option<&str>, o: Option<&str>| -> (usize, u64) {
8053 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
8054 let lazy = Rete::open_ranged_lazy(reader.clone()).unwrap();
8055 let before = reader.bytes_read();
8056 let mut n = 0usize;
8057 lazy.dump_filtered_each(None, s, p, o, |_, _, _| n += 1);
8058 assert!(!lazy.index_incomplete());
8059 (n, reader.bytes_read() - before)
8060 };
8061
8062 let (full_rows, full_bytes) = probe(None, None, None);
8063 let (pred_rows, pred_bytes) = probe(None, Some("<http://ex/p/07>"), None);
8064 let (subj_rows, subj_bytes) = probe(Some("<http://ex/s/0021>"), None, None);
8065 assert_eq!(full_rows, 4000);
8066 assert_eq!(pred_rows, 200, "one predicate covers every subject once");
8067 assert_eq!(subj_rows, 20);
8068 assert!(
8070 pred_bytes < full_bytes && subj_bytes < full_bytes,
8071 "a slice fetched {pred_bytes} / {subj_bytes} B where the whole graph fetched \
8072 {full_bytes} B — the scan is not being pruned at all"
8073 );
8074 #[cfg(feature = "compression")]
8085 assert!(
8086 pred_bytes * 4 < full_bytes && subj_bytes * 4 < full_bytes,
8087 "a 1-in-20 predicate slice fetched {pred_bytes} B and one subject \
8088 {subj_bytes} B where the whole graph fetched {full_bytes} B"
8089 );
8090
8091 let reader = std::sync::Arc::new(CountingReader::new(SliceReader::new(leaked)));
8093 let lazy = Rete::open_ranged_lazy(reader.clone()).unwrap();
8094 let before = reader.bytes_read();
8095 let plan = lazy
8096 .dump_plan(None, None, Some("<http://ex/p/07>"), None)
8097 .scan
8098 .expect("a known predicate plans a scan");
8099 assert!(
8100 reader.bytes_read() - before < 64 * 1024,
8101 "planning fetched {} B — it is supposed to read directories, not tiles",
8102 reader.bytes_read() - before
8103 );
8104 assert!(plan.tiles_admitted < plan.tiles_total);
8105 assert!(
8111 plan.tile_bytes * 5 < plan.section_bytes,
8112 "one predicate of twenty planned {} B of a {} B section",
8113 plan.tile_bytes,
8114 plan.section_bytes
8115 );
8116 }
8117
8118 #[test]
8119 fn query_quads_tags_every_graph() {
8120 let mut db = DictionaryBuilder::new();
8121 for (s, p, o) in [("Alice", "knows", "Bob"), ("Alice", "knows", "Dave")] {
8122 db.observe(s, p, o);
8123 }
8124 let dict = db.build();
8125 let mut def = GraphIndexBuilder::new();
8126 def.push(dict.encode("Alice", "knows", "Bob").unwrap());
8127 let mut g1 = GraphIndexBuilder::new();
8128 g1.push(dict.encode("Alice", "knows", "Dave").unwrap());
8129 let named = vec![("http://ex/g1".to_string(), g1.build())];
8130 let bytes = write_dataset(&dict, &def.build(), &named, true, &[], 0);
8131 let rete = Rete::open(&bytes).unwrap();
8132
8133 let quads = rete.query_quads(Some("Alice"), Some("knows"), None);
8135 assert_eq!(quads.len(), 2);
8136 assert_eq!(
8137 quads[0],
8138 (("Alice".into(), "knows".into(), "Bob".into()), None)
8139 );
8140 assert_eq!(
8141 quads[1],
8142 (
8143 ("Alice".into(), "knows".into(), "Dave".into()),
8144 Some("http://ex/g1".to_string())
8145 )
8146 );
8147
8148 assert!(rete.query_quads(Some("Nobody"), None, None).is_empty());
8150 }
8151
8152 #[test]
8153 fn pyramid_meta_round_trips_in_file() {
8154 let rete = Rete::open(&build_image()).unwrap();
8155 let pyr = rete.pyramid().expect("file has a pyramid");
8156 let total: u32 = pyr.summary.iter().map(|e| e.count).sum();
8158 assert_eq!(total, 3);
8159 assert!(!pyr.summary.is_empty());
8160 assert!(pyr.tiles.is_empty());
8161 }
8162
8163 #[test]
8164 fn schema_pyramid_round_trips_through_file_index_free() {
8165 use crate::reader::{CountingReader, SliceReader};
8166 let sub = "<http://www.w3.org/2000/01/rdf-schema#subClassOf>";
8167 let q = |s: &str, p: &str, o: &str| {
8168 (s.to_string(), p.to_string(), o.to_string(), None::<String>)
8169 };
8170 let quads = vec![
8172 q("<a>", RDF_TYPE, "<Astronomer>"),
8173 q("<b>", RDF_TYPE, "<Astronomer>"),
8174 q("<c>", RDF_TYPE, "<Person>"),
8175 q("<Astronomer>", sub, "<Scientist>"),
8176 q("<Scientist>", sub, "<Person>"),
8177 q("<Person>", sub, "<Agent>"),
8178 q("<a>", "<knows>", "<b>"),
8179 q("<b>", "<knows>", "<c>"),
8180 ];
8181 let (bytes, _) =
8182 crate::ingest::assemble_dataset_with_opts(quads, true, false, None, |_, _| Vec::new());
8183
8184 let rete = Rete::open(&bytes).unwrap();
8186 let pyr = rete.pyramid().expect("pyramid present");
8187 assert!(!pyr.level_rollups.is_empty(), "schema pyramid shipped");
8188 assert!(pyr
8189 .class_hierarchy
8190 .iter()
8191 .any(|n| n.class == "<Agent>" && n.depth == 0));
8192
8193 let r = CountingReader::new(SliceReader::new(&bytes));
8195 let view = SummaryView::open_ranged(&r).unwrap().unwrap();
8196 assert!(view.level_count() >= 2, "multi-level pyramid");
8197 let coarse = view.level_rollup(0).unwrap();
8198 assert!(
8199 coarse.classes.iter().any(|(c, _)| c == "<Agent>"),
8200 "coarsest level rolls up to the root Agent"
8201 );
8202 let h = Header::from_bytes(&bytes).unwrap();
8203 assert!(
8204 r.bytes_read() <= bytes.len() as u64 - h.root_dir_len,
8205 "summary read {} bytes; the {}-byte index section must be skipped",
8206 r.bytes_read(),
8207 h.root_dir_len
8208 );
8209 }
8210
8211 #[test]
8212 fn predicate_totals_from_summary_only() {
8213 use crate::reader::SliceReader;
8214 let bytes = build_image();
8216 let reader = SliceReader::new(&bytes);
8217 let view = SummaryView::open_ranged(&reader).unwrap().unwrap();
8218 assert_eq!(view.predicate_total("knows"), 2);
8219 assert_eq!(view.predicate_total("age"), 1);
8220 assert_eq!(view.predicate_total("missing"), 0);
8221 let totals = view.predicate_totals();
8222 assert_eq!(totals[0], ("knows".to_string(), 2)); }
8224
8225 #[test]
8226 fn query_patterns_resolve_to_terms() {
8227 let rete = Rete::open(&build_image()).unwrap();
8228
8229 assert_eq!(rete.query(None, None, None).len(), 3);
8231
8232 let mut alice = rete.query(Some("Alice"), None, None);
8234 alice.sort();
8235 assert_eq!(
8236 alice,
8237 vec![
8238 ("Alice".into(), "age".into(), "30".into()),
8239 ("Alice".into(), "knows".into(), "Bob".into()),
8240 ]
8241 );
8242
8243 assert_eq!(rete.query(None, Some("knows"), None).len(), 2);
8245
8246 assert_eq!(
8248 rete.query(Some("Bob"), Some("knows"), Some("Carol")),
8249 vec![("Bob".into(), "knows".into(), "Carol".into())]
8250 );
8251 assert!(rete.query(Some("Nobody"), None, None).is_empty());
8252 assert!(rete.query(None, Some("likes"), None).is_empty());
8253 }
8254
8255 #[test]
8256 fn query_provenance_reports_terms_ids_sections_and_index_choice() {
8257 let bytes = build_image();
8258 let rete = Rete::open(&bytes).unwrap();
8259
8260 let mut matches = rete.query_with_provenance(None, Some("knows"), None);
8261 matches.sort_by(|a, b| a.terms.cmp(&b.terms));
8262
8263 assert_eq!(matches.len(), 2);
8264 assert_eq!(
8265 matches[0].terms,
8266 ("Alice".into(), "knows".into(), "Bob".into())
8267 );
8268 assert_eq!(
8269 matches[0].ids,
8270 rete.dictionary().encode("Alice", "knows", "Bob").unwrap()
8271 );
8272 assert_eq!(matches[0].graph.as_deref(), None);
8273 assert_eq!(
8274 matches[0].matched_pattern,
8275 (None, Some(matches[0].ids.1), None)
8276 );
8277 assert_eq!(
8278 matches[0].index_permutation,
8279 crate::index::IndexPermutation::Pos
8280 );
8281
8282 let h = rete.header();
8283 assert_eq!(matches[0].dictionary_range.offset, h.dictionary_offset);
8284 assert_eq!(matches[0].dictionary_range.len, h.dictionary_len);
8285 assert_eq!(matches[0].index_range.offset, h.root_dir_offset);
8286 assert_eq!(matches[0].index_range.len, h.root_dir_len);
8287 assert!(
8288 matches[0].index_section_range.offset > h.root_dir_offset,
8289 "POS is section 1, so its payload starts after the container header and SPO payload"
8290 );
8291 assert!(matches[0].index_section_range.len > 0);
8292 assert!(matches[0].index_section_range.end() <= matches[0].index_range.end());
8293 assert!(matches[0].index_section_range.len < matches[0].index_range.len);
8294 assert_eq!(
8295 matches[0].pyramid_range.as_ref().map(|r| (r.offset, r.len)),
8296 Some((h.pyramid_meta_offset, h.pyramid_meta_len))
8297 );
8298 let tile_range = matches[0].tile_range.expect("tiled file reports a tile");
8301 assert!(matches[0]
8302 .tile
8303 .as_deref()
8304 .unwrap()
8305 .starts_with(matches[0].index_permutation.name()));
8306 assert!(matches[0].index_section_range.offset <= tile_range.offset);
8307 assert!(tile_range.end() <= matches[0].index_section_range.end());
8308 }
8309
8310 fn build_labeled(n: usize) -> Vec<u8> {
8314 const LABEL: &str = "<http://www.w3.org/2000/01/rdf-schema#label>";
8315 const WORDS: &[&str] = &[
8316 "alanine",
8317 "benzene",
8318 "glucose",
8319 "dextrose",
8320 "ethanol",
8321 "formate",
8322 "heptane",
8323 "isoleucine",
8324 ];
8325 let triples: Vec<(String, String, String)> = (0..n)
8326 .flat_map(|i| {
8327 let s = format!("<http://ex/e{i}>");
8328 let w = WORDS[i % WORDS.len()];
8329 [
8330 (s.clone(), LABEL.to_string(), format!("\"{w}-{i:06}\"")),
8331 (
8332 s,
8333 "<http://ex/p>".to_string(),
8334 format!("<http://ex/c{}>", i % 64),
8335 ),
8336 ]
8337 })
8338 .collect();
8339 let mut db = DictionaryBuilder::new();
8340 for (s, p, o) in &triples {
8341 db.observe(s, p, o);
8342 }
8343 let dict = db.build();
8344 let ids: Vec<(u32, u32, u32)> = triples
8345 .iter()
8346 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
8347 .collect();
8348 let mut ib = GraphIndexBuilder::new();
8349 for &t in &ids {
8350 ib.push(t);
8351 }
8352 let (meta, levels) = build_pyramid_meta(&dict, &ids, DEFAULT_TILE_BUDGET);
8353 write_file(&dict, &ib.build(), false, &meta, levels)
8354 }
8355
8356 #[test]
8357 fn prefix_search_matches_a_filter_scan() {
8358 let bytes = build_labeled(800);
8361 let rete = Rete::open(&bytes).unwrap();
8362 let idx_subjects: std::collections::BTreeSet<String> = rete
8363 .prefix_search("glucose", 10_000)
8364 .into_iter()
8365 .map(|(_label, subject)| subject)
8366 .collect();
8367 let q = "SELECT ?s WHERE { ?s <http://www.w3.org/2000/01/rdf-schema#label> ?l \
8369 FILTER(STRSTARTS(LCASE(?l), \"glucose\")) }";
8370 let crate::QueryOutput::Select(_, rows) = crate::eval_query(&rete, q).unwrap() else {
8371 panic!("expected SELECT");
8372 };
8373 let scan_subjects: std::collections::BTreeSet<String> =
8374 rows.iter().map(|r| r.get("s").cloned().unwrap()).collect();
8375 assert_eq!(idx_subjects, scan_subjects, "index agrees with the scan");
8376 assert_eq!(
8377 idx_subjects.len(),
8378 100,
8379 "800/8 words = 100 glucose-* labels"
8380 );
8381 }
8382
8383 #[test]
8387 #[ignore]
8388 fn bench_prefix_search_vs_filter_scan() {
8389 use std::time::Instant;
8390 let n = 6000; let bytes = build_labeled(n);
8392 let rete = Rete::open(&bytes).unwrap();
8393 let q = "SELECT ?s WHERE { ?s <http://www.w3.org/2000/01/rdf-schema#label> ?l \
8394 FILTER(STRSTARTS(LCASE(?l), \"glucose\")) }";
8395 let reps = 200;
8396 let idx_n = rete.prefix_search("glucose", 100_000).len();
8397 let t = Instant::now();
8398 for _ in 0..reps {
8399 std::hint::black_box(rete.prefix_search("glucose", 100_000));
8400 }
8401 let idx_ms = t.elapsed().as_secs_f64() * 1000.0 / reps as f64;
8402 let t = Instant::now();
8403 for _ in 0..reps {
8404 let _ = std::hint::black_box(crate::eval_query(&rete, q).unwrap());
8405 }
8406 let scan_ms = t.elapsed().as_secs_f64() * 1000.0 / reps as f64;
8407 println!(
8408 "label prefix search over {n} labeled subjects ({idx_n} matches): \
8409 index {idx_ms:.4} ms vs FILTER scan {scan_ms:.3} ms ({:.0}× faster)",
8410 scan_ms / idx_ms
8411 );
8412 }
8413
8414 #[test]
8418 #[ignore]
8419 fn bench_text_search_vs_contains_scan() {
8420 use std::time::Instant;
8421 const LABEL: &str = "<http://www.w3.org/2000/01/rdf-schema#label>";
8422 const WORDS: &[&str] = &[
8423 "alanine",
8424 "benzene",
8425 "glucose",
8426 "dextrose",
8427 "ethanol",
8428 "formate",
8429 "heptane",
8430 "isoleucine",
8431 ];
8432 let n = 6000;
8433 let triples: Vec<(String, String, String)> = (0..n)
8434 .map(|i| {
8435 (
8436 format!("<http://ex/e{i}>"),
8437 LABEL.to_string(),
8438 format!("\"{} sample number {i:06}\"", WORDS[i % WORDS.len()]),
8439 )
8440 })
8441 .collect();
8442 let mut db = DictionaryBuilder::new();
8443 for (s, p, o) in &triples {
8444 db.observe(s, p, o);
8445 }
8446 let dict = db.build();
8447 let ids: Vec<(u32, u32, u32)> = triples
8448 .iter()
8449 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
8450 .collect();
8451 let mut ib = GraphIndexBuilder::new();
8452 for &t in &ids {
8453 ib.push(t);
8454 }
8455 let ti = compute_text_index(&dict, &ids);
8456 let bytes = write_dataset_with_metadata(&dict, &ib.build(), &[], false, &[], 0, &[], &ti);
8457 let rete = Rete::open(&bytes).unwrap();
8458
8459 let q = "SELECT ?s WHERE { ?s <http://www.w3.org/2000/01/rdf-schema#label> ?l \
8460 FILTER(CONTAINS(LCASE(?l), \"glucose\")) }";
8461 let reps = 200;
8462 let idx_n = rete.text_search(&["glucose"], None, 100_000).len();
8463 let t = Instant::now();
8464 for _ in 0..reps {
8465 std::hint::black_box(rete.text_search(&["glucose"], None, 100_000));
8466 }
8467 let idx_ms = t.elapsed().as_secs_f64() * 1000.0 / reps as f64;
8468 let t = Instant::now();
8469 for _ in 0..reps {
8470 let _ = std::hint::black_box(crate::eval_query(&rete, q).unwrap());
8471 }
8472 let scan_ms = t.elapsed().as_secs_f64() * 1000.0 / reps as f64;
8473 println!(
8474 "text search over {n} literals ({idx_n} matches): \
8475 index {idx_ms:.4} ms vs FILTER(CONTAINS) scan {scan_ms:.3} ms ({:.0}× faster)",
8476 scan_ms / idx_ms
8477 );
8478 }
8479
8480 #[test]
8484 #[ignore = "operational tool, driven by RETE_DEBUG_* env vars"]
8485 fn debug_bound_po_routing() {
8486 struct FR(std::fs::File);
8487 impl crate::RangeReader for FR {
8488 fn len(&self) -> u64 {
8489 self.0.metadata().map(|m| m.len()).unwrap_or(0)
8490 }
8491 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
8492 use std::os::unix::fs::FileExt;
8493 let mut buf = vec![0u8; len as usize];
8494 self.0.read_exact_at(&mut buf, offset)?;
8495 Ok(buf)
8496 }
8497 }
8498 let path = std::env::var("RETE_DEBUG_FILE").expect("RETE_DEBUG_FILE");
8499 let p_iri = std::env::var("RETE_DEBUG_P").expect("RETE_DEBUG_P");
8500 let o_iri = std::env::var("RETE_DEBUG_O").expect("RETE_DEBUG_O");
8501 let rete =
8502 Rete::open_ranged_lazy(std::sync::Arc::new(FR(std::fs::File::open(&path).unwrap())))
8503 .unwrap();
8504 let pid = rete.dict.predicate_id(&p_iri).expect("p resolves");
8505 let oid = rete.dict.object_id(&o_iri).expect("o resolves");
8506 eprintln!("pid={pid} oid={oid}");
8507 let pattern = (None, Some(pid), Some(oid));
8508 let perm = GraphIndex::best_permutation(pattern);
8509 eprintln!("best_permutation = {}", perm.name());
8510 let si = perm.section_index();
8511 let tiles = &rete.index.sections[si];
8512 eprintln!("section {} tiles = {}", perm.name(), tiles.len());
8513 let [pa, pb, pc] = perm.order_pattern(pattern);
8514 eprintln!("permuted pattern pa={pa:?} pb={pb:?} pc={pc:?}");
8515 let (start, end) = rete.index.tile_span(si, pa);
8516 eprintln!("tile_span = [{start}, {end}) -> {} tiles", end - start);
8517 let mut admitted = 0usize;
8518 for (ti, t) in tiles.iter().enumerate().take(end).skip(start) {
8519 if t.syn_admits(pb, pc) {
8520 admitted += 1;
8521 if admitted <= 10 {
8522 let (lo, hi) = t.leading_range();
8523 eprintln!(" admit tile {ti}: a=[{lo},{hi}] syn={:?}", t.syn);
8524 }
8525 }
8526 }
8527 eprintln!("admitted {admitted} tile(s) by synopsis");
8528 let n = rete.index.scan_iter(pattern).count();
8529 eprintln!("scan_iter matches = {n}");
8530 let hi_res = rete.query(None, Some(&p_iri), Some(&o_iri));
8531 eprintln!("high-level query matches = {}", hi_res.len());
8532 }
8533
8534 #[test]
8542 fn layout_header_label_matches_its_span() {
8543 let mut db = DictionaryBuilder::new();
8544 let mut triples = Vec::new();
8545 for i in 0..24u32 {
8546 let (s, p, o) = (
8547 format!("<http://ex/s{}>", i % 5),
8548 format!("<http://ex/p{}>", i % 3),
8549 format!("<http://ex/o{i}>"),
8550 );
8551 db.observe(&s, &p, &o);
8552 triples.push((s, p, o));
8553 }
8554 let dict = db.build();
8555 let ids: Vec<(u32, u32, u32)> = triples
8556 .iter()
8557 .map(|(s, p, o)| dict.encode(s, p, o).unwrap())
8558 .collect();
8559 let index = GraphIndexBuilder::from_triples(ids).build();
8560 let image = write_dataset_with_metadata(
8561 &dict,
8562 &index,
8563 &[],
8564 false,
8565 &[],
8566 0,
8567 br#"{"name":"layout-label-fixture"}"#,
8568 &[],
8569 );
8570 let rete = Rete::open(&image).unwrap();
8571
8572 let layout = rete.file_layout();
8573 let header = layout
8574 .iter()
8575 .find(|s| s.kind == "header")
8576 .expect("every layout starts with a header segment");
8577
8578 assert_eq!(
8579 header.offset, 0,
8580 "the header is the first thing in the file"
8581 );
8582 assert_eq!(
8583 header.len, HEADER_LEN as u64,
8584 "the header segment must span exactly HEADER_LEN"
8585 );
8586
8587 let digits: String = header
8590 .label
8591 .chars()
8592 .skip_while(|c| !c.is_ascii_digit())
8593 .take_while(|c| c.is_ascii_digit())
8594 .collect();
8595 assert!(
8596 !digits.is_empty(),
8597 "header label {:?} states no size at all",
8598 header.label
8599 );
8600 assert_eq!(
8601 digits.parse::<u64>().unwrap(),
8602 header.len,
8603 "header label {:?} disagrees with the {}-byte span it describes",
8604 header.label,
8605 header.len
8606 );
8607
8608 assert!(
8611 layout.windows(2).all(|w| w[0].offset <= w[1].offset),
8612 "file_layout must return segments sorted by offset"
8613 );
8614 }
8615}