Skip to main content

terminus_store/layer/internal/
predicate_iterator.rs

1use crate::layer::*;
2use std::convert::TryInto;
3use tdb_succinct::*;
4
5#[derive(Clone)]
6pub struct InternalLayerTriplePredicateIterator {
7    len: usize,
8    predicate_wavelet_lookup: WaveletLookup,
9    subject_iterator: InternalLayerTripleSubjectIterator,
10    predicate_pos: u64,
11    sp_boundary: bool,
12    peeked: Option<IdTriple>,
13}
14
15impl InternalLayerTriplePredicateIterator {
16    pub fn new(
17        predicate_wavelet_lookup: WaveletLookup,
18        subjects: Option<MonotonicLogArray>,
19        s_p_adjacency_list: AdjacencyList,
20        sp_o_adjacency_list: AdjacencyList,
21    ) -> Self {
22        let len = predicate_wavelet_lookup.len();
23        let subject_iterator = InternalLayerTripleSubjectIterator::new(
24            subjects,
25            s_p_adjacency_list,
26            sp_o_adjacency_list,
27        );
28
29        Self {
30            len,
31            predicate_wavelet_lookup,
32            subject_iterator,
33            predicate_pos: 0,
34            sp_boundary: true,
35            peeked: None,
36        }
37    }
38
39    fn next_pos(&mut self) -> bool {
40        if self.predicate_pos >= self.len as u64 {
41            return false;
42        }
43
44        let s_p_pos = self
45            .predicate_wavelet_lookup
46            .entry(self.predicate_pos.try_into().unwrap());
47        self.subject_iterator.seek_s_p_pos(s_p_pos);
48
49        self.predicate_pos += 1;
50
51        true
52    }
53
54    pub fn peek(&mut self) -> Option<&IdTriple> {
55        self.peeked = self.next();
56
57        self.peeked.as_ref()
58    }
59}
60
61impl Iterator for InternalLayerTriplePredicateIterator {
62    type Item = IdTriple;
63
64    fn next(&mut self) -> Option<IdTriple> {
65        if self.peeked.is_some() {
66            let peeked = self.peeked;
67            self.peeked = None;
68
69            return peeked;
70        }
71        if self.sp_boundary {
72            // We have reached the end of our previous lookup.
73            // we need to look up the next predicate.
74            if !self.next_pos() {
75                // but there is no next pedicate, so return here.
76                return None;
77            }
78        }
79        let result = self.subject_iterator.next();
80
81        // Check the next entry of the subject iterator.
82        // If it has a different subject and predicate than before, we
83        // set the sp_boundary flag, ensuring that upon the subsequent
84        // next() call, we move on to the next predicate.
85        let next = self.subject_iterator.peek();
86        if next.is_none()
87            || next.map(|t| (t.subject, t.predicate)) != result.map(|t| (t.subject, t.predicate))
88        {
89            self.sp_boundary = true;
90        } else {
91            self.sp_boundary = false;
92        }
93
94        result
95    }
96}
97
98#[derive(Clone)]
99pub struct OptInternalLayerTriplePredicateIterator(
100    pub Option<InternalLayerTriplePredicateIterator>,
101);
102
103impl OptInternalLayerTriplePredicateIterator {
104    pub fn peek(&mut self) -> Option<&IdTriple> {
105        self.0.as_mut().and_then(|i| i.peek())
106    }
107}
108
109impl Iterator for OptInternalLayerTriplePredicateIterator {
110    type Item = IdTriple;
111
112    fn next(&mut self) -> Option<IdTriple> {
113        self.0.as_mut().and_then(|i| i.next())
114    }
115}
116
117pub struct InternalTriplePredicateIterator {
118    positives: Vec<OptInternalLayerTriplePredicateIterator>,
119    negatives: Vec<OptInternalLayerTriplePredicateIterator>,
120}
121
122impl InternalTriplePredicateIterator {
123    pub fn from_layer(layer: &InternalLayer, predicate: u64) -> Self {
124        let stack_size = layer.layer_stack_size();
125        let mut positives = Vec::with_capacity(stack_size);
126        let mut negatives = Vec::with_capacity(stack_size);
127        positives.push(layer.internal_triple_additions_p(predicate));
128        negatives.push(layer.internal_triple_removals_p(predicate));
129
130        let mut layer_opt = layer.immediate_parent();
131
132        while layer_opt.is_some() {
133            positives.push(layer_opt.unwrap().internal_triple_additions_p(predicate));
134            negatives.push(layer_opt.unwrap().internal_triple_removals_p(predicate));
135
136            layer_opt = layer_opt.unwrap().immediate_parent();
137        }
138
139        Self {
140            positives,
141            negatives,
142        }
143    }
144}
145
146impl Iterator for InternalTriplePredicateIterator {
147    type Item = IdTriple;
148
149    fn next(&mut self) -> Option<IdTriple> {
150        'outer: loop {
151            // find the lowest triple.
152            // if that triple appears multiple times, we want the most recent one, which should be the one appearing the earliest in the positives list.
153            let lowest_index = self
154                .positives
155                .iter_mut()
156                .map(|p| p.peek())
157                .enumerate()
158                .filter(|(_, elt)| elt.is_some())
159                .min_by_key(|(_, elt)| elt.unwrap())
160                .map(|(index, _)| index);
161
162            match lowest_index {
163                None => return None,
164                Some(lowest_index) => {
165                    let lowest = self.positives[lowest_index].next().unwrap();
166                    // check all negative layers below the lowest_index for a removal
167                    // if there's a removal, we continue after advancing. if not, it is the result.
168                    // we can be sure that there's only one removal, or we'd have found another addition.
169                    for iter in self.negatives[0..lowest_index].iter_mut() {
170                        if iter.peek() == Some(&lowest) {
171                            iter.next().unwrap();
172                            continue 'outer;
173                        }
174                    }
175
176                    return Some(lowest);
177                }
178            }
179        }
180    }
181}
182
183#[cfg(test)]
184mod tests {
185    use crate::layer::base::base_tests::*;
186    use crate::layer::child::child_tests::*;
187    use crate::layer::*;
188
189    use std::sync::Arc;
190
191    #[tokio::test]
192    async fn base_triple_predicate_iterator() {
193        let base_layer: InternalLayer = example_base_layer().await.into();
194
195        let triples: Vec<_> = base_layer.internal_triple_additions_p(3).collect();
196        let expected = vec![
197            IdTriple::new(2, 3, 6),
198            IdTriple::new(3, 3, 6),
199            IdTriple::new(4, 3, 6),
200        ];
201
202        assert_eq!(expected, triples);
203    }
204
205    async fn child_layer() -> InternalLayer {
206        let base_layer = example_base_layer().await;
207        let parent: Arc<InternalLayer> = Arc::new(base_layer.into());
208
209        let child_files = child_layer_files();
210
211        let child_builder = ChildLayerFileBuilder::from_files(parent.clone(), &child_files)
212            .await
213            .unwrap();
214        let mut builder = child_builder.into_phase2().await.unwrap();
215        builder.add_triple(1, 2, 3).await.unwrap();
216        builder.add_triple(3, 3, 4).await.unwrap();
217        builder.add_triple(3, 5, 6).await.unwrap();
218        builder.remove_triple(1, 1, 1).await.unwrap();
219        builder.remove_triple(2, 1, 3).await.unwrap();
220        builder.remove_triple(2, 3, 6).await.unwrap();
221        builder.remove_triple(4, 3, 6).await.unwrap();
222        builder.finalize().await.unwrap();
223
224        ChildLayer::load_from_files([5, 4, 3, 2, 1], parent, &child_files)
225            .await
226            .unwrap()
227            .into()
228    }
229
230    #[tokio::test]
231    async fn child_triple_addition_iterator() {
232        let layer = child_layer().await;
233
234        let triples: Vec<_> = layer.internal_triple_additions_p(3).collect();
235
236        let expected = vec![IdTriple::new(3, 3, 4)];
237
238        assert_eq!(expected, triples);
239    }
240
241    #[tokio::test]
242    async fn child_triple_removal_iterator() {
243        let layer = child_layer().await;
244
245        let triples: Vec<_> = layer.internal_triple_removals_p(3).collect();
246
247        let expected = vec![IdTriple::new(2, 3, 6), IdTriple::new(4, 3, 6)];
248
249        assert_eq!(expected, triples);
250    }
251
252    use crate::storage::memory::*;
253    use crate::storage::LayerStore;
254    #[tokio::test]
255    async fn combined_iterator_for_predicate() {
256        let store = MemoryLayerStore::new();
257        let mut builder = store.create_base_layer().await.unwrap();
258        let base_name = builder.name();
259
260        builder.add_value_triple(ValueTriple::new_string_value("cow", "says", "moo"));
261        builder.add_value_triple(ValueTriple::new_string_value("duck", "says", "quack"));
262        builder.add_value_triple(ValueTriple::new_node("cow", "likes", "duck"));
263        builder.add_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
264        builder.commit_boxed().await.unwrap();
265
266        builder = store.create_child_layer(base_name).await.unwrap();
267        let child1_name = builder.name();
268
269        builder.add_value_triple(ValueTriple::new_string_value("horse", "says", "neigh"));
270        builder.add_value_triple(ValueTriple::new_node("horse", "likes", "horse"));
271        builder.commit_boxed().await.unwrap();
272
273        builder = store.create_child_layer(child1_name).await.unwrap();
274        let child2_name = builder.name();
275
276        builder.remove_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
277        builder.add_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
278        builder.commit_boxed().await.unwrap();
279
280        builder = store.create_child_layer(child2_name).await.unwrap();
281        let child3_name = builder.name();
282
283        builder.remove_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
284        builder.add_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
285        builder.commit_boxed().await.unwrap();
286
287        builder = store.create_child_layer(child3_name).await.unwrap();
288        let child4_name = builder.name();
289
290        builder.remove_value_triple(ValueTriple::new_node("duck", "hates", "cow"));
291        builder.add_value_triple(ValueTriple::new_node("duck", "likes", "cow"));
292        builder.commit_boxed().await.unwrap();
293
294        let layer = store.get_layer(child4_name).await.unwrap().unwrap();
295
296        let predicate_id = layer.predicate_id("likes").unwrap();
297        let triples: Vec<_> = layer
298            .triples_p(predicate_id)
299            .map(|t| layer.id_triple_to_string(&t).unwrap())
300            .collect();
301
302        let expected = vec![
303            ValueTriple::new_node("cow", "likes", "duck"),
304            ValueTriple::new_node("duck", "likes", "cow"),
305            ValueTriple::new_node("horse", "likes", "horse"),
306        ];
307
308        assert_eq!(expected, triples);
309    }
310
311    #[tokio::test]
312    async fn one_subject_two_objects() {
313        let store = MemoryLayerStore::new();
314        let mut builder = store.create_base_layer().await.unwrap();
315        let base_name = builder.name();
316
317        builder.add_value_triple(ValueTriple::new_node("cow", "says", "moo"));
318        builder.add_value_triple(ValueTriple::new_node("cow", "says", "quack"));
319        builder.commit_boxed().await.unwrap();
320
321        let layer = store.get_layer(base_name).await.unwrap().unwrap();
322        let predicate_id = layer.predicate_id("says").unwrap();
323        let triples: Vec<_> = layer
324            .triples_p(predicate_id)
325            .map(|t| layer.id_triple_to_string(&t).unwrap())
326            .collect();
327
328        let expected = vec![
329            ValueTriple::new_node("cow", "says", "moo"),
330            ValueTriple::new_node("cow", "says", "quack"),
331        ];
332
333        assert_eq!(expected, triples);
334    }
335}