terminus_store/layer/internal/
predicate_iterator.rs1use 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 if !self.next_pos() {
75 return None;
77 }
78 }
79 let result = self.subject_iterator.next();
80
81 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 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 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}