Skip to main content

rete_core/
triples.rs

1//! Integer triple blocks: grouped, delta-coded adjacency with a zone map
2//! (SPEC.md §6.1).
3//!
4//! A block stores triples `(a, b, c)` of dictionary IDs for one permutation,
5//! sorted ascending. `a`/`b`/`c` are the permuted roles — for an SPO block
6//! `a=subject, b=predicate, c=object`; for POS, `a=predicate, b=object,
7//! c=subject`; and so on. The encoding is role-agnostic.
8
9use crate::varint::{read_uvarint, write_uvarint};
10
11/// A triple of dictionary IDs in some permutation's component order.
12pub type Triple = (u32, u32, u32);
13
14/// Per-block summary statistics enabling block-skipping.
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
16pub struct ZoneMap {
17    pub min_a: u32,
18    pub max_a: u32,
19    pub min_b: u32,
20    pub max_b: u32,
21    pub min_c: u32,
22    pub max_c: u32,
23    pub count: u32,
24}
25
26impl ZoneMap {
27    /// Could a triple with the given bound components possibly live in this
28    /// block? `None` means "unbound — don't constrain on this component".
29    pub fn may_contain(&self, a: Option<u32>, b: Option<u32>, c: Option<u32>) -> bool {
30        let in_range = |v: Option<u32>, lo: u32, hi: u32| v.is_none_or(|x| lo <= x && x <= hi);
31        in_range(a, self.min_a, self.max_a)
32            && in_range(b, self.min_b, self.max_b)
33            && in_range(c, self.min_c, self.max_c)
34    }
35}
36
37#[derive(Debug, thiserror::Error)]
38pub enum TripleError {
39    #[error("malformed triple block: {0}")]
40    Malformed(&'static str),
41}
42
43/// Accumulates triples and serializes a block.
44#[derive(Default)]
45pub struct TripleBlockBuilder {
46    triples: Vec<Triple>,
47}
48
49impl TripleBlockBuilder {
50    pub fn new() -> Self {
51        Self::default()
52    }
53
54    pub fn push(&mut self, t: Triple) {
55        self.triples.push(t);
56    }
57
58    pub fn len(&self) -> usize {
59        self.triples.len()
60    }
61
62    pub fn is_empty(&self) -> bool {
63        self.triples.is_empty()
64    }
65
66    /// Sort, dedup, and encode. Returns an empty body for zero triples.
67    pub fn build(mut self) -> Vec<u8> {
68        self.triples.sort_unstable();
69        self.triples.dedup();
70        let t = &self.triples;
71
72        let mut out = Vec::new();
73        if t.is_empty() {
74            // zone map of zeros + count 0, no body.
75            for _ in 0..7 {
76                write_uvarint(&mut out, 0);
77            }
78            write_uvarint(&mut out, 0); // num_a
79            return out;
80        }
81
82        // Zone map.
83        let (mut min_a, mut max_a) = (u32::MAX, 0u32);
84        let (mut min_b, mut max_b) = (u32::MAX, 0u32);
85        let (mut min_c, mut max_c) = (u32::MAX, 0u32);
86        for &(a, b, c) in t {
87            min_a = min_a.min(a);
88            max_a = max_a.max(a);
89            min_b = min_b.min(b);
90            max_b = max_b.max(b);
91            min_c = min_c.min(c);
92            max_c = max_c.max(c);
93        }
94        for v in [min_a, max_a, min_b, max_b, min_c, max_c, t.len() as u32] {
95            write_uvarint(&mut out, v as u64);
96        }
97
98        // Body: grouped delta adjacency.
99        // First pass: collect a-groups -> b-groups -> c-list.
100        // (b, c-list) for one b under an a; and (a, its b-groups).
101        type BGroup = (u32, Vec<u32>);
102        type AGroup = (u32, Vec<BGroup>);
103        let mut i = 0;
104        let mut a_groups: Vec<AGroup> = Vec::new();
105        while i < t.len() {
106            let a = t[i].0;
107            let mut b_groups: Vec<(u32, Vec<u32>)> = Vec::new();
108            while i < t.len() && t[i].0 == a {
109                let b = t[i].1;
110                let mut cs = Vec::new();
111                while i < t.len() && t[i].0 == a && t[i].1 == b {
112                    cs.push(t[i].2);
113                    i += 1;
114                }
115                b_groups.push((b, cs));
116            }
117            a_groups.push((a, b_groups));
118        }
119
120        write_uvarint(&mut out, a_groups.len() as u64);
121        let mut prev_a = 0u32;
122        for (a, b_groups) in &a_groups {
123            write_uvarint(&mut out, (a - prev_a) as u64);
124            prev_a = *a;
125            write_uvarint(&mut out, b_groups.len() as u64);
126            let mut prev_b = 0u32;
127            for (b, cs) in b_groups {
128                write_uvarint(&mut out, (b - prev_b) as u64);
129                prev_b = *b;
130                write_uvarint(&mut out, cs.len() as u64);
131                let mut prev_c = 0u32;
132                for c in cs {
133                    write_uvarint(&mut out, (c - prev_c) as u64);
134                    prev_c = *c;
135                }
136            }
137        }
138        out
139    }
140}
141
142/// A parsed triple block.
143pub struct TripleBlock<'a> {
144    bytes: &'a [u8],
145    zone: ZoneMap,
146    body_start: usize,
147}
148
149impl<'a> TripleBlock<'a> {
150    pub fn parse(bytes: &'a [u8]) -> Result<Self, TripleError> {
151        let mut pos = 0;
152        let take = |pos: &mut usize| -> Result<u32, TripleError> {
153            let (v, n) = read_uvarint(&bytes[*pos..]).ok_or(TripleError::Malformed("truncated"))?;
154            *pos += n;
155            Ok(v as u32)
156        };
157        let zone = ZoneMap {
158            min_a: take(&mut pos)?,
159            max_a: take(&mut pos)?,
160            min_b: take(&mut pos)?,
161            max_b: take(&mut pos)?,
162            min_c: take(&mut pos)?,
163            max_c: take(&mut pos)?,
164            count: take(&mut pos)?,
165        };
166        Ok(Self {
167            bytes,
168            zone,
169            body_start: pos,
170        })
171    }
172
173    pub fn zone(&self) -> &ZoneMap {
174        &self.zone
175    }
176
177    /// Decode all triples in ascending order. The bytes may be corrupt (a block
178    /// from an untrusted file), so decoding is bounds-safe and stops gracefully
179    /// at the first malformed varint rather than panicking — returning whatever
180    /// prefix decoded cleanly.
181    pub fn triples(&self) -> Vec<Triple> {
182        self.try_triples().unwrap_or_default()
183    }
184
185    fn try_triples(&self) -> Option<Vec<Triple>> {
186        // `zone.count` is untrusted; each pushed triple consumes ≥1 byte, so the
187        // buffer length is a safe capacity ceiling (avoids an OOM on a bogus count).
188        let mut out = Vec::with_capacity((self.zone.count as usize).min(self.bytes.len()));
189        let mut pos = self.body_start;
190        let g = |pos: &mut usize| -> Option<u32> {
191            let (v, n) = read_uvarint(self.bytes.get(*pos..)?)?;
192            *pos += n;
193            Some(v as u32)
194        };
195        let num_a = g(&mut pos)?;
196        let mut a = 0u32;
197        for _ in 0..num_a {
198            // wrapping_add: corrupt deltas must not overflow-panic in debug builds.
199            a = a.wrapping_add(g(&mut pos)?);
200            let num_b = g(&mut pos)?;
201            let mut b = 0u32;
202            for _ in 0..num_b {
203                b = b.wrapping_add(g(&mut pos)?);
204                let num_c = g(&mut pos)?;
205                let mut c = 0u32;
206                for _ in 0..num_c {
207                    c = c.wrapping_add(g(&mut pos)?);
208                    out.push((a, b, c));
209                }
210            }
211        }
212        Some(out)
213    }
214
215    /// Build the byte-offset directory of this block's a-groups (one header
216    /// walk), enabling binary-search probes via [`scan_from`](Self::scan_from).
217    /// On corrupt bytes the walk stops early — the directory is a prefix, and
218    /// the bounds-checked cursor degrades gracefully like every other reader.
219    pub fn group_directory(&self) -> GroupDirectory {
220        let bytes = self.bytes;
221        let mut entries = Vec::new();
222        let mut p = self.body_start;
223        let mut walk = || -> Option<()> {
224            let num_a = rd(bytes, &mut p)?;
225            // `num_a` is untrusted; each group consumes ≥2 bytes, so the buffer
226            // length caps the allocation.
227            entries.reserve((num_a as usize).min(bytes.len()));
228            let mut a = 0u32;
229            for i in 0..num_a {
230                a = a.wrapping_add(rd(bytes, &mut p)?);
231                let num_b = rd(bytes, &mut p)?;
232                entries.push(DirEntry {
233                    a,
234                    pos: p,
235                    num_b,
236                    a_rem_after: num_a - 1 - i,
237                });
238                for _ in 0..num_b {
239                    rd(bytes, &mut p)?; // delta_b
240                    let nc = rd(bytes, &mut p)?;
241                    for _ in 0..nc {
242                        rd(bytes, &mut p)?;
243                    }
244                }
245            }
246            Some(())
247        };
248        let _ = walk();
249        GroupDirectory { entries }
250    }
251
252    /// Probe the block for a **bound leading component** `pa`, jumping straight
253    /// to its a-group through the directory (binary search) instead of walking
254    /// every preceding group header. Yields exactly what
255    /// `scan(Some(pa), pb, pc)` would.
256    pub fn scan_from(
257        &self,
258        dir: &GroupDirectory,
259        pa: u32,
260        pb: Option<u32>,
261        pc: Option<u32>,
262    ) -> BlockCursor<'a> {
263        let mut cursor = BlockCursor {
264            bytes: self.bytes,
265            pos: self.body_start,
266            a: 0,
267            b: 0,
268            c: 0,
269            a_rem: 0,
270            b_rem: 0,
271            c_rem: 0,
272            started: true, // a dead cursor unless the probe below arms it
273            pa: Some(pa),
274            pb,
275            pc,
276        };
277        if let Ok(i) = dir.entries.binary_search_by_key(&pa, |e| e.a) {
278            let e = &dir.entries[i];
279            // State as if the main cursor had just consumed this group's
280            // delta_a + num_b header: positioned at the first b-group.
281            cursor.pos = e.pos;
282            cursor.a = e.a;
283            cursor.a_rem = e.a_rem_after;
284            cursor.b_rem = e.num_b;
285        }
286        cursor
287    }
288
289    /// Stream the triples matching a (permuted) pattern, *without* decoding the
290    /// whole block. `pa`/`pb`/`pc` are the bound components in this block's stored
291    /// order (`None` = wildcard). The cursor walks the grouped body and:
292    ///
293    /// * **range-stops** once the leading component `a` exceeds a bound `pa` — the
294    ///   a-groups are stored ascending, so nothing later can match (the early-out
295    ///   that makes a leading-bound lookup `O(matches + preceding groups)` instead
296    ///   of `O(whole block)`);
297    /// * **group-skips** a/b groups that can't match (decoding their headers to
298    ///   advance, but never building or emitting their triples);
299    /// * **equality-filters** `pb`/`pc` without ever early-breaking inside a
300    ///   c-list, so on a *valid* block the yielded set equals what
301    ///   [`triples`](Self::triples) would yield filtered — even on corrupt bytes
302    ///   it only ever yields fewer, never panics (every read is bounds-checked).
303    ///
304    /// Yields triples in this block's stored `(a, b, c)` order; callers map back
305    /// to canonical `(s, p, o)` themselves.
306    pub fn scan(&self, pa: Option<u32>, pb: Option<u32>, pc: Option<u32>) -> BlockCursor<'a> {
307        BlockCursor {
308            bytes: self.bytes,
309            pos: self.body_start,
310            a: 0,
311            b: 0,
312            c: 0,
313            a_rem: 0,
314            b_rem: 0,
315            c_rem: 0,
316            started: false,
317            pa,
318            pb,
319            pc,
320        }
321    }
322}
323
324/// Read one uvarint at `*pos`, advancing it; `None` if truncated. Panic-free,
325/// mirroring the decoder inside [`TripleBlock::try_triples`].
326#[inline]
327fn rd(bytes: &[u8], pos: &mut usize) -> Option<u32> {
328    let (v, n) = read_uvarint(bytes.get(*pos..)?)?;
329    *pos += n;
330    Some(v as u32)
331}
332
333/// A byte-offset directory of a block's a-groups: one entry per group, sorted
334/// by leading id (the storage order). Built once per block with
335/// [`TripleBlock::group_directory`]; [`TripleBlock::scan_from`] then
336/// binary-searches it to jump a probe straight to its group.
337pub struct GroupDirectory {
338    entries: Vec<DirEntry>,
339}
340
341impl GroupDirectory {
342    /// Number of a-groups indexed.
343    pub fn len(&self) -> usize {
344        self.entries.len()
345    }
346
347    pub fn is_empty(&self) -> bool {
348        self.entries.is_empty()
349    }
350}
351
352/// One a-group: its leading id, the byte offset of its first b-group header
353/// (right after `num_b`), its b-group count, and how many a-groups follow it.
354struct DirEntry {
355    a: u32,
356    pos: usize,
357    num_b: u32,
358    a_rem_after: u32,
359}
360
361/// A lazy cursor over a [`TripleBlock`] body produced by [`TripleBlock::scan`].
362/// Holds only the block bytes and the delta-decode accumulators, so it borrows
363/// the block's bytes but allocates nothing.
364pub struct BlockCursor<'a> {
365    bytes: &'a [u8],
366    pos: usize,
367    // Running delta accumulators for the current (a, b, c).
368    a: u32,
369    b: u32,
370    c: u32,
371    // Groups/items not yet consumed at each level.
372    a_rem: u32,
373    b_rem: u32,
374    c_rem: u32,
375    started: bool,
376    pa: Option<u32>,
377    pb: Option<u32>,
378    pc: Option<u32>,
379}
380
381impl Iterator for BlockCursor<'_> {
382    type Item = Triple;
383
384    fn next(&mut self) -> Option<Triple> {
385        let bytes = self.bytes;
386        if !self.started {
387            self.a_rem = rd(bytes, &mut self.pos)?; // num_a
388            self.started = true;
389        }
390        loop {
391            // (1) Drain the c-list of the current matched (a, b) group. Every c is
392            // decoded to keep the delta chain correct; only matches are emitted.
393            while self.c_rem > 0 {
394                self.c_rem -= 1;
395                self.c = self.c.wrapping_add(rd(bytes, &mut self.pos)?);
396                if self.pc.is_none_or(|z| z == self.c) {
397                    return Some((self.a, self.b, self.c));
398                }
399            }
400            // (2) Advance to the next b-group within the current a-group.
401            while self.b_rem > 0 {
402                self.b_rem -= 1;
403                self.b = self.b.wrapping_add(rd(bytes, &mut self.pos)?);
404                let num_c = rd(bytes, &mut self.pos)?;
405                if self.pb.is_some_and(|y| y != self.b) {
406                    for _ in 0..num_c {
407                        rd(bytes, &mut self.pos)?; // group-skip: advance, never emit
408                    }
409                    continue;
410                }
411                self.c = 0; // the encoder resets prev_c per b-group
412                self.c_rem = num_c;
413                break;
414            }
415            if self.c_rem > 0 {
416                continue; // re-enter (1) to drain the matched c-list
417            }
418            // (3) Advance to the next a-group.
419            if self.a_rem == 0 {
420                return None;
421            }
422            self.a_rem -= 1;
423            self.a = self.a.wrapping_add(rd(bytes, &mut self.pos)?);
424            let num_b = rd(bytes, &mut self.pos)?;
425            self.b = 0; // the encoder resets prev_b per a-group
426            if let Some(x) = self.pa {
427                if self.a > x {
428                    return None; // range-stop: a-groups are ascending
429                }
430                if self.a < x {
431                    // skip this whole a-group (b-group headers + c-lists)
432                    for _ in 0..num_b {
433                        rd(bytes, &mut self.pos)?; // delta_b
434                        let nc = rd(bytes, &mut self.pos)?;
435                        for _ in 0..nc {
436                            rd(bytes, &mut self.pos)?;
437                        }
438                    }
439                    continue;
440                }
441            }
442            self.b_rem = num_b;
443        }
444    }
445}
446
447#[cfg(test)]
448mod tests {
449    use super::*;
450
451    fn sample() -> Vec<Triple> {
452        // Unsorted, with a duplicate, multiple b's per a and c's per b.
453        vec![
454            (5, 2, 9),
455            (1, 1, 1),
456            (1, 1, 4),
457            (1, 3, 2),
458            (5, 2, 7),
459            (1, 1, 1), // dup
460            (2, 9, 9),
461        ]
462    }
463
464    #[test]
465    fn round_trip_sorted_dedup() {
466        let mut b = TripleBlockBuilder::new();
467        for t in sample() {
468            b.push(t);
469        }
470        let bytes = b.build();
471        let blk = TripleBlock::parse(&bytes).unwrap();
472
473        let mut expected = sample();
474        expected.sort_unstable();
475        expected.dedup();
476        assert_eq!(blk.zone().count as usize, expected.len());
477        assert_eq!(blk.triples(), expected);
478    }
479
480    #[test]
481    fn zone_map_bounds_and_skipping() {
482        let mut b = TripleBlockBuilder::new();
483        for t in sample() {
484            b.push(t);
485        }
486        let bytes = b.build();
487        let blk = TripleBlock::parse(&bytes).unwrap();
488        let z = blk.zone();
489        assert_eq!((z.min_a, z.max_a), (1, 5));
490        assert_eq!((z.min_b, z.max_b), (1, 9));
491        assert_eq!((z.min_c, z.max_c), (1, 9));
492        // a=3 is within [1,5] so "maybe"; a=99 is out so skippable.
493        assert!(z.may_contain(Some(3), None, None));
494        assert!(!z.may_contain(Some(99), None, None));
495        assert!(z.may_contain(None, None, None)); // fully unbound
496    }
497
498    #[test]
499    fn empty_block() {
500        let blk_bytes = TripleBlockBuilder::new().build();
501        let blk = TripleBlock::parse(&blk_bytes).unwrap();
502        assert_eq!(blk.zone().count, 0);
503        assert!(blk.triples().is_empty());
504        assert!(blk.scan(None, None, None).next().is_none());
505        assert!(blk.scan(Some(1), None, None).next().is_none());
506    }
507
508    /// The streaming `scan` cursor must, for every bound/unbound shape, yield
509    /// exactly the full-decode result filtered by the same bounds.
510    #[test]
511    fn scan_matches_full_decode_every_shape() {
512        let mut b = TripleBlockBuilder::new();
513        for t in sample() {
514            b.push(t);
515        }
516        let bytes = b.build();
517        let blk = TripleBlock::parse(&bytes).unwrap();
518        let all = blk.triples(); // sorted, deduped, ascending (a, b, c)
519
520        let opt = |v: u32| [None, Some(v)];
521        // Probe present values and an absent one in each position.
522        for pa in opt(1).into_iter().chain([Some(5), Some(99)]) {
523            for pb in opt(1).into_iter().chain([Some(2), Some(99)]) {
524                for pc in opt(1).into_iter().chain([Some(9), Some(99)]) {
525                    let want: Vec<Triple> = all
526                        .iter()
527                        .copied()
528                        .filter(|&(a, bb, c)| {
529                            pa.is_none_or(|x| x == a)
530                                && pb.is_none_or(|x| x == bb)
531                                && pc.is_none_or(|x| x == c)
532                        })
533                        .collect();
534                    // The cursor preserves stored ascending order, so no re-sort.
535                    let got: Vec<Triple> = blk.scan(pa, pb, pc).collect();
536                    assert_eq!(got, want, "scan({pa:?},{pb:?},{pc:?})");
537                }
538            }
539        }
540    }
541
542    /// A range-stop on a bound leading component must not over-read: once `a`
543    /// passes the bound the cursor returns `None` and stops decoding.
544    #[test]
545    fn scan_range_stops_on_leading_bound() {
546        let mut b = TripleBlockBuilder::new();
547        for t in [(1, 1, 1), (1, 2, 2), (3, 1, 1), (5, 1, 1)] {
548            b.push(t);
549        }
550        let bytes = b.build();
551        let blk = TripleBlock::parse(&bytes).unwrap();
552        let got: Vec<Triple> = blk.scan(Some(1), None, None).collect();
553        assert_eq!(got, vec![(1, 1, 1), (1, 2, 2)]);
554        // A bound `a` between stored groups yields nothing (and stops early).
555        assert!(blk.scan(Some(2), None, None).next().is_none());
556        assert!(blk.scan(Some(99), None, None).next().is_none());
557    }
558
559    /// Truncations and byte corruptions must never panic the cursor — it only
560    /// ever yields a clean prefix (every read is bounds-checked).
561    #[test]
562    fn scan_never_panics_on_bad_bytes() {
563        let mut b = TripleBlockBuilder::new();
564        for t in sample() {
565            b.push(t);
566        }
567        let bytes = b.build();
568        for len in 0..bytes.len() {
569            if let Ok(blk) = TripleBlock::parse(&bytes[..len]) {
570                for pat in [(None, None, None), (Some(1u32), Some(1u32), Some(1u32))] {
571                    let _ = blk.scan(pat.0, pat.1, pat.2).count();
572                }
573            }
574        }
575        for i in 0..bytes.len() {
576            for v in [0x00u8, 0xff, 0x80, 0x7f] {
577                let mut bad = bytes.clone();
578                bad[i] = v;
579                if let Ok(blk) = TripleBlock::parse(&bad) {
580                    let _ = blk.scan(None, None, None).count();
581                    let _ = blk.scan(Some(1), None, None).count();
582                }
583            }
584        }
585    }
586}