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    /// **Resume** a leading-*unbound* scan at the first a-group whose leading id
290    /// is `>= from_a`, jumping there through the directory instead of re-walking
291    /// the groups already consumed. Yields exactly the tail of
292    /// `scan(None, pb, pc)` that begins at that group; `from_a = 0` is the whole
293    /// block.
294    ///
295    /// This is what lets a *batched* cursor stop in the middle of a block and
296    /// pick up where it left off without rescanning — the difference between an
297    /// O(n) resumable scan and an O(n²/batch) skip-and-take. Its bound-leading
298    /// twin is [`scan_from`](Self::scan_from), which pins `pa` to one group;
299    /// here `pa` stays unbound and `from_a` is only a starting point, so the
300    /// cursor runs on through every later group.
301    pub fn scan_resume(
302        &self,
303        dir: &GroupDirectory,
304        from_a: u32,
305        pb: Option<u32>,
306        pc: Option<u32>,
307    ) -> BlockCursor<'a> {
308        let mut cursor = BlockCursor {
309            bytes: self.bytes,
310            pos: self.body_start,
311            a: 0,
312            b: 0,
313            c: 0,
314            a_rem: 0,
315            b_rem: 0,
316            c_rem: 0,
317            started: true, // a dead cursor unless the probe below arms it
318            pa: None,
319            pb,
320            pc,
321        };
322        // Groups are stored ascending by `a`, so the first entry not less than
323        // `from_a` is the resume point (and `Err`/`Ok` of a binary search are the
324        // same answer here — unlike `scan_from`, a miss is not "no matches").
325        let i = dir.entries.partition_point(|e| e.a < from_a);
326        if let Some(e) = dir.entries.get(i) {
327            // State as if the main cursor had just consumed this group's
328            // delta_a + num_b header: positioned at the first b-group.
329            cursor.pos = e.pos;
330            cursor.a = e.a;
331            cursor.a_rem = e.a_rem_after;
332            cursor.b_rem = e.num_b;
333        }
334        cursor
335    }
336
337    /// Stream the triples matching a (permuted) pattern, *without* decoding the
338    /// whole block. `pa`/`pb`/`pc` are the bound components in this block's stored
339    /// order (`None` = wildcard). The cursor walks the grouped body and:
340    ///
341    /// * **range-stops** once the leading component `a` exceeds a bound `pa` — the
342    ///   a-groups are stored ascending, so nothing later can match (the early-out
343    ///   that makes a leading-bound lookup `O(matches + preceding groups)` instead
344    ///   of `O(whole block)`);
345    /// * **group-skips** a/b groups that can't match (decoding their headers to
346    ///   advance, but never building or emitting their triples);
347    /// * **equality-filters** `pb`/`pc` without ever early-breaking inside a
348    ///   c-list, so on a *valid* block the yielded set equals what
349    ///   [`triples`](Self::triples) would yield filtered — even on corrupt bytes
350    ///   it only ever yields fewer, never panics (every read is bounds-checked).
351    ///
352    /// Yields triples in this block's stored `(a, b, c)` order; callers map back
353    /// to canonical `(s, p, o)` themselves.
354    pub fn scan(&self, pa: Option<u32>, pb: Option<u32>, pc: Option<u32>) -> BlockCursor<'a> {
355        BlockCursor {
356            bytes: self.bytes,
357            pos: self.body_start,
358            a: 0,
359            b: 0,
360            c: 0,
361            a_rem: 0,
362            b_rem: 0,
363            c_rem: 0,
364            started: false,
365            pa,
366            pb,
367            pc,
368        }
369    }
370}
371
372/// Read one uvarint at `*pos`, advancing it; `None` if truncated. Panic-free,
373/// mirroring the decoder inside [`TripleBlock::try_triples`].
374#[inline]
375fn rd(bytes: &[u8], pos: &mut usize) -> Option<u32> {
376    let (v, n) = read_uvarint(bytes.get(*pos..)?)?;
377    *pos += n;
378    Some(v as u32)
379}
380
381/// A byte-offset directory of a block's a-groups: one entry per group, sorted
382/// by leading id (the storage order). Built once per block with
383/// [`TripleBlock::group_directory`]; [`TripleBlock::scan_from`] then
384/// binary-searches it to jump a probe straight to its group.
385pub struct GroupDirectory {
386    entries: Vec<DirEntry>,
387}
388
389impl GroupDirectory {
390    /// Number of a-groups indexed.
391    pub fn len(&self) -> usize {
392        self.entries.len()
393    }
394
395    pub fn is_empty(&self) -> bool {
396        self.entries.is_empty()
397    }
398}
399
400/// One a-group: its leading id, the byte offset of its first b-group header
401/// (right after `num_b`), its b-group count, and how many a-groups follow it.
402struct DirEntry {
403    a: u32,
404    pos: usize,
405    num_b: u32,
406    a_rem_after: u32,
407}
408
409/// A lazy cursor over a [`TripleBlock`] body produced by [`TripleBlock::scan`].
410/// Holds only the block bytes and the delta-decode accumulators, so it borrows
411/// the block's bytes but allocates nothing.
412pub struct BlockCursor<'a> {
413    bytes: &'a [u8],
414    pos: usize,
415    // Running delta accumulators for the current (a, b, c).
416    a: u32,
417    b: u32,
418    c: u32,
419    // Groups/items not yet consumed at each level.
420    a_rem: u32,
421    b_rem: u32,
422    c_rem: u32,
423    started: bool,
424    pa: Option<u32>,
425    pb: Option<u32>,
426    pc: Option<u32>,
427}
428
429impl Iterator for BlockCursor<'_> {
430    type Item = Triple;
431
432    fn next(&mut self) -> Option<Triple> {
433        let bytes = self.bytes;
434        if !self.started {
435            self.a_rem = rd(bytes, &mut self.pos)?; // num_a
436            self.started = true;
437        }
438        loop {
439            // (1) Drain the c-list of the current matched (a, b) group. Every c is
440            // decoded to keep the delta chain correct; only matches are emitted.
441            while self.c_rem > 0 {
442                self.c_rem -= 1;
443                self.c = self.c.wrapping_add(rd(bytes, &mut self.pos)?);
444                if self.pc.is_none_or(|z| z == self.c) {
445                    return Some((self.a, self.b, self.c));
446                }
447            }
448            // (2) Advance to the next b-group within the current a-group.
449            while self.b_rem > 0 {
450                self.b_rem -= 1;
451                self.b = self.b.wrapping_add(rd(bytes, &mut self.pos)?);
452                let num_c = rd(bytes, &mut self.pos)?;
453                if self.pb.is_some_and(|y| y != self.b) {
454                    for _ in 0..num_c {
455                        rd(bytes, &mut self.pos)?; // group-skip: advance, never emit
456                    }
457                    continue;
458                }
459                self.c = 0; // the encoder resets prev_c per b-group
460                self.c_rem = num_c;
461                break;
462            }
463            if self.c_rem > 0 {
464                continue; // re-enter (1) to drain the matched c-list
465            }
466            // (3) Advance to the next a-group.
467            if self.a_rem == 0 {
468                return None;
469            }
470            self.a_rem -= 1;
471            self.a = self.a.wrapping_add(rd(bytes, &mut self.pos)?);
472            let num_b = rd(bytes, &mut self.pos)?;
473            self.b = 0; // the encoder resets prev_b per a-group
474            if let Some(x) = self.pa {
475                if self.a > x {
476                    return None; // range-stop: a-groups are ascending
477                }
478                if self.a < x {
479                    // skip this whole a-group (b-group headers + c-lists)
480                    for _ in 0..num_b {
481                        rd(bytes, &mut self.pos)?; // delta_b
482                        let nc = rd(bytes, &mut self.pos)?;
483                        for _ in 0..nc {
484                            rd(bytes, &mut self.pos)?;
485                        }
486                    }
487                    continue;
488                }
489            }
490            self.b_rem = num_b;
491        }
492    }
493}
494
495#[cfg(test)]
496mod tests {
497    use super::*;
498
499    fn sample() -> Vec<Triple> {
500        // Unsorted, with a duplicate, multiple b's per a and c's per b.
501        vec![
502            (5, 2, 9),
503            (1, 1, 1),
504            (1, 1, 4),
505            (1, 3, 2),
506            (5, 2, 7),
507            (1, 1, 1), // dup
508            (2, 9, 9),
509        ]
510    }
511
512    #[test]
513    fn round_trip_sorted_dedup() {
514        let mut b = TripleBlockBuilder::new();
515        for t in sample() {
516            b.push(t);
517        }
518        let bytes = b.build();
519        let blk = TripleBlock::parse(&bytes).unwrap();
520
521        let mut expected = sample();
522        expected.sort_unstable();
523        expected.dedup();
524        assert_eq!(blk.zone().count as usize, expected.len());
525        assert_eq!(blk.triples(), expected);
526    }
527
528    #[test]
529    fn zone_map_bounds_and_skipping() {
530        let mut b = TripleBlockBuilder::new();
531        for t in sample() {
532            b.push(t);
533        }
534        let bytes = b.build();
535        let blk = TripleBlock::parse(&bytes).unwrap();
536        let z = blk.zone();
537        assert_eq!((z.min_a, z.max_a), (1, 5));
538        assert_eq!((z.min_b, z.max_b), (1, 9));
539        assert_eq!((z.min_c, z.max_c), (1, 9));
540        // a=3 is within [1,5] so "maybe"; a=99 is out so skippable.
541        assert!(z.may_contain(Some(3), None, None));
542        assert!(!z.may_contain(Some(99), None, None));
543        assert!(z.may_contain(None, None, None)); // fully unbound
544    }
545
546    #[test]
547    fn empty_block() {
548        let blk_bytes = TripleBlockBuilder::new().build();
549        let blk = TripleBlock::parse(&blk_bytes).unwrap();
550        assert_eq!(blk.zone().count, 0);
551        assert!(blk.triples().is_empty());
552        assert!(blk.scan(None, None, None).next().is_none());
553        assert!(blk.scan(Some(1), None, None).next().is_none());
554    }
555
556    /// The streaming `scan` cursor must, for every bound/unbound shape, yield
557    /// exactly the full-decode result filtered by the same bounds.
558    #[test]
559    fn scan_matches_full_decode_every_shape() {
560        let mut b = TripleBlockBuilder::new();
561        for t in sample() {
562            b.push(t);
563        }
564        let bytes = b.build();
565        let blk = TripleBlock::parse(&bytes).unwrap();
566        let all = blk.triples(); // sorted, deduped, ascending (a, b, c)
567
568        let opt = |v: u32| [None, Some(v)];
569        // Probe present values and an absent one in each position.
570        for pa in opt(1).into_iter().chain([Some(5), Some(99)]) {
571            for pb in opt(1).into_iter().chain([Some(2), Some(99)]) {
572                for pc in opt(1).into_iter().chain([Some(9), Some(99)]) {
573                    let want: Vec<Triple> = all
574                        .iter()
575                        .copied()
576                        .filter(|&(a, bb, c)| {
577                            pa.is_none_or(|x| x == a)
578                                && pb.is_none_or(|x| x == bb)
579                                && pc.is_none_or(|x| x == c)
580                        })
581                        .collect();
582                    // The cursor preserves stored ascending order, so no re-sort.
583                    let got: Vec<Triple> = blk.scan(pa, pb, pc).collect();
584                    assert_eq!(got, want, "scan({pa:?},{pb:?},{pc:?})");
585                }
586            }
587        }
588    }
589
590    /// A range-stop on a bound leading component must not over-read: once `a`
591    /// passes the bound the cursor returns `None` and stops decoding.
592    #[test]
593    fn scan_range_stops_on_leading_bound() {
594        let mut b = TripleBlockBuilder::new();
595        for t in [(1, 1, 1), (1, 2, 2), (3, 1, 1), (5, 1, 1)] {
596            b.push(t);
597        }
598        let bytes = b.build();
599        let blk = TripleBlock::parse(&bytes).unwrap();
600        let got: Vec<Triple> = blk.scan(Some(1), None, None).collect();
601        assert_eq!(got, vec![(1, 1, 1), (1, 2, 2)]);
602        // A bound `a` between stored groups yields nothing (and stops early).
603        assert!(blk.scan(Some(2), None, None).next().is_none());
604        assert!(blk.scan(Some(99), None, None).next().is_none());
605    }
606
607    /// Truncations and byte corruptions must never panic the cursor — it only
608    /// ever yields a clean prefix (every read is bounds-checked).
609    #[test]
610    fn scan_never_panics_on_bad_bytes() {
611        let mut b = TripleBlockBuilder::new();
612        for t in sample() {
613            b.push(t);
614        }
615        let bytes = b.build();
616        for len in 0..bytes.len() {
617            if let Ok(blk) = TripleBlock::parse(&bytes[..len]) {
618                for pat in [(None, None, None), (Some(1u32), Some(1u32), Some(1u32))] {
619                    let _ = blk.scan(pat.0, pat.1, pat.2).count();
620                }
621            }
622        }
623        for i in 0..bytes.len() {
624            for v in [0x00u8, 0xff, 0x80, 0x7f] {
625                let mut bad = bytes.clone();
626                bad[i] = v;
627                if let Ok(blk) = TripleBlock::parse(&bad) {
628                    let _ = blk.scan(None, None, None).count();
629                    let _ = blk.scan(Some(1), None, None).count();
630                }
631            }
632        }
633    }
634}