Skip to main content

terminus_store/layer/internal/
object_iterator.rs

1use crate::layer::*;
2use std::convert::TryInto;
3use tdb_succinct::*;
4
5#[derive(Clone)]
6pub struct InternalLayerTripleObjectIterator {
7    subjects: Option<MonotonicLogArray>,
8    objects: Option<MonotonicLogArray>,
9    o_ps_adjacency_list: AdjacencyList,
10    s_p_adjacency_list: AdjacencyList,
11    stop_at_boundary: bool,
12
13    o_position: u64,
14    o_ps_position: u64,
15    peeked: Option<IdTriple>,
16    hit_boundary: bool,
17}
18
19impl InternalLayerTripleObjectIterator {
20    pub fn new(
21        subjects: Option<MonotonicLogArray>,
22        objects: Option<MonotonicLogArray>,
23        o_ps_adjacency_list: AdjacencyList,
24        s_p_adjacency_list: AdjacencyList,
25        stop_at_boundary: bool,
26    ) -> Self {
27        Self {
28            subjects,
29            objects,
30            o_ps_adjacency_list,
31            s_p_adjacency_list,
32            stop_at_boundary,
33            o_position: 0,
34            o_ps_position: 0,
35            peeked: None,
36            hit_boundary: false,
37        }
38    }
39
40    pub fn seek_object(mut self, object: u64) -> Self {
41        self.seek_object_ref(object);
42
43        self
44    }
45
46    pub fn seek_object_ref(&mut self, object: u64) {
47        self.peeked = None;
48        self.hit_boundary = false;
49
50        if object == 0 {
51            self.o_position = 0;
52            self.o_ps_position = 0;
53
54            return;
55        }
56
57        self.o_position = match self.objects.as_ref() {
58            None => object - 1,
59            Some(objects) => objects.nearest_index_of(object) as u64,
60        };
61
62        if self.o_position >= self.o_ps_adjacency_list.left_count() as u64 {
63            self.o_ps_position = self.o_ps_adjacency_list.right_count() as u64;
64        } else {
65            self.o_ps_position = self.o_ps_adjacency_list.offset_for(self.o_position + 1);
66        }
67    }
68
69    pub fn peek(&mut self) -> Option<&IdTriple> {
70        self.peeked = self.next();
71
72        self.peeked.as_ref()
73    }
74
75    pub fn stop_at_boundary(mut self, stop: bool) -> Self {
76        self.stop_at_boundary_ref(stop);
77
78        self
79    }
80
81    pub fn stop_at_boundary_ref(&mut self, stop: bool) {
82        self.stop_at_boundary = stop;
83    }
84
85    pub fn continue_from_boundary(&mut self) {
86        self.hit_boundary = false;
87    }
88}
89
90impl Iterator for InternalLayerTripleObjectIterator {
91    type Item = IdTriple;
92
93    fn next(&mut self) -> Option<IdTriple> {
94        if self.peeked.is_some() {
95            let peeked = self.peeked;
96            self.peeked = None;
97
98            return peeked;
99        }
100
101        loop {
102            if self.stop_at_boundary && self.hit_boundary {
103                return None;
104            }
105
106            if self.o_ps_position >= self.o_ps_adjacency_list.right_count() as u64 {
107                return None;
108            } else {
109                let o_pos = self.o_position;
110                let o_ps_bit = self.o_ps_adjacency_list.bit_at_pos(self.o_ps_position);
111                let sp_pair_num = self.o_ps_adjacency_list.num_at_pos(self.o_ps_position);
112
113                if o_ps_bit {
114                    self.o_position += 1;
115                    if self.stop_at_boundary {
116                        self.hit_boundary = true;
117                    }
118                }
119                self.o_ps_position += 1;
120
121                if sp_pair_num == 0 {
122                    continue;
123                }
124
125                let (mapped_subject, predicate) =
126                    self.s_p_adjacency_list.pair_at_pos(sp_pair_num - 1);
127
128                let subject = match self.subjects.as_ref() {
129                    Some(subjects) => subjects.entry(mapped_subject as usize - 1),
130                    None => mapped_subject,
131                };
132
133                let object = match self.objects.as_ref() {
134                    Some(objects) => objects.entry(o_pos.try_into().unwrap()),
135                    None => o_pos + 1,
136                };
137
138                return Some(IdTriple::new(subject, predicate, object));
139            }
140        }
141    }
142}
143
144pub struct OptInternalLayerTripleObjectIterator(pub Option<InternalLayerTripleObjectIterator>);
145
146impl OptInternalLayerTripleObjectIterator {
147    pub fn seek_object(self, object: u64) -> Self {
148        OptInternalLayerTripleObjectIterator(self.0.map(|i| i.seek_object(object)))
149    }
150
151    pub fn seek_object_ref(&mut self, object: u64) {
152        if let Some(i) = self.0.as_mut() {
153            i.seek_object_ref(object)
154        }
155    }
156
157    pub fn stop_at_boundary(self, stop: bool) -> Self {
158        OptInternalLayerTripleObjectIterator(self.0.map(|i| i.stop_at_boundary(stop)))
159    }
160
161    pub fn stop_at_boundary_ref(&mut self, stop: bool) {
162        if let Some(i) = self.0.as_mut() {
163            i.stop_at_boundary_ref(stop)
164        }
165    }
166
167    pub fn peek(&mut self) -> Option<&IdTriple> {
168        self.0.as_mut().and_then(|i| i.peek())
169    }
170}
171
172impl Iterator for OptInternalLayerTripleObjectIterator {
173    type Item = IdTriple;
174
175    fn next(&mut self) -> Option<IdTriple> {
176        match self.0.as_mut() {
177            Some(i) => i.next(),
178            None => None,
179        }
180    }
181}
182
183pub struct InternalTripleObjectIterator {
184    positives: Vec<OptInternalLayerTripleObjectIterator>,
185    negatives: Vec<OptInternalLayerTripleObjectIterator>,
186}
187
188impl InternalTripleObjectIterator {
189    pub fn from_layer(layer: &InternalLayer) -> Self {
190        let stack_size = layer.layer_stack_size();
191        let mut positives = Vec::with_capacity(stack_size);
192        let mut negatives = Vec::with_capacity(stack_size);
193        positives.push(layer.internal_triple_additions_by_object());
194        negatives.push(layer.internal_triple_removals_by_object());
195
196        let mut layer_opt = layer.immediate_parent();
197
198        while layer_opt.is_some() {
199            positives.push(layer_opt.unwrap().internal_triple_additions_by_object());
200            negatives.push(layer_opt.unwrap().internal_triple_removals_by_object());
201
202            layer_opt = layer_opt.unwrap().immediate_parent();
203        }
204
205        Self {
206            positives,
207            negatives,
208        }
209    }
210
211    pub fn seek_object(mut self, object: u64) -> Self {
212        for p in self.positives.iter_mut() {
213            p.seek_object_ref(object);
214        }
215
216        for n in self.negatives.iter_mut() {
217            n.seek_object_ref(object);
218        }
219
220        self
221    }
222}
223
224impl Iterator for InternalTripleObjectIterator {
225    type Item = IdTriple;
226
227    fn next(&mut self) -> Option<IdTriple> {
228        'outer: loop {
229            // find the lowest triple.
230            // if that triple appears multiple times, we want the most recent one, which should be the one appearing the earliest in the positives list.
231            let lowest_index = self
232                .positives
233                .iter_mut()
234                .map(|p| p.peek())
235                .enumerate()
236                .filter(|(_, elt)| elt.is_some())
237                .min_by_key(|(_, elt)| {
238                    let e = elt.unwrap();
239                    // we need to restructure because we need to order by object
240                    (e.object, e.subject, e.predicate)
241                })
242                .map(|(index, _)| index);
243
244            match lowest_index {
245                None => return None,
246                Some(lowest_index) => {
247                    let lowest = self.positives[lowest_index].next().unwrap();
248                    // check all negative layers below the lowest_index for a removal
249                    // if there's a removal, we continue after advancing. if not, it is the result.
250                    // we can be sure that there's only one removal, or we'd have found another addition.
251                    for iter in self.negatives[0..lowest_index].iter_mut() {
252                        if iter.peek() == Some(&lowest) {
253                            iter.next().unwrap();
254                            continue 'outer;
255                        }
256                    }
257
258                    return Some(lowest);
259                }
260            }
261        }
262    }
263}
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268    use crate::layer::base::base_tests::*;
269    use crate::storage::memory::*;
270    use crate::storage::*;
271
272    async fn example_base_layer_files() -> BaseLayerFiles<MemoryBackedStore> {
273        let nodes = vec!["aaaaa", "baa", "bbbbb", "ccccc", "mooo"];
274        let predicates = vec!["abcde", "fghij", "klmno", "lll"];
275        let values = vec!["chicken", "cow", "dog", "pig", "zebra"];
276
277        let base_layer_files = base_layer_files();
278
279        let mut builder = BaseLayerFileBuilder::from_files(&base_layer_files)
280            .await
281            .unwrap();
282
283        builder.add_nodes(nodes.into_iter().map(|s| s.to_string()));
284        builder.add_predicates(predicates.into_iter().map(|s| s.to_string()));
285        builder.add_values(values.into_iter().map(|s| String::make_entry(&s)));
286        let mut builder = builder.into_phase2().await.unwrap();
287
288        builder.add_triple(1, 1, 2).await.unwrap();
289        builder.add_triple(2, 1, 2).await.unwrap();
290        builder.add_triple(2, 1, 3).await.unwrap();
291        builder.add_triple(2, 1, 5).await.unwrap();
292        builder.add_triple(2, 3, 6).await.unwrap();
293        builder.add_triple(3, 2, 5).await.unwrap();
294        builder.add_triple(3, 3, 6).await.unwrap();
295        builder.add_triple(4, 1, 5).await.unwrap();
296        builder.add_triple(4, 3, 6).await.unwrap();
297        builder.finalize().await.unwrap();
298
299        base_layer_files
300    }
301
302    async fn example_base_layer() -> InternalLayer {
303        let base_layer_files = example_base_layer_files().await;
304
305        let layer = BaseLayer::load_from_files([1, 2, 3, 4, 5], &base_layer_files)
306            .await
307            .unwrap();
308
309        layer
310    }
311
312    #[tokio::test]
313    async fn object_iterator() {
314        let base_layer = example_base_layer().await;
315
316        let iterator = base_layer.internal_triple_additions_by_object();
317        let triples: Vec<_> = iterator.collect();
318
319        let expected = vec![
320            IdTriple::new(1, 1, 2),
321            IdTriple::new(2, 1, 2),
322            IdTriple::new(2, 1, 3),
323            IdTriple::new(2, 1, 5),
324            IdTriple::new(3, 2, 5),
325            IdTriple::new(4, 1, 5),
326            IdTriple::new(2, 3, 6),
327            IdTriple::new(3, 3, 6),
328            IdTriple::new(4, 3, 6),
329        ];
330        assert_eq!(expected, triples);
331    }
332
333    #[tokio::test]
334    async fn object_iterator_seek() {
335        let base_layer = example_base_layer().await;
336
337        let iterator = base_layer.internal_triple_additions_by_object();
338        let triples: Vec<_> = iterator.seek_object(5).collect();
339
340        let expected = vec![
341            IdTriple::new(2, 1, 5),
342            IdTriple::new(3, 2, 5),
343            IdTriple::new(4, 1, 5),
344            IdTriple::new(2, 3, 6),
345            IdTriple::new(3, 3, 6),
346            IdTriple::new(4, 3, 6),
347        ];
348        assert_eq!(expected, triples);
349    }
350
351    #[tokio::test]
352    async fn object_iterator_seek_0() {
353        let base_layer = example_base_layer().await;
354
355        let iterator = base_layer.internal_triple_additions_by_object();
356        let triples: Vec<_> = iterator.seek_object(0).collect();
357
358        let expected = vec![
359            IdTriple::new(1, 1, 2),
360            IdTriple::new(2, 1, 2),
361            IdTriple::new(2, 1, 3),
362            IdTriple::new(2, 1, 5),
363            IdTriple::new(3, 2, 5),
364            IdTriple::new(4, 1, 5),
365            IdTriple::new(2, 3, 6),
366            IdTriple::new(3, 3, 6),
367            IdTriple::new(4, 3, 6),
368        ];
369        assert_eq!(expected, triples);
370    }
371
372    #[tokio::test]
373    async fn object_iterator_seek_before_begin() {
374        let base_layer = example_base_layer().await;
375
376        let iterator = base_layer.internal_triple_additions_by_object();
377        let triples: Vec<_> = iterator.seek_object(1).collect();
378
379        let expected = vec![
380            IdTriple::new(1, 1, 2),
381            IdTriple::new(2, 1, 2),
382            IdTriple::new(2, 1, 3),
383            IdTriple::new(2, 1, 5),
384            IdTriple::new(3, 2, 5),
385            IdTriple::new(4, 1, 5),
386            IdTriple::new(2, 3, 6),
387            IdTriple::new(3, 3, 6),
388            IdTriple::new(4, 3, 6),
389        ];
390        assert_eq!(expected, triples);
391    }
392
393    #[tokio::test]
394    async fn object_iterator_seek_nonexistent() {
395        let base_layer = example_base_layer().await;
396
397        let iterator = base_layer.internal_triple_additions_by_object();
398        let triples: Vec<_> = iterator.seek_object(4).collect();
399
400        let expected = vec![
401            IdTriple::new(2, 1, 5),
402            IdTriple::new(3, 2, 5),
403            IdTriple::new(4, 1, 5),
404            IdTriple::new(2, 3, 6),
405            IdTriple::new(3, 3, 6),
406            IdTriple::new(4, 3, 6),
407        ];
408        assert_eq!(expected, triples);
409    }
410
411    #[tokio::test]
412    async fn object_iterator_seek_past_end() {
413        let base_layer = example_base_layer().await;
414
415        let iterator = base_layer.internal_triple_additions_by_object();
416        let triples: Vec<_> = iterator.seek_object(7).collect();
417        assert!(triples.is_empty());
418    }
419
420    #[tokio::test]
421    async fn object_additions_iterator_for_object() {
422        let base_layer = example_base_layer().await;
423
424        let triples: Vec<_> = base_layer.internal_triple_additions_o(5).collect();
425
426        let expected = vec![
427            IdTriple::new(2, 1, 5),
428            IdTriple::new(3, 2, 5),
429            IdTriple::new(4, 1, 5),
430        ];
431
432        assert_eq!(expected, triples);
433    }
434
435    #[tokio::test]
436    async fn object_additions_iterator_for_nonexistent_object() {
437        let base_layer = example_base_layer().await;
438
439        let triples: Vec<_> = base_layer.internal_triple_additions_o(4).collect();
440
441        assert!(triples.is_empty());
442    }
443
444    #[tokio::test]
445    async fn combined_iterator_for_object() {
446        let store = MemoryLayerStore::new();
447        let mut builder = store.create_base_layer().await.unwrap();
448        let base_name = builder.name();
449
450        builder.add_value_triple(ValueTriple::new_string_value("cow", "says", "moo"));
451        builder.add_value_triple(ValueTriple::new_string_value("duck", "says", "quack"));
452        builder.add_value_triple(ValueTriple::new_node("cow", "likes", "duck"));
453        builder.add_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
454        builder.commit_boxed().await.unwrap();
455
456        builder = store.create_child_layer(base_name).await.unwrap();
457        let child1_name = builder.name();
458
459        builder.add_value_triple(ValueTriple::new_string_value("horse", "says", "neigh"));
460        builder.add_value_triple(ValueTriple::new_node("horse", "likes", "horse"));
461        builder.add_value_triple(ValueTriple::new_node("horse", "likes", "cow"));
462        builder.commit_boxed().await.unwrap();
463
464        builder = store.create_child_layer(child1_name).await.unwrap();
465        let child2_name = builder.name();
466
467        builder.remove_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
468        builder.add_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
469        builder.commit_boxed().await.unwrap();
470
471        builder = store.create_child_layer(child2_name).await.unwrap();
472        let child3_name = builder.name();
473
474        builder.remove_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
475        builder.add_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
476        builder.commit_boxed().await.unwrap();
477
478        builder = store.create_child_layer(child3_name).await.unwrap();
479        let child4_name = builder.name();
480
481        builder.remove_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
482        builder.add_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
483        builder.add_value_triple(ValueTriple::new_node("field", "contains", "cow"));
484        builder.commit_boxed().await.unwrap();
485
486        let layer = store.get_layer(child4_name).await.unwrap().unwrap();
487
488        let object_id = layer.object_node_id("cow").unwrap();
489        let triples: Vec<_> = layer
490            .triples_o(object_id)
491            .map(|t| layer.id_triple_to_string(&t).unwrap())
492            .collect();
493
494        let expected = vec![
495            ValueTriple::new_node("duck", "likes", "cow"),
496            ValueTriple::new_node("horse", "likes", "cow"),
497            ValueTriple::new_node("field", "contains", "cow"),
498        ];
499
500        assert_eq!(expected, triples);
501    }
502}