Skip to main content

terminus_store/storage/
layer.rs

1use super::cache::*;
2use super::consts::FILENAMES;
3use super::delta::*;
4use super::file::*;
5use super::pack::Packable;
6use crate::layer::base_merge::merge_base_layers;
7use crate::layer::builder::DictionarySetFileBuilder;
8use crate::layer::BaseLayerFileBuilder;
9use crate::layer::ChildLayerFileBuilderPhase2;
10use crate::layer::TripleChange;
11use crate::layer::{
12    layer_triple_exists, BaseLayer, ChildLayer, IdMap, IdTriple, InternalLayer,
13    InternalLayerTripleObjectIterator, InternalLayerTriplePredicateIterator,
14    InternalLayerTripleSubjectIterator, InternalTripleStackIterator, LayerBuilder,
15    OptInternalLayerTriplePredicateIterator, OptInternalLayerTripleSubjectIterator, RollupLayer,
16    SimpleLayerBuilder,
17};
18use crate::Layer;
19use tdb_succinct::bitarray::bitarray_len_from_file;
20use tdb_succinct::dict_file_get_count;
21use tdb_succinct::logarray::logarray_file_get_length_and_width;
22use tdb_succinct::StringDict;
23use tdb_succinct::TypedDict;
24use tdb_succinct::{util, AdjacencyList, BitIndex, LogArray, MonotonicLogArray, WaveletTree};
25
26use bitvec::prelude::*;
27use std::convert::TryInto;
28use std::io;
29use std::path::Path;
30use std::sync::Arc;
31
32use tokio::io::{AsyncReadExt, AsyncWriteExt};
33
34use async_trait::async_trait;
35
36macro_rules! walk_backwards_from_disk {
37    ($store:ident, $name:ident, $current:ident, $body:block) => {
38        let mut $current = $name;
39        loop {
40            $body
41
42            if let Some(parent) = $store.get_layer_parent_name($current).await? {
43                $current = parent;
44            }
45            else {
46                break;
47            }
48        }
49    }
50}
51
52macro_rules! walk_backwards_from_disk_upto {
53    ($store:ident, $name:ident, $upto:ident, $current:ident, $body:block) => {
54        let mut $current = $name;
55        loop {
56            if $current == $upto {
57                break Ok(());
58            }
59
60            $body
61
62            if let Some(parent) = $store.get_layer_parent_name($current).await? {
63                $current = parent;
64            }
65            else {
66                break Err(io::Error::new(io::ErrorKind::NotFound, "expected a parent layer, but it was not found"));
67            }
68        }?
69    }
70}
71
72#[async_trait]
73pub trait LayerStore: 'static + Packable + Send + Sync {
74    async fn layers(&self) -> io::Result<Vec<[u32; 5]>>;
75    async fn get_layer_with_cache(
76        &self,
77        name: [u32; 5],
78        cache: Arc<dyn LayerCache>,
79    ) -> io::Result<Option<Arc<InternalLayer>>>;
80    async fn get_layer(&self, name: [u32; 5]) -> io::Result<Option<Arc<InternalLayer>>> {
81        self.get_layer_with_cache(name, NOCACHE.clone()).await
82    }
83
84    async fn finalize_layer(&self, _name: [u32; 5]) -> io::Result<()> {
85        Ok(())
86    }
87
88    async fn get_layer_parent_name(&self, name: [u32; 5]) -> io::Result<Option<[u32; 5]>>;
89
90    async fn get_node_dictionary(&self, name: [u32; 5]) -> io::Result<Option<StringDict>>;
91
92    async fn get_predicate_dictionary(&self, name: [u32; 5]) -> io::Result<Option<StringDict>>;
93
94    async fn get_value_dictionary(&self, name: [u32; 5]) -> io::Result<Option<TypedDict>>;
95
96    async fn get_node_count(&self, name: [u32; 5]) -> io::Result<Option<u64>>;
97
98    async fn get_predicate_count(&self, name: [u32; 5]) -> io::Result<Option<u64>>;
99
100    async fn get_value_count(&self, name: [u32; 5]) -> io::Result<Option<u64>>;
101
102    async fn get_node_value_idmap(&self, name: [u32; 5]) -> io::Result<Option<IdMap>>;
103
104    async fn get_predicate_idmap(&self, name: [u32; 5]) -> io::Result<Option<IdMap>>;
105
106    async fn create_base_layer(&self) -> io::Result<Box<dyn LayerBuilder>>;
107    async fn create_child_layer_with_cache(
108        &self,
109        parent: [u32; 5],
110        cache: Arc<dyn LayerCache>,
111    ) -> io::Result<Box<dyn LayerBuilder>>;
112    async fn create_child_layer(&self, parent: [u32; 5]) -> io::Result<Box<dyn LayerBuilder>> {
113        self.create_child_layer_with_cache(parent, NOCACHE.clone())
114            .await
115    }
116
117    async fn perform_rollup(&self, layer: Arc<InternalLayer>) -> io::Result<[u32; 5]>;
118    async fn perform_rollup_upto_with_cache(
119        &self,
120        layer: Arc<InternalLayer>,
121        upto: [u32; 5],
122        cache: Arc<dyn LayerCache>,
123    ) -> io::Result<[u32; 5]>;
124    async fn perform_rollup_upto(
125        &self,
126        layer: Arc<InternalLayer>,
127        upto: [u32; 5],
128    ) -> io::Result<[u32; 5]> {
129        self.perform_rollup_upto_with_cache(layer, upto, NOCACHE.clone())
130            .await
131    }
132    async fn perform_imprecise_rollup_upto_with_cache(
133        &self,
134        layer: Arc<InternalLayer>,
135        upto: [u32; 5],
136        cache: Arc<dyn LayerCache>,
137    ) -> io::Result<[u32; 5]>;
138    async fn perform_imprecise_rollup_upto(
139        &self,
140        layer: Arc<InternalLayer>,
141        upto: [u32; 5],
142    ) -> io::Result<[u32; 5]> {
143        self.perform_rollup_upto_with_cache(layer, upto, NOCACHE.clone())
144            .await
145    }
146    async fn register_rollup(&self, layer: [u32; 5], rollup: [u32; 5]) -> io::Result<()>;
147
148    /// Create a new rollup layer which rolls up all triples in the given layer, as well as all its ancestors.
149    ///
150    /// It is a good idea to keep layer stacks small, meaning, to only
151    /// have a handful of ancestors for a layer. The more layers there
152    /// are, the longer queries take. Rollup is one approach of
153    /// accomplishing this. Squash is another. Rollup is the better
154    /// option if you need to retain history.
155    async fn rollup(self: Arc<Self>, layer: Arc<InternalLayer>) -> io::Result<[u32; 5]> {
156        let name = layer.name();
157        let rollup = self.perform_rollup(layer).await?;
158        self.register_rollup(name, rollup).await?;
159
160        Ok(rollup)
161    }
162
163    async fn rollup_upto_with_cache(
164        &self,
165        layer: Arc<InternalLayer>,
166        upto: [u32; 5],
167        cache: Arc<dyn LayerCache>,
168    ) -> io::Result<[u32; 5]> {
169        let name = layer.name();
170        let rollup = self
171            .perform_rollup_upto_with_cache(layer, upto, cache)
172            .await?;
173        self.register_rollup(name, rollup).await?;
174
175        Ok(rollup)
176    }
177
178    /// Create a new rollup layer which rolls up all triples in the given layer, as well as all ancestors up to (but not including) the given ancestor.
179    ///
180    /// It is a good idea to keep layer stacks small, meaning, to only
181    /// have a handful of ancestors for a layer. The more layers there
182    /// are, the longer queries take. Rollup is one approach of
183    /// accomplishing this. Squash is another. Rollup is the better
184    /// option if you need to retain history.
185    async fn rollup_upto(&self, layer: Arc<InternalLayer>, upto: [u32; 5]) -> io::Result<[u32; 5]> {
186        self.rollup_upto_with_cache(layer, upto, NOCACHE.clone())
187            .await
188    }
189
190    async fn imprecise_rollup_upto_with_cache(
191        &self,
192        layer: Arc<InternalLayer>,
193        upto: [u32; 5],
194        cache: Arc<dyn LayerCache>,
195    ) -> io::Result<[u32; 5]> {
196        let name = layer.name();
197        let rollup = self
198            .perform_imprecise_rollup_upto_with_cache(layer, upto, cache)
199            .await?;
200        self.register_rollup(name, rollup).await?;
201
202        Ok(rollup)
203    }
204
205    /// Create a new rollup layer which rolls up all triples in the given layer, as well as all ancestors up to (but not including) the given ancestor.
206    ///
207    /// It is a good idea to keep layer stacks small, meaning, to only
208    /// have a handful of ancestors for a layer. The more layers there
209    /// are, the longer queries take. Rollup is one approach of
210    /// accomplishing this. Squash is another. Rollup is the better
211    /// option if you need to retain history.
212    async fn imprecise_rollup_upto(
213        &self,
214        layer: Arc<InternalLayer>,
215        upto: [u32; 5],
216    ) -> io::Result<[u32; 5]> {
217        self.imprecise_rollup_upto_with_cache(layer, upto, NOCACHE.clone())
218            .await
219    }
220
221    async fn squash(&self, layer: Arc<InternalLayer>) -> io::Result<[u32; 5]>;
222    async fn squash_upto(&self, layer: Arc<InternalLayer>, upto: [u32; 5]) -> io::Result<[u32; 5]>;
223
224    async fn merge_base_layer(&self, layers: &[[u32; 5]], temp_dir: &Path) -> io::Result<[u32; 5]>;
225
226    async fn layer_is_ancestor_of(
227        &self,
228        descendant: [u32; 5],
229        ancestor: [u32; 5],
230    ) -> io::Result<bool>;
231
232    async fn triple_addition_exists(
233        &self,
234        layer: [u32; 5],
235        subject: u64,
236        predicate: u64,
237        object: u64,
238    ) -> io::Result<bool>;
239
240    async fn triple_removal_exists(
241        &self,
242        layer: [u32; 5],
243        subject: u64,
244        predicate: u64,
245        object: u64,
246    ) -> io::Result<bool>;
247
248    async fn triple_additions(
249        &self,
250        layer: [u32; 5],
251    ) -> io::Result<OptInternalLayerTripleSubjectIterator>;
252
253    async fn triple_removals(
254        &self,
255        layer: [u32; 5],
256    ) -> io::Result<OptInternalLayerTripleSubjectIterator>;
257
258    async fn triple_additions_s(
259        &self,
260        layer: [u32; 5],
261        subject: u64,
262    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
263
264    async fn triple_removals_s(
265        &self,
266        layer: [u32; 5],
267        subject: u64,
268    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
269
270    async fn triple_additions_sp(
271        &self,
272        layer: [u32; 5],
273        subject: u64,
274        predicate: u64,
275    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
276
277    async fn triple_removals_sp(
278        &self,
279        layer: [u32; 5],
280        subject: u64,
281        predicate: u64,
282    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
283
284    async fn triple_additions_p(
285        &self,
286        layer: [u32; 5],
287        predicate: u64,
288    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
289
290    async fn triple_removals_o(
291        &self,
292        layer: [u32; 5],
293        object: u64,
294    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
295
296    async fn triple_additions_o(
297        &self,
298        layer: [u32; 5],
299        object: u64,
300    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
301
302    async fn triple_removals_p(
303        &self,
304        layer: [u32; 5],
305        predicate: u64,
306    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>>;
307
308    async fn triple_layer_addition_count(&self, layer: [u32; 5]) -> io::Result<usize>;
309
310    async fn triple_layer_removal_count(&self, layer: [u32; 5]) -> io::Result<usize>;
311
312    async fn retrieve_layer_stack_names(&self, name: [u32; 5]) -> io::Result<Vec<[u32; 5]>>;
313
314    async fn retrieve_layer_stack_names_upto(
315        &self,
316        name: [u32; 5],
317        upto: [u32; 5],
318    ) -> io::Result<Vec<[u32; 5]>>;
319
320    async fn layer_changes(&self, name: [u32; 5]) -> io::Result<InternalTripleStackIterator> {
321        let mut positives = Vec::new();
322        let mut negatives = Vec::new();
323        walk_backwards_from_disk!(self, name, current, {
324            positives.push(self.triple_additions(name).await?);
325            negatives.push(self.triple_removals(name).await?);
326        });
327        positives.reverse();
328        negatives.reverse();
329
330        Ok(InternalTripleStackIterator::from_parts(
331            positives, negatives,
332        ))
333    }
334
335    async fn layer_changes_upto(
336        &self,
337        name: [u32; 5],
338        upto: [u32; 5],
339    ) -> io::Result<InternalTripleStackIterator> {
340        let mut positives = Vec::new();
341        let mut negatives = Vec::new();
342        walk_backwards_from_disk_upto!(self, name, upto, current, {
343            positives.push(self.triple_additions(current).await?);
344            negatives.push(self.triple_removals(current).await?);
345        });
346        positives.reverse();
347        negatives.reverse();
348
349        Ok(InternalTripleStackIterator::from_parts(
350            positives, negatives,
351        ))
352    }
353}
354
355#[async_trait]
356pub trait PersistentLayerStore: 'static + Send + Sync + Clone {
357    type File: FileLoad + FileStore + Clone;
358    async fn directories(&self) -> io::Result<Vec<[u32; 5]>>;
359    async fn create_named_directory(&self, id: [u32; 5]) -> io::Result<[u32; 5]>;
360    async fn create_directory(&self) -> io::Result<[u32; 5]> {
361        let name = rand::random();
362        self.create_named_directory(name).await
363    }
364
365    async fn directory_exists(&self, name: [u32; 5]) -> io::Result<bool>;
366    async fn get_file(&self, directory: [u32; 5], name: &str) -> io::Result<Self::File>;
367    async fn file_exists(&self, directory: [u32; 5], file: &str) -> io::Result<bool>;
368
369    async fn finalize(&self, _directory: [u32; 5]) -> io::Result<()> {
370        Ok(())
371    }
372
373    async fn layer_has_rollup(&self, name: [u32; 5]) -> io::Result<bool> {
374        self.file_exists(name, FILENAMES.rollup).await
375    }
376
377    async fn layer_has_parent(&self, name: [u32; 5]) -> io::Result<bool> {
378        self.file_exists(name, FILENAMES.parent).await
379    }
380
381    async fn layer_parent(&self, name: [u32; 5]) -> io::Result<Option<[u32; 5]>> {
382        if self.directory_exists(name).await? {
383            if self.layer_has_parent(name).await? {
384                let parent = self.read_parent_file(name).await?;
385                Ok(Some(parent))
386            } else {
387                Ok(None)
388            }
389        } else {
390            Err(io::Error::new(
391                io::ErrorKind::NotFound,
392                "parent layer not found",
393            ))
394        }
395    }
396
397    async fn base_layer_files(&self, name: [u32; 5]) -> io::Result<BaseLayerFiles<Self::File>> {
398        let filenames = vec![
399            FILENAMES.node_dictionary_blocks,
400            FILENAMES.node_dictionary_offsets,
401            FILENAMES.predicate_dictionary_blocks,
402            FILENAMES.predicate_dictionary_offsets,
403            FILENAMES.value_dictionary_types_present,
404            FILENAMES.value_dictionary_type_offsets,
405            FILENAMES.value_dictionary_offsets,
406            FILENAMES.value_dictionary_blocks,
407            FILENAMES.node_value_idmap_bits,
408            FILENAMES.node_value_idmap_bit_index_blocks,
409            FILENAMES.node_value_idmap_bit_index_sblocks,
410            FILENAMES.predicate_idmap_bits,
411            FILENAMES.predicate_idmap_bit_index_blocks,
412            FILENAMES.predicate_idmap_bit_index_sblocks,
413            FILENAMES.base_subjects,
414            FILENAMES.base_objects,
415            FILENAMES.base_s_p_adjacency_list_bits,
416            FILENAMES.base_s_p_adjacency_list_bit_index_blocks,
417            FILENAMES.base_s_p_adjacency_list_bit_index_sblocks,
418            FILENAMES.base_s_p_adjacency_list_nums,
419            FILENAMES.base_sp_o_adjacency_list_bits,
420            FILENAMES.base_sp_o_adjacency_list_bit_index_blocks,
421            FILENAMES.base_sp_o_adjacency_list_bit_index_sblocks,
422            FILENAMES.base_sp_o_adjacency_list_nums,
423            FILENAMES.base_o_ps_adjacency_list_bits,
424            FILENAMES.base_o_ps_adjacency_list_bit_index_blocks,
425            FILENAMES.base_o_ps_adjacency_list_bit_index_sblocks,
426            FILENAMES.base_o_ps_adjacency_list_nums,
427            FILENAMES.base_predicate_wavelet_tree_bits,
428            FILENAMES.base_predicate_wavelet_tree_bit_index_blocks,
429            FILENAMES.base_predicate_wavelet_tree_bit_index_sblocks,
430        ];
431
432        let mut files = Vec::with_capacity(filenames.len());
433
434        for filename in filenames {
435            files.push(self.get_file(name, filename).await?);
436        }
437
438        Ok(BaseLayerFiles {
439            node_dictionary_files: DictionaryFiles {
440                blocks_file: files[0].clone(),
441                offsets_file: files[1].clone(),
442            },
443            predicate_dictionary_files: DictionaryFiles {
444                blocks_file: files[2].clone(),
445                offsets_file: files[3].clone(),
446            },
447            value_dictionary_files: TypedDictionaryFiles {
448                types_present_file: files[4].clone(),
449                type_offsets_file: files[5].clone(),
450                offsets_file: files[6].clone(),
451                blocks_file: files[7].clone(),
452            },
453
454            id_map_files: IdMapFiles {
455                node_value_idmap_files: BitIndexFiles {
456                    bits_file: files[8].clone(),
457                    blocks_file: files[9].clone(),
458                    sblocks_file: files[10].clone(),
459                },
460                predicate_idmap_files: BitIndexFiles {
461                    bits_file: files[11].clone(),
462                    blocks_file: files[12].clone(),
463                    sblocks_file: files[13].clone(),
464                },
465            },
466
467            subjects_file: files[14].clone(),
468            objects_file: files[15].clone(),
469
470            s_p_adjacency_list_files: AdjacencyListFiles {
471                bitindex_files: BitIndexFiles {
472                    bits_file: files[16].clone(),
473                    blocks_file: files[17].clone(),
474                    sblocks_file: files[18].clone(),
475                },
476                nums_file: files[19].clone(),
477            },
478            sp_o_adjacency_list_files: AdjacencyListFiles {
479                bitindex_files: BitIndexFiles {
480                    bits_file: files[20].clone(),
481                    blocks_file: files[21].clone(),
482                    sblocks_file: files[22].clone(),
483                },
484                nums_file: files[23].clone(),
485            },
486            o_ps_adjacency_list_files: AdjacencyListFiles {
487                bitindex_files: BitIndexFiles {
488                    bits_file: files[24].clone(),
489                    blocks_file: files[25].clone(),
490                    sblocks_file: files[26].clone(),
491                },
492                nums_file: files[27].clone(),
493            },
494            predicate_wavelet_tree_files: BitIndexFiles {
495                bits_file: files[28].clone(),
496                blocks_file: files[29].clone(),
497                sblocks_file: files[30].clone(),
498            },
499        })
500    }
501
502    async fn child_layer_files(&self, name: [u32; 5]) -> io::Result<ChildLayerFiles<Self::File>> {
503        let filenames = vec![
504            FILENAMES.node_dictionary_blocks,
505            FILENAMES.node_dictionary_offsets,
506            FILENAMES.predicate_dictionary_blocks,
507            FILENAMES.predicate_dictionary_offsets,
508            FILENAMES.value_dictionary_types_present,
509            FILENAMES.value_dictionary_type_offsets,
510            FILENAMES.value_dictionary_offsets,
511            FILENAMES.value_dictionary_blocks,
512            FILENAMES.node_value_idmap_bits,
513            FILENAMES.node_value_idmap_bit_index_blocks,
514            FILENAMES.node_value_idmap_bit_index_sblocks,
515            FILENAMES.predicate_idmap_bits,
516            FILENAMES.predicate_idmap_bit_index_blocks,
517            FILENAMES.predicate_idmap_bit_index_sblocks,
518            FILENAMES.pos_subjects,
519            FILENAMES.pos_objects,
520            FILENAMES.neg_subjects,
521            FILENAMES.neg_objects,
522            FILENAMES.pos_s_p_adjacency_list_bits,
523            FILENAMES.pos_s_p_adjacency_list_bit_index_blocks,
524            FILENAMES.pos_s_p_adjacency_list_bit_index_sblocks,
525            FILENAMES.pos_s_p_adjacency_list_nums,
526            FILENAMES.pos_sp_o_adjacency_list_bits,
527            FILENAMES.pos_sp_o_adjacency_list_bit_index_blocks,
528            FILENAMES.pos_sp_o_adjacency_list_bit_index_sblocks,
529            FILENAMES.pos_sp_o_adjacency_list_nums,
530            FILENAMES.pos_o_ps_adjacency_list_bits,
531            FILENAMES.pos_o_ps_adjacency_list_bit_index_blocks,
532            FILENAMES.pos_o_ps_adjacency_list_bit_index_sblocks,
533            FILENAMES.pos_o_ps_adjacency_list_nums,
534            FILENAMES.neg_s_p_adjacency_list_bits,
535            FILENAMES.neg_s_p_adjacency_list_bit_index_blocks,
536            FILENAMES.neg_s_p_adjacency_list_bit_index_sblocks,
537            FILENAMES.neg_s_p_adjacency_list_nums,
538            FILENAMES.neg_sp_o_adjacency_list_bits,
539            FILENAMES.neg_sp_o_adjacency_list_bit_index_blocks,
540            FILENAMES.neg_sp_o_adjacency_list_bit_index_sblocks,
541            FILENAMES.neg_sp_o_adjacency_list_nums,
542            FILENAMES.neg_o_ps_adjacency_list_bits,
543            FILENAMES.neg_o_ps_adjacency_list_bit_index_blocks,
544            FILENAMES.neg_o_ps_adjacency_list_bit_index_sblocks,
545            FILENAMES.neg_o_ps_adjacency_list_nums,
546            FILENAMES.pos_predicate_wavelet_tree_bits,
547            FILENAMES.pos_predicate_wavelet_tree_bit_index_blocks,
548            FILENAMES.pos_predicate_wavelet_tree_bit_index_sblocks,
549            FILENAMES.neg_predicate_wavelet_tree_bits,
550            FILENAMES.neg_predicate_wavelet_tree_bit_index_blocks,
551            FILENAMES.neg_predicate_wavelet_tree_bit_index_sblocks,
552        ];
553
554        let mut files = Vec::with_capacity(filenames.len());
555        for filename in filenames {
556            files.push(self.get_file(name, filename).await?);
557        }
558
559        Ok(ChildLayerFiles {
560            node_dictionary_files: DictionaryFiles {
561                blocks_file: files[0].clone(),
562                offsets_file: files[1].clone(),
563            },
564            predicate_dictionary_files: DictionaryFiles {
565                blocks_file: files[2].clone(),
566                offsets_file: files[3].clone(),
567            },
568            value_dictionary_files: TypedDictionaryFiles {
569                types_present_file: files[4].clone(),
570                type_offsets_file: files[5].clone(),
571                offsets_file: files[6].clone(),
572                blocks_file: files[7].clone(),
573            },
574
575            id_map_files: IdMapFiles {
576                node_value_idmap_files: BitIndexFiles {
577                    bits_file: files[8].clone(),
578                    blocks_file: files[9].clone(),
579                    sblocks_file: files[10].clone(),
580                },
581                predicate_idmap_files: BitIndexFiles {
582                    bits_file: files[11].clone(),
583                    blocks_file: files[12].clone(),
584                    sblocks_file: files[13].clone(),
585                },
586            },
587
588            pos_subjects_file: files[14].clone(),
589            pos_objects_file: files[15].clone(),
590            neg_subjects_file: files[16].clone(),
591            neg_objects_file: files[17].clone(),
592
593            pos_s_p_adjacency_list_files: AdjacencyListFiles {
594                bitindex_files: BitIndexFiles {
595                    bits_file: files[18].clone(),
596                    blocks_file: files[19].clone(),
597                    sblocks_file: files[20].clone(),
598                },
599                nums_file: files[21].clone(),
600            },
601            pos_sp_o_adjacency_list_files: AdjacencyListFiles {
602                bitindex_files: BitIndexFiles {
603                    bits_file: files[22].clone(),
604                    blocks_file: files[23].clone(),
605                    sblocks_file: files[24].clone(),
606                },
607                nums_file: files[25].clone(),
608            },
609            pos_o_ps_adjacency_list_files: AdjacencyListFiles {
610                bitindex_files: BitIndexFiles {
611                    bits_file: files[26].clone(),
612                    blocks_file: files[27].clone(),
613                    sblocks_file: files[28].clone(),
614                },
615                nums_file: files[29].clone(),
616            },
617            neg_s_p_adjacency_list_files: AdjacencyListFiles {
618                bitindex_files: BitIndexFiles {
619                    bits_file: files[30].clone(),
620                    blocks_file: files[31].clone(),
621                    sblocks_file: files[32].clone(),
622                },
623                nums_file: files[33].clone(),
624            },
625            neg_sp_o_adjacency_list_files: AdjacencyListFiles {
626                bitindex_files: BitIndexFiles {
627                    bits_file: files[34].clone(),
628                    blocks_file: files[35].clone(),
629                    sblocks_file: files[36].clone(),
630                },
631                nums_file: files[37].clone(),
632            },
633            neg_o_ps_adjacency_list_files: AdjacencyListFiles {
634                bitindex_files: BitIndexFiles {
635                    bits_file: files[38].clone(),
636                    blocks_file: files[39].clone(),
637                    sblocks_file: files[40].clone(),
638                },
639                nums_file: files[41].clone(),
640            },
641            pos_predicate_wavelet_tree_files: BitIndexFiles {
642                bits_file: files[42].clone(),
643                blocks_file: files[43].clone(),
644                sblocks_file: files[44].clone(),
645            },
646            neg_predicate_wavelet_tree_files: BitIndexFiles {
647                bits_file: files[45].clone(),
648                blocks_file: files[46].clone(),
649                sblocks_file: files[47].clone(),
650            },
651        })
652    }
653
654    async fn write_parent_file(&self, dir_name: [u32; 5], parent_name: [u32; 5]) -> io::Result<()> {
655        let parent_string = name_to_string(parent_name);
656
657        let file = self.get_file(dir_name, FILENAMES.parent).await?;
658        let mut writer = file.open_write().await?;
659
660        writer.write_all(parent_string.as_bytes()).await?;
661        writer.flush().await?;
662        writer.sync_all().await?;
663
664        Ok(())
665    }
666
667    async fn read_parent_file(&self, dir_name: [u32; 5]) -> io::Result<[u32; 5]> {
668        let file = self.get_file(dir_name, FILENAMES.parent).await?;
669        let mut reader = file.open_read().await?;
670
671        let mut buf = [0; 40];
672        reader.read_exact(&mut buf).await?;
673
674        bytes_to_name(&buf)
675    }
676
677    // TODO this should check if the rollup is better than what is there
678    async fn write_rollup_file(&self, dir_name: [u32; 5], rollup_name: [u32; 5]) -> io::Result<()> {
679        let rollup_string = name_to_string(rollup_name);
680
681        let file = self.get_file(dir_name, FILENAMES.rollup).await?;
682        let mut writer = file.open_write().await?;
683
684        let contents = format!("{}\n{}\n", 1, rollup_string);
685        writer.write_all(contents.as_bytes()).await?;
686        writer.flush().await?;
687        writer.sync_all().await?;
688
689        Ok(())
690    }
691
692    async fn read_rollup_file(&self, dir_name: [u32; 5]) -> io::Result<[u32; 5]> {
693        let file = self.get_file(dir_name, FILENAMES.rollup).await?;
694        let mut reader = file.open_read().await?;
695
696        let mut data = Vec::new();
697        reader.read_to_end(&mut data).await?;
698
699        let s = String::from_utf8_lossy(&data);
700        let lines: Vec<&str> = s.lines().collect();
701        if lines.len() != 2 {
702            return Err(io::Error::new(
703                io::ErrorKind::InvalidData,
704                format!(
705                    "expected rollup file to have two lines. contents were ({:?})",
706                    lines
707                ),
708            ));
709        }
710
711        let _version_str = &lines[0];
712        let layer_str = &lines[1];
713
714        string_to_name(layer_str)
715    }
716
717    async fn create_child_layer_files_with_cache(
718        &self,
719        parent: [u32; 5],
720        cache: Arc<dyn LayerCache>,
721    ) -> io::Result<([u32; 5], Arc<InternalLayer>, ChildLayerFiles<Self::File>)> {
722        let parent_layer = match self.get_layer_with_cache(parent, cache).await? {
723            None => {
724                return Err(io::Error::new(
725                    io::ErrorKind::NotFound,
726                    "parent layer not found",
727                ))
728            }
729            Some(parent_layer) => Ok::<_, io::Error>(parent_layer),
730        }?;
731
732        let layer_dir = self.create_directory().await?;
733        self.write_parent_file(layer_dir, parent).await?;
734        let child_layer_files = self.child_layer_files(layer_dir).await?;
735
736        Ok((layer_dir, parent_layer, child_layer_files))
737    }
738
739    async fn node_dictionary_files(
740        &self,
741        layer: [u32; 5],
742    ) -> io::Result<DictionaryFiles<Self::File>> {
743        // does layer exist?
744        if self.directory_exists(layer).await? {
745            let offsets_file = self
746                .get_file(layer, FILENAMES.node_dictionary_offsets)
747                .await?;
748            let blocks_file = self
749                .get_file(layer, FILENAMES.node_dictionary_blocks)
750                .await?;
751
752            Ok(DictionaryFiles {
753                blocks_file,
754                offsets_file,
755            })
756        } else {
757            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
758        }
759    }
760
761    async fn predicate_dictionary_files(
762        &self,
763        layer: [u32; 5],
764    ) -> io::Result<DictionaryFiles<Self::File>> {
765        // does layer exist?
766        if self.directory_exists(layer).await? {
767            let offsets_file = self
768                .get_file(layer, FILENAMES.predicate_dictionary_offsets)
769                .await?;
770            let blocks_file = self
771                .get_file(layer, FILENAMES.predicate_dictionary_blocks)
772                .await?;
773
774            Ok(DictionaryFiles {
775                blocks_file,
776                offsets_file,
777            })
778        } else {
779            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
780        }
781    }
782
783    async fn value_dictionary_files(
784        &self,
785        layer: [u32; 5],
786    ) -> io::Result<TypedDictionaryFiles<Self::File>> {
787        // does layer exist?
788        if self.directory_exists(layer).await? {
789            let types_present_file = self
790                .get_file(layer, FILENAMES.value_dictionary_types_present)
791                .await?;
792            let type_offsets_file = self
793                .get_file(layer, FILENAMES.value_dictionary_type_offsets)
794                .await?;
795            let offsets_file = self
796                .get_file(layer, FILENAMES.value_dictionary_offsets)
797                .await?;
798            let blocks_file = self
799                .get_file(layer, FILENAMES.value_dictionary_blocks)
800                .await?;
801
802            Ok(TypedDictionaryFiles {
803                types_present_file,
804                type_offsets_file,
805                blocks_file,
806                offsets_file,
807            })
808        } else {
809            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
810        }
811    }
812
813    async fn node_value_idmap_files(
814        &self,
815        layer: [u32; 5],
816    ) -> io::Result<BitIndexFiles<Self::File>> {
817        // does layer exist?
818        if self.directory_exists(layer).await? {
819            let bits_file = self
820                .get_file(layer, FILENAMES.node_value_idmap_bits)
821                .await?;
822            let blocks_file = self
823                .get_file(layer, FILENAMES.node_value_idmap_bit_index_blocks)
824                .await?;
825            let sblocks_file = self
826                .get_file(layer, FILENAMES.node_value_idmap_bit_index_sblocks)
827                .await?;
828
829            Ok(BitIndexFiles {
830                bits_file,
831                blocks_file,
832                sblocks_file,
833            })
834        } else {
835            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
836        }
837    }
838
839    async fn predicate_idmap_files(
840        &self,
841        layer: [u32; 5],
842    ) -> io::Result<BitIndexFiles<Self::File>> {
843        // does layer exist?
844        if self.directory_exists(layer).await? {
845            let bits_file = self.get_file(layer, FILENAMES.predicate_idmap_bits).await?;
846            let blocks_file = self
847                .get_file(layer, FILENAMES.predicate_idmap_bit_index_blocks)
848                .await?;
849            let sblocks_file = self
850                .get_file(layer, FILENAMES.predicate_idmap_bit_index_sblocks)
851                .await?;
852
853            Ok(BitIndexFiles {
854                bits_file,
855                blocks_file,
856                sblocks_file,
857            })
858        } else {
859            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
860        }
861    }
862
863    async fn triple_addition_files(
864        &self,
865        layer: [u32; 5],
866    ) -> io::Result<(
867        Self::File,
868        AdjacencyListFiles<Self::File>,
869        AdjacencyListFiles<Self::File>,
870    )> {
871        // does layer exist?
872        if self.directory_exists(layer).await? {
873            let (
874                s_p_aj_nums_file,
875                s_p_aj_bits_file,
876                s_p_aj_bit_index_blocks_file,
877                s_p_aj_bit_index_sblocks_file,
878                subjects_file,
879            );
880            let (
881                sp_o_aj_nums_file,
882                sp_o_aj_bits_file,
883                sp_o_aj_bit_index_blocks_file,
884                sp_o_aj_bit_index_sblocks_file,
885            );
886            if self.layer_has_parent(layer).await? {
887                // this is a child layer
888                s_p_aj_nums_file = self
889                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_nums)
890                    .await?;
891                s_p_aj_bits_file = self
892                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bits)
893                    .await?;
894                s_p_aj_bit_index_blocks_file = self
895                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bit_index_blocks)
896                    .await?;
897                s_p_aj_bit_index_sblocks_file = self
898                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bit_index_sblocks)
899                    .await?;
900
901                sp_o_aj_nums_file = self
902                    .get_file(layer, FILENAMES.pos_sp_o_adjacency_list_nums)
903                    .await?;
904                sp_o_aj_bits_file = self
905                    .get_file(layer, FILENAMES.pos_sp_o_adjacency_list_bits)
906                    .await?;
907                sp_o_aj_bit_index_blocks_file = self
908                    .get_file(layer, FILENAMES.pos_sp_o_adjacency_list_bit_index_blocks)
909                    .await?;
910                sp_o_aj_bit_index_sblocks_file = self
911                    .get_file(layer, FILENAMES.pos_sp_o_adjacency_list_bit_index_sblocks)
912                    .await?;
913
914                subjects_file = self.get_file(layer, FILENAMES.pos_subjects).await?;
915            } else {
916                // this is a base layer
917                s_p_aj_nums_file = self
918                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_nums)
919                    .await?;
920                s_p_aj_bits_file = self
921                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bits)
922                    .await?;
923                s_p_aj_bit_index_blocks_file = self
924                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bit_index_blocks)
925                    .await?;
926                s_p_aj_bit_index_sblocks_file = self
927                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bit_index_sblocks)
928                    .await?;
929
930                sp_o_aj_nums_file = self
931                    .get_file(layer, FILENAMES.base_sp_o_adjacency_list_nums)
932                    .await?;
933                sp_o_aj_bits_file = self
934                    .get_file(layer, FILENAMES.base_sp_o_adjacency_list_bits)
935                    .await?;
936                sp_o_aj_bit_index_blocks_file = self
937                    .get_file(layer, FILENAMES.base_sp_o_adjacency_list_bit_index_blocks)
938                    .await?;
939                sp_o_aj_bit_index_sblocks_file = self
940                    .get_file(layer, FILENAMES.base_sp_o_adjacency_list_bit_index_sblocks)
941                    .await?;
942
943                subjects_file = self.get_file(layer, FILENAMES.base_subjects).await?;
944            }
945
946            let s_p_aj_files = AdjacencyListFiles {
947                bitindex_files: BitIndexFiles {
948                    bits_file: s_p_aj_bits_file,
949                    blocks_file: s_p_aj_bit_index_blocks_file,
950                    sblocks_file: s_p_aj_bit_index_sblocks_file,
951                },
952                nums_file: s_p_aj_nums_file,
953            };
954            let sp_o_aj_files = AdjacencyListFiles {
955                bitindex_files: BitIndexFiles {
956                    bits_file: sp_o_aj_bits_file,
957                    blocks_file: sp_o_aj_bit_index_blocks_file,
958                    sblocks_file: sp_o_aj_bit_index_sblocks_file,
959                },
960                nums_file: sp_o_aj_nums_file,
961            };
962
963            Ok((subjects_file, s_p_aj_files, sp_o_aj_files))
964        } else {
965            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
966        }
967    }
968
969    async fn triple_removal_files(
970        &self,
971        layer: [u32; 5],
972    ) -> io::Result<
973        Option<(
974            Self::File,
975            AdjacencyListFiles<Self::File>,
976            AdjacencyListFiles<Self::File>,
977        )>,
978    > {
979        // does layer exist?
980        if self.directory_exists(layer).await? {
981            if self.layer_has_parent(layer).await? {
982                let (
983                    s_p_aj_nums_file,
984                    s_p_aj_bits_file,
985                    s_p_aj_bit_index_blocks_file,
986                    s_p_aj_bit_index_sblocks_file,
987                    subjects_file,
988                );
989                let (
990                    sp_o_aj_nums_file,
991                    sp_o_aj_bits_file,
992                    sp_o_aj_bit_index_blocks_file,
993                    sp_o_aj_bit_index_sblocks_file,
994                );
995
996                s_p_aj_nums_file = self
997                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_nums)
998                    .await?;
999                s_p_aj_bits_file = self
1000                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bits)
1001                    .await?;
1002                s_p_aj_bit_index_blocks_file = self
1003                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bit_index_blocks)
1004                    .await?;
1005                s_p_aj_bit_index_sblocks_file = self
1006                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bit_index_sblocks)
1007                    .await?;
1008
1009                sp_o_aj_nums_file = self
1010                    .get_file(layer, FILENAMES.neg_sp_o_adjacency_list_nums)
1011                    .await?;
1012                sp_o_aj_bits_file = self
1013                    .get_file(layer, FILENAMES.neg_sp_o_adjacency_list_bits)
1014                    .await?;
1015                sp_o_aj_bit_index_blocks_file = self
1016                    .get_file(layer, FILENAMES.neg_sp_o_adjacency_list_bit_index_blocks)
1017                    .await?;
1018                sp_o_aj_bit_index_sblocks_file = self
1019                    .get_file(layer, FILENAMES.neg_sp_o_adjacency_list_bit_index_sblocks)
1020                    .await?;
1021
1022                subjects_file = self.get_file(layer, FILENAMES.neg_subjects).await?;
1023                let s_p_aj_files = AdjacencyListFiles {
1024                    bitindex_files: BitIndexFiles {
1025                        bits_file: s_p_aj_bits_file,
1026                        blocks_file: s_p_aj_bit_index_blocks_file,
1027                        sblocks_file: s_p_aj_bit_index_sblocks_file,
1028                    },
1029                    nums_file: s_p_aj_nums_file,
1030                };
1031                let sp_o_aj_files = AdjacencyListFiles {
1032                    bitindex_files: BitIndexFiles {
1033                        bits_file: sp_o_aj_bits_file,
1034                        blocks_file: sp_o_aj_bit_index_blocks_file,
1035                        sblocks_file: sp_o_aj_bit_index_sblocks_file,
1036                    },
1037                    nums_file: sp_o_aj_nums_file,
1038                };
1039
1040                Ok(Some((subjects_file, s_p_aj_files, sp_o_aj_files)))
1041            } else {
1042                // base layer, so removal does not exist
1043                Ok(None)
1044            }
1045        } else {
1046            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1047        }
1048    }
1049
1050    async fn predicate_wavelet_addition_files(
1051        &self,
1052        layer: [u32; 5],
1053    ) -> io::Result<BitIndexFiles<Self::File>> {
1054        // does layer exist?
1055        if self.directory_exists(layer).await? {
1056            let (wavelet_bits_file, wavelet_bit_index_blocks_file, wavelet_bit_index_sblocks_file);
1057            if self.layer_has_parent(layer).await? {
1058                // this is a child layer
1059                wavelet_bits_file = self
1060                    .get_file(layer, FILENAMES.pos_predicate_wavelet_tree_bits)
1061                    .await?;
1062                wavelet_bit_index_blocks_file = self
1063                    .get_file(layer, FILENAMES.pos_predicate_wavelet_tree_bit_index_blocks)
1064                    .await?;
1065                wavelet_bit_index_sblocks_file = self
1066                    .get_file(
1067                        layer,
1068                        FILENAMES.pos_predicate_wavelet_tree_bit_index_sblocks,
1069                    )
1070                    .await?;
1071            } else {
1072                // this is a base layer
1073                wavelet_bits_file = self
1074                    .get_file(layer, FILENAMES.base_predicate_wavelet_tree_bits)
1075                    .await?;
1076                wavelet_bit_index_blocks_file = self
1077                    .get_file(
1078                        layer,
1079                        FILENAMES.base_predicate_wavelet_tree_bit_index_blocks,
1080                    )
1081                    .await?;
1082                wavelet_bit_index_sblocks_file = self
1083                    .get_file(
1084                        layer,
1085                        FILENAMES.base_predicate_wavelet_tree_bit_index_sblocks,
1086                    )
1087                    .await?;
1088            }
1089
1090            let bitindex_files = BitIndexFiles {
1091                bits_file: wavelet_bits_file,
1092                blocks_file: wavelet_bit_index_blocks_file,
1093                sblocks_file: wavelet_bit_index_sblocks_file,
1094            };
1095            Ok(bitindex_files)
1096        } else {
1097            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1098        }
1099    }
1100
1101    async fn predicate_wavelet_removal_files(
1102        &self,
1103        layer: [u32; 5],
1104    ) -> io::Result<Option<BitIndexFiles<Self::File>>> {
1105        // does layer exist?
1106        if self.directory_exists(layer).await? {
1107            if self.layer_has_parent(layer).await? {
1108                // this is a child layer
1109                let wavelet_bits_file = self
1110                    .get_file(layer, FILENAMES.neg_predicate_wavelet_tree_bits)
1111                    .await?;
1112                let wavelet_bit_index_blocks_file = self
1113                    .get_file(layer, FILENAMES.neg_predicate_wavelet_tree_bit_index_blocks)
1114                    .await?;
1115                let wavelet_bit_index_sblocks_file = self
1116                    .get_file(
1117                        layer,
1118                        FILENAMES.neg_predicate_wavelet_tree_bit_index_sblocks,
1119                    )
1120                    .await?;
1121                let bitindex_files = BitIndexFiles {
1122                    bits_file: wavelet_bits_file,
1123                    blocks_file: wavelet_bit_index_blocks_file,
1124                    sblocks_file: wavelet_bit_index_sblocks_file,
1125                };
1126                Ok(Some(bitindex_files))
1127            } else {
1128                // this is a base layer
1129                Ok(None)
1130            }
1131        } else {
1132            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1133        }
1134    }
1135
1136    async fn triple_addition_files_by_object(
1137        &self,
1138        layer: [u32; 5],
1139    ) -> io::Result<(
1140        Self::File,
1141        Self::File,
1142        AdjacencyListFiles<Self::File>,
1143        AdjacencyListFiles<Self::File>,
1144    )> {
1145        // does layer exist?
1146        if self.directory_exists(layer).await? {
1147            let (
1148                subjects_file,
1149                objects_file,
1150                o_ps_aj_nums_file,
1151                o_ps_aj_bits_file,
1152                o_ps_aj_bit_index_blocks_file,
1153                o_ps_aj_bit_index_sblocks_file,
1154            );
1155            let (
1156                s_p_aj_nums_file,
1157                s_p_aj_bits_file,
1158                s_p_aj_bit_index_blocks_file,
1159                s_p_aj_bit_index_sblocks_file,
1160            );
1161            if self.layer_has_parent(layer).await? {
1162                // this is a child layer
1163                o_ps_aj_nums_file = self
1164                    .get_file(layer, FILENAMES.pos_o_ps_adjacency_list_nums)
1165                    .await?;
1166                o_ps_aj_bits_file = self
1167                    .get_file(layer, FILENAMES.pos_o_ps_adjacency_list_bits)
1168                    .await?;
1169                o_ps_aj_bit_index_blocks_file = self
1170                    .get_file(layer, FILENAMES.pos_o_ps_adjacency_list_bit_index_blocks)
1171                    .await?;
1172                o_ps_aj_bit_index_sblocks_file = self
1173                    .get_file(layer, FILENAMES.pos_o_ps_adjacency_list_bit_index_sblocks)
1174                    .await?;
1175
1176                s_p_aj_nums_file = self
1177                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_nums)
1178                    .await?;
1179                s_p_aj_bits_file = self
1180                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bits)
1181                    .await?;
1182                s_p_aj_bit_index_blocks_file = self
1183                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bit_index_blocks)
1184                    .await?;
1185                s_p_aj_bit_index_sblocks_file = self
1186                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_bit_index_sblocks)
1187                    .await?;
1188
1189                subjects_file = self.get_file(layer, FILENAMES.pos_subjects).await?;
1190                objects_file = self.get_file(layer, FILENAMES.pos_objects).await?;
1191            } else {
1192                // this is a base layer
1193                o_ps_aj_nums_file = self
1194                    .get_file(layer, FILENAMES.base_o_ps_adjacency_list_nums)
1195                    .await?;
1196                o_ps_aj_bits_file = self
1197                    .get_file(layer, FILENAMES.base_o_ps_adjacency_list_bits)
1198                    .await?;
1199                o_ps_aj_bit_index_blocks_file = self
1200                    .get_file(layer, FILENAMES.base_o_ps_adjacency_list_bit_index_blocks)
1201                    .await?;
1202                o_ps_aj_bit_index_sblocks_file = self
1203                    .get_file(layer, FILENAMES.base_o_ps_adjacency_list_bit_index_sblocks)
1204                    .await?;
1205
1206                s_p_aj_nums_file = self
1207                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_nums)
1208                    .await?;
1209                s_p_aj_bits_file = self
1210                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bits)
1211                    .await?;
1212                s_p_aj_bit_index_blocks_file = self
1213                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bit_index_blocks)
1214                    .await?;
1215                s_p_aj_bit_index_sblocks_file = self
1216                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_bit_index_sblocks)
1217                    .await?;
1218
1219                subjects_file = self.get_file(layer, FILENAMES.base_subjects).await?;
1220                objects_file = self.get_file(layer, FILENAMES.base_objects).await?;
1221            }
1222
1223            let o_ps_aj_files = AdjacencyListFiles {
1224                bitindex_files: BitIndexFiles {
1225                    bits_file: o_ps_aj_bits_file,
1226                    blocks_file: o_ps_aj_bit_index_blocks_file,
1227                    sblocks_file: o_ps_aj_bit_index_sblocks_file,
1228                },
1229                nums_file: o_ps_aj_nums_file,
1230            };
1231            let s_p_aj_files = AdjacencyListFiles {
1232                bitindex_files: BitIndexFiles {
1233                    bits_file: s_p_aj_bits_file,
1234                    blocks_file: s_p_aj_bit_index_blocks_file,
1235                    sblocks_file: s_p_aj_bit_index_sblocks_file,
1236                },
1237                nums_file: s_p_aj_nums_file,
1238            };
1239
1240            Ok((subjects_file, objects_file, o_ps_aj_files, s_p_aj_files))
1241        } else {
1242            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1243        }
1244    }
1245
1246    async fn triple_removal_files_by_object(
1247        &self,
1248        layer: [u32; 5],
1249    ) -> io::Result<
1250        Option<(
1251            Self::File,
1252            Self::File,
1253            AdjacencyListFiles<Self::File>,
1254            AdjacencyListFiles<Self::File>,
1255        )>,
1256    > {
1257        // does layer exist?
1258        if self.directory_exists(layer).await? {
1259            if self.layer_has_parent(layer).await? {
1260                // this is a child layer
1261                let o_ps_aj_nums_file = self
1262                    .get_file(layer, FILENAMES.neg_o_ps_adjacency_list_nums)
1263                    .await?;
1264                let o_ps_aj_bits_file = self
1265                    .get_file(layer, FILENAMES.neg_o_ps_adjacency_list_bits)
1266                    .await?;
1267                let o_ps_aj_bit_index_blocks_file = self
1268                    .get_file(layer, FILENAMES.neg_o_ps_adjacency_list_bit_index_blocks)
1269                    .await?;
1270                let o_ps_aj_bit_index_sblocks_file = self
1271                    .get_file(layer, FILENAMES.neg_o_ps_adjacency_list_bit_index_sblocks)
1272                    .await?;
1273
1274                let s_p_aj_nums_file = self
1275                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_nums)
1276                    .await?;
1277                let s_p_aj_bits_file = self
1278                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bits)
1279                    .await?;
1280                let s_p_aj_bit_index_blocks_file = self
1281                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bit_index_blocks)
1282                    .await?;
1283                let s_p_aj_bit_index_sblocks_file = self
1284                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_bit_index_sblocks)
1285                    .await?;
1286
1287                let subjects_file = self.get_file(layer, FILENAMES.neg_subjects).await?;
1288                let objects_file = self.get_file(layer, FILENAMES.neg_objects).await?;
1289
1290                let o_ps_aj_files = AdjacencyListFiles {
1291                    bitindex_files: BitIndexFiles {
1292                        bits_file: o_ps_aj_bits_file,
1293                        blocks_file: o_ps_aj_bit_index_blocks_file,
1294                        sblocks_file: o_ps_aj_bit_index_sblocks_file,
1295                    },
1296                    nums_file: o_ps_aj_nums_file,
1297                };
1298                let s_p_aj_files = AdjacencyListFiles {
1299                    bitindex_files: BitIndexFiles {
1300                        bits_file: s_p_aj_bits_file,
1301                        blocks_file: s_p_aj_bit_index_blocks_file,
1302                        sblocks_file: s_p_aj_bit_index_sblocks_file,
1303                    },
1304                    nums_file: s_p_aj_nums_file,
1305                };
1306
1307                Ok(Some((
1308                    subjects_file,
1309                    objects_file,
1310                    o_ps_aj_files,
1311                    s_p_aj_files,
1312                )))
1313            } else {
1314                // this is a base layer
1315                Ok(None)
1316            }
1317        } else {
1318            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1319        }
1320    }
1321
1322    async fn triple_layer_addition_count_files(
1323        &self,
1324        layer: [u32; 5],
1325    ) -> io::Result<(Self::File, Self::File, BitIndexFiles<Self::File>)> {
1326        // does layer exist?
1327        if self.directory_exists(layer).await? {
1328            let (s_p_nums_file, sp_o_bits_file);
1329            if self.layer_has_parent(layer).await? {
1330                // this is a child layer
1331                s_p_nums_file = self
1332                    .get_file(layer, FILENAMES.pos_s_p_adjacency_list_nums)
1333                    .await?;
1334                sp_o_bits_file = self
1335                    .get_file(layer, FILENAMES.pos_sp_o_adjacency_list_bits)
1336                    .await?;
1337            } else {
1338                // this is a base layer
1339                s_p_nums_file = self
1340                    .get_file(layer, FILENAMES.base_s_p_adjacency_list_nums)
1341                    .await?;
1342                sp_o_bits_file = self
1343                    .get_file(layer, FILENAMES.base_sp_o_adjacency_list_bits)
1344                    .await?;
1345            }
1346
1347            let predicate_wavelet_files = self.predicate_wavelet_addition_files(layer).await?;
1348            Ok((s_p_nums_file, sp_o_bits_file, predicate_wavelet_files))
1349        } else {
1350            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1351        }
1352    }
1353
1354    async fn triple_layer_removal_count_files(
1355        &self,
1356        layer: [u32; 5],
1357    ) -> io::Result<Option<(Self::File, Self::File, BitIndexFiles<Self::File>)>> {
1358        // does layer exist?
1359        if self.directory_exists(layer).await? {
1360            if self.layer_has_parent(layer).await? {
1361                // this is a child layer
1362                let s_p_nums_file = self
1363                    .get_file(layer, FILENAMES.neg_s_p_adjacency_list_nums)
1364                    .await?;
1365                let sp_o_bits_file = self
1366                    .get_file(layer, FILENAMES.neg_sp_o_adjacency_list_bits)
1367                    .await?;
1368                let predicate_wavelet_files = self
1369                    .predicate_wavelet_removal_files(layer)
1370                    .await?
1371                    .expect("expected wavelet removal files to exist");
1372                Ok(Some((
1373                    s_p_nums_file,
1374                    sp_o_bits_file,
1375                    predicate_wavelet_files,
1376                )))
1377            } else {
1378                // this is a base layer
1379                Ok(None)
1380            }
1381        } else {
1382            Err(io::Error::new(io::ErrorKind::NotFound, "layer not found"))
1383        }
1384    }
1385}
1386
1387pub fn name_to_string(name: [u32; 5]) -> String {
1388    format!(
1389        "{:08x}{:08x}{:08x}{:08x}{:08x}",
1390        name[0], name[1], name[2], name[3], name[4]
1391    )
1392}
1393
1394pub fn string_to_name(string: &str) -> Result<[u32; 5], std::io::Error> {
1395    if string.len() != 40 {
1396        return Err(io::Error::new(
1397            io::ErrorKind::Other,
1398            format!("string not len 40: {}", string),
1399        ));
1400    }
1401    let n1 = u32::from_str_radix(&string[..8], 16)
1402        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
1403    let n2 = u32::from_str_radix(&string[8..16], 16)
1404        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
1405    let n3 = u32::from_str_radix(&string[16..24], 16)
1406        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
1407    let n4 = u32::from_str_radix(&string[24..32], 16)
1408        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
1409    let n5 = u32::from_str_radix(&string[32..40], 16)
1410        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
1411
1412    Ok([n1, n2, n3, n4, n5])
1413}
1414
1415pub fn bytes_to_name(bytes: &[u8]) -> Result<[u32; 5], std::io::Error> {
1416    if bytes.len() != 40 {
1417        Err(io::Error::new(io::ErrorKind::Other, "bytes not len 40"))
1418    } else {
1419        let string = String::from_utf8_lossy(bytes);
1420
1421        string_to_name(&string)
1422    }
1423}
1424
1425#[async_trait]
1426impl<F: 'static + FileLoad + FileStore + Clone, T: 'static + PersistentLayerStore<File = F>>
1427    LayerStore for T
1428{
1429    async fn layers(&self) -> io::Result<Vec<[u32; 5]>> {
1430        self.directories().await
1431    }
1432
1433    async fn get_layer_with_cache(
1434        &self,
1435        name: [u32; 5],
1436        cache: Arc<dyn LayerCache>,
1437    ) -> io::Result<Option<Arc<InternalLayer>>> {
1438        if let Some(layer) = cache.get_layer_from_cache(name) {
1439            return Ok(Some(layer));
1440        }
1441
1442        let mut layers_to_load: Vec<([u32; 5], Option<([u32; 5], Option<[u32; 5]>)>)> =
1443            vec![(name, None)];
1444
1445        if !self.directory_exists(name).await? {
1446            return Ok(None);
1447        }
1448
1449        // find an ancestor in cache
1450        let mut ancestor = None;
1451        loop {
1452            let (original_layer, rollup_option) = *layers_to_load.last().unwrap();
1453            let current_layer;
1454            if let Some((rollup, _original_parent)) = rollup_option {
1455                current_layer = rollup;
1456            } else {
1457                current_layer = original_layer;
1458            }
1459
1460            match cache.get_layer_from_cache(current_layer) {
1461                Some(layer) => {
1462                    // remove found cached layer from ids to retrieve
1463                    layers_to_load.pop().unwrap();
1464
1465                    // if this is a rollup, the behavior has to be slightly different
1466                    if let Some((_, original_parent)) = rollup_option {
1467                        if layer.immediate_parent().is_some() {
1468                            ancestor = Some(Arc::new(
1469                                RollupLayer::from_child_layer(
1470                                    layer,
1471                                    original_layer,
1472                                    original_parent.unwrap(),
1473                                )
1474                                .into(),
1475                            ));
1476                        } else {
1477                            ancestor = Some(Arc::new(
1478                                RollupLayer::from_base_layer(
1479                                    layer,
1480                                    original_layer,
1481                                    original_parent,
1482                                )
1483                                .into(),
1484                            ));
1485                        }
1486                    } else {
1487                        ancestor = Some(layer);
1488                    }
1489                    break;
1490                }
1491                None => {
1492                    if self.layer_has_rollup(current_layer).await? {
1493                        let rollup = self.read_rollup_file(current_layer).await?;
1494                        if rollup == current_layer {
1495                            panic!("infinite rollup loop for layer {:?}", rollup);
1496                        }
1497
1498                        let original_parent;
1499                        if self.layer_has_parent(current_layer).await? {
1500                            original_parent = Some(self.read_parent_file(current_layer).await?);
1501                        } else {
1502                            original_parent = None;
1503                        }
1504
1505                        layers_to_load.pop().unwrap(); // we don't want to load this, we want to load the rollup instead!
1506                        layers_to_load.push((current_layer, Some((rollup, original_parent))));
1507                    } else if self.layer_has_parent(current_layer).await? {
1508                        let parent = self.read_parent_file(current_layer).await?;
1509                        layers_to_load.push((parent, None));
1510                    } else {
1511                        break;
1512                    }
1513                }
1514            }
1515        }
1516
1517        if ancestor.is_none() {
1518            // load the base layer
1519            let (base_id, rollup) = layers_to_load.pop().unwrap();
1520            let layer: Arc<InternalLayer>;
1521            match rollup {
1522                None => {
1523                    let files = self.base_layer_files(base_id).await?;
1524                    let base_layer = BaseLayer::load_from_files(base_id, &files).await?;
1525
1526                    layer = Arc::new(base_layer.into());
1527                }
1528                Some((rollup_id, original_parent_id_option)) => {
1529                    let files = self.base_layer_files(rollup_id).await?;
1530                    let base_layer: Arc<InternalLayer> =
1531                        Arc::new(BaseLayer::load_from_files(rollup_id, &files).await?.into());
1532                    cache.cache_layer(base_layer.clone());
1533
1534                    layer = Arc::new(
1535                        RollupLayer::from_base_layer(
1536                            base_layer,
1537                            base_id,
1538                            original_parent_id_option,
1539                        )
1540                        .into(),
1541                    );
1542                }
1543            }
1544
1545            cache.cache_layer(layer.clone());
1546            ancestor = Some(layer);
1547        }
1548
1549        let mut ancestor = ancestor.unwrap();
1550        layers_to_load.reverse();
1551
1552        for (layer_id, rollup) in layers_to_load {
1553            let layer: Arc<InternalLayer>;
1554            match rollup {
1555                None => {
1556                    let files = self.child_layer_files(layer_id).await?;
1557                    let child_layer =
1558                        ChildLayer::load_from_files(layer_id, ancestor, &files).await?;
1559                    layer = Arc::new(child_layer.into());
1560                }
1561                Some((rollup_id, original_parent_id_option)) => {
1562                    let original_parent_id = original_parent_id_option
1563                        .expect("child rollup layer should always have original parent id");
1564
1565                    let files = self.child_layer_files(rollup_id).await?;
1566                    let child_layer: Arc<InternalLayer> = Arc::new(
1567                        ChildLayer::load_from_files(rollup_id, ancestor, &files)
1568                            .await?
1569                            .into(),
1570                    );
1571                    cache.cache_layer(child_layer.clone());
1572
1573                    layer = Arc::new(
1574                        RollupLayer::from_child_layer(child_layer, layer_id, original_parent_id)
1575                            .into(),
1576                    );
1577                }
1578            }
1579
1580            cache.cache_layer(layer.clone());
1581            ancestor = layer;
1582        }
1583
1584        debug_assert_eq!(name, ancestor.name());
1585
1586        Ok(Some(ancestor))
1587    }
1588
1589    async fn finalize_layer(&self, name: [u32; 5]) -> io::Result<()> {
1590        self.finalize(name).await
1591    }
1592
1593    async fn get_layer_parent_name(&self, name: [u32; 5]) -> io::Result<Option<[u32; 5]>> {
1594        self.layer_parent(name).await
1595    }
1596
1597    async fn get_node_dictionary(&self, name: [u32; 5]) -> io::Result<Option<StringDict>> {
1598        if self.directory_exists(name).await? {
1599            let files = self.node_dictionary_files(name).await?;
1600            let maps = files.map_all().await?;
1601
1602            Ok(Some(StringDict::parse(maps.offsets_map, maps.blocks_map)))
1603        } else {
1604            Ok(None)
1605        }
1606    }
1607
1608    async fn get_predicate_dictionary(&self, name: [u32; 5]) -> io::Result<Option<StringDict>> {
1609        if self.directory_exists(name).await? {
1610            let files = self.predicate_dictionary_files(name).await?;
1611            let maps = files.map_all().await?;
1612
1613            Ok(Some(StringDict::parse(maps.offsets_map, maps.blocks_map)))
1614        } else {
1615            Ok(None)
1616        }
1617    }
1618
1619    async fn get_value_dictionary(&self, name: [u32; 5]) -> io::Result<Option<TypedDict>> {
1620        if self.directory_exists(name).await? {
1621            let files = self.value_dictionary_files(name).await?;
1622            let maps = files.map_all().await?;
1623
1624            Ok(Some(TypedDict::from_parts(
1625                maps.types_present_map,
1626                maps.type_offsets_map,
1627                maps.offsets_map,
1628                maps.blocks_map,
1629            )))
1630        } else {
1631            Ok(None)
1632        }
1633    }
1634
1635    async fn get_node_count(&self, name: [u32; 5]) -> io::Result<Option<u64>> {
1636        if self.directory_exists(name).await? {
1637            let file = self.node_dictionary_files(name).await?.blocks_file;
1638            Ok(Some(dict_file_get_count(file).await?))
1639        } else {
1640            Ok(None)
1641        }
1642    }
1643
1644    async fn get_predicate_count(&self, name: [u32; 5]) -> io::Result<Option<u64>> {
1645        if self.directory_exists(name).await? {
1646            let file = self.predicate_dictionary_files(name).await?.blocks_file;
1647            Ok(Some(dict_file_get_count(file).await?))
1648        } else {
1649            Ok(None)
1650        }
1651    }
1652
1653    async fn get_value_count(&self, name: [u32; 5]) -> io::Result<Option<u64>> {
1654        if self.directory_exists(name).await? {
1655            let file = self.value_dictionary_files(name).await?.blocks_file;
1656            Ok(Some(dict_file_get_count(file).await?))
1657        } else {
1658            Ok(None)
1659        }
1660    }
1661
1662    async fn get_node_value_idmap(&self, name: [u32; 5]) -> io::Result<Option<IdMap>> {
1663        if self.directory_exists(name).await? {
1664            let size = self.get_node_count(name).await?.unwrap()
1665                + self.get_value_count(name).await?.unwrap();
1666            let width = util::calculate_width(size);
1667            let files = self.node_value_idmap_files(name).await?;
1668            let maps = files.map_all_if_exists().await?;
1669
1670            let wtree = maps.map(|m| {
1671                WaveletTree::from_parts(
1672                    BitIndex::from_maps(m.bits_map, m.blocks_map, m.sblocks_map),
1673                    width,
1674                )
1675            });
1676            let idmap = IdMap::from_parts(wtree);
1677
1678            Ok(Some(idmap))
1679        } else {
1680            Ok(None)
1681        }
1682    }
1683
1684    async fn get_predicate_idmap(&self, name: [u32; 5]) -> io::Result<Option<IdMap>> {
1685        if self.directory_exists(name).await? {
1686            let size = self.get_predicate_count(name).await?.unwrap();
1687            let width = util::calculate_width(size);
1688            let files = self.predicate_idmap_files(name).await?;
1689            let maps = files.map_all_if_exists().await?;
1690
1691            let wtree = maps.map(|m| {
1692                WaveletTree::from_parts(
1693                    BitIndex::from_maps(m.bits_map, m.blocks_map, m.sblocks_map),
1694                    width,
1695                )
1696            });
1697            let idmap = IdMap::from_parts(wtree);
1698
1699            Ok(Some(idmap))
1700        } else {
1701            Ok(None)
1702        }
1703    }
1704
1705    async fn create_base_layer(&self) -> io::Result<Box<dyn LayerBuilder>> {
1706        let dir_name = self.create_directory().await?;
1707        let files = self.base_layer_files(dir_name).await?;
1708        Ok(Box::new(SimpleLayerBuilder::new(dir_name, files)) as Box<dyn LayerBuilder>)
1709    }
1710
1711    async fn create_child_layer_with_cache(
1712        &self,
1713        parent: [u32; 5],
1714        cache: Arc<dyn LayerCache>,
1715    ) -> io::Result<Box<dyn LayerBuilder>> {
1716        let (layer_dir, parent_layer, child_layer_files) = self
1717            .create_child_layer_files_with_cache(parent, cache)
1718            .await?;
1719
1720        Ok(Box::new(SimpleLayerBuilder::from_parent(
1721            layer_dir,
1722            parent_layer,
1723            child_layer_files,
1724        )) as Box<dyn LayerBuilder>)
1725    }
1726
1727    async fn perform_rollup(&self, layer: Arc<InternalLayer>) -> io::Result<[u32; 5]> {
1728        if layer.parent_name().is_none() {
1729            // we're already a base layer. there's nothing that can be rolled up.
1730            // returning our own name will inhibit writing a rollup file.
1731            return Ok(layer.name());
1732        }
1733
1734        // check if there's already an equivalent rollup
1735        if self.layer_has_rollup(layer.name()).await? {
1736            let rollup = self.read_rollup_file(layer.name()).await?;
1737            // the rollup is equivalent if it is a base layer
1738            if !self.layer_has_parent(rollup).await? {
1739                // yup, equivalent. let's just return the rollup we know about.
1740                return Ok(rollup);
1741            }
1742        }
1743
1744        let dir_name = self.create_directory().await?;
1745        let files = self.base_layer_files(dir_name).await?;
1746        delta_rollup(&layer, files).await?;
1747        self.finalize(dir_name).await?;
1748
1749        Ok(dir_name)
1750    }
1751
1752    async fn perform_rollup_upto_with_cache(
1753        &self,
1754        layer: Arc<InternalLayer>,
1755        upto: [u32; 5],
1756        cache: Arc<dyn LayerCache>,
1757    ) -> io::Result<[u32; 5]> {
1758        if layer.name() == upto {
1759            // rolling up upto ourselves is pretty pointless. Let's not do that.
1760            return Ok(layer.name());
1761        }
1762
1763        if layer.parent_name() == Some(upto) {
1764            // rolling up to our parent is just going to create a clone of this child layer. Let's not do that.
1765            return Ok(layer.name());
1766        }
1767
1768        // check if there's already an equivalent rollup
1769        if self.layer_has_rollup(layer.name()).await? {
1770            let rollup = self.read_rollup_file(layer.name()).await?;
1771
1772            // get rollup parent. if it is the same as upto, we're requesting something equivalent.
1773            if upto == self.read_parent_file(rollup).await? {
1774                // yup, equivalent. Let's just return the rollup we know about.
1775                return Ok(rollup);
1776            }
1777        }
1778
1779        let (layer_dir, _parent_layer, child_layer_files) = self
1780            .create_child_layer_files_with_cache(upto, cache)
1781            .await?;
1782        delta_rollup_upto(self, &layer, upto, child_layer_files).await?;
1783        self.finalize(layer_dir).await?;
1784        Ok(layer_dir)
1785    }
1786
1787    async fn perform_imprecise_rollup_upto_with_cache(
1788        &self,
1789        layer: Arc<InternalLayer>,
1790        upto: [u32; 5],
1791        cache: Arc<dyn LayerCache>,
1792    ) -> io::Result<[u32; 5]> {
1793        if layer.name() == upto {
1794            // rolling up upto ourselves is pretty pointless. Let's not do that.
1795            return Ok(layer.name());
1796        }
1797
1798        if layer.parent_name() == Some(upto) {
1799            // rolling up to our parent is just going to create a clone of this child layer. Let's not do that.
1800            return Ok(layer.name());
1801        }
1802
1803        // check if there's already an equivalent rollup
1804        if self.layer_has_rollup(layer.name()).await? {
1805            let rollup = self.read_rollup_file(layer.name()).await?;
1806
1807            // get rollup parent. if it is the same as upto, we're requesting something equivalent.
1808            if upto == self.read_parent_file(rollup).await? {
1809                // yup, equivalent. Let's just return the rollup we know about.
1810                return Ok(rollup);
1811            }
1812        }
1813
1814        let (layer_dir, _parent_layer, child_layer_files) = self
1815            .create_child_layer_files_with_cache(upto, cache)
1816            .await?;
1817        imprecise_delta_rollup_upto(self, &layer, upto, child_layer_files).await?;
1818        self.finalize(layer_dir).await?;
1819        Ok(layer_dir)
1820    }
1821
1822    async fn register_rollup(&self, layer: [u32; 5], rollup: [u32; 5]) -> io::Result<()> {
1823        if layer == rollup {
1824            // let's not create a loop
1825            Ok(())
1826        } else {
1827            self.write_rollup_file(layer, rollup).await
1828        }
1829    }
1830
1831    async fn squash(&self, layer: Arc<InternalLayer>) -> io::Result<[u32; 5]> {
1832        // we create a new base layer
1833        // we then build a new set of dictionaries by sorting what we got in all layers
1834        // keep track of how the ids move so we have a mapping
1835        // then iterate over all id triples and map them
1836        // sort the mapped triples, insert them
1837        // finalize
1838
1839        let stack = layer.immediate_layers();
1840        let node_count: usize = stack
1841            .iter()
1842            .map(|l| l.node_dictionary().num_entries())
1843            .sum();
1844        let predicate_count: usize = stack
1845            .iter()
1846            .map(|l| l.predicate_dictionary().num_entries())
1847            .sum();
1848        let value_count: usize = stack
1849            .iter()
1850            .map(|l| l.value_dictionary().num_entries())
1851            .sum();
1852
1853        // figure out what keys are actually used through a prepass over all triples
1854        let mut node_value_existences = bitvec![0;node_count+value_count+1];
1855        let mut predicate_existences = bitvec![0;predicate_count+1];
1856        let mut num_triples = 0;
1857        for triple in layer.triples() {
1858            num_triples += 1;
1859            node_value_existences.set(triple.subject as usize, true);
1860            predicate_existences.set(triple.predicate as usize, true);
1861            node_value_existences.set(triple.object as usize, true);
1862        }
1863
1864        let mut nodes = Vec::with_capacity(node_count);
1865        let mut predicates = Vec::with_capacity(predicate_count);
1866        let mut values = Vec::with_capacity(value_count);
1867        let mut node_value_count = 0;
1868        let mut pred_count = 0;
1869
1870        for layer in stack {
1871            nodes.extend(
1872                layer
1873                    .node_dictionary()
1874                    .iter()
1875                    .enumerate()
1876                    .flat_map(|(i, x)| {
1877                        let mapped_node = layer.node_value_id_map().inner_to_outer(i as u64 + 1)
1878                            + node_value_count;
1879                        if node_value_existences[mapped_node as usize] {
1880                            Some((x, mapped_node))
1881                        } else {
1882                            None
1883                        }
1884                    }),
1885            );
1886            predicates.extend(layer.predicate_dictionary().iter().enumerate().flat_map(
1887                |(i, x)| {
1888                    let mapped_predicate =
1889                        layer.predicate_id_map().inner_to_outer(i as u64 + 1) + pred_count;
1890                    if predicate_existences[mapped_predicate as usize] {
1891                        Some((x, mapped_predicate))
1892                    } else {
1893                        None
1894                    }
1895                },
1896            ));
1897            values.extend(
1898                layer
1899                    .value_dictionary()
1900                    .iter()
1901                    .enumerate()
1902                    .flat_map(|(i, x)| {
1903                        let mapped_value = layer
1904                            .node_value_id_map()
1905                            .inner_to_outer(i as u64 + layer.node_dict_len() as u64 + 1)
1906                            + node_value_count;
1907                        if node_value_existences[mapped_value as usize] {
1908                            Some((x, mapped_value))
1909                        } else {
1910                            None
1911                        }
1912                    }),
1913            );
1914
1915            node_value_count += (layer.node_dictionary().num_entries()
1916                + layer.value_dictionary().num_entries()) as u64;
1917            pred_count += layer.predicate_dictionary().num_entries() as u64;
1918        }
1919        nodes.sort();
1920        predicates.sort();
1921        values.sort();
1922
1923        let mut node_value_map: Vec<u64> = vec![0; node_count + value_count + 1];
1924        for (remapped, (_, old)) in nodes.iter().enumerate() {
1925            node_value_map[*old as usize] = remapped as u64 + 1;
1926        }
1927
1928        let node_count = nodes.len();
1929        for (remapped, (_, old)) in values.iter().enumerate() {
1930            node_value_map[*old as usize] = (remapped + node_count + 1) as u64;
1931        }
1932
1933        let mut pred_map: Vec<u64> = vec![0; predicate_count + 1];
1934        for (remapped, (_, old)) in predicates.iter().enumerate() {
1935            pred_map[*old as usize] = remapped as u64 + 1;
1936        }
1937
1938        let layer_name = self.create_directory().await?;
1939        let base_layer_files = self.base_layer_files(layer_name).await?;
1940        let mut builder = BaseLayerFileBuilder::from_files(&base_layer_files).await?;
1941        builder.add_nodes_bytes(nodes.into_iter().map(|(x, _)| x.to_bytes()));
1942        builder.add_predicates_bytes(predicates.into_iter().map(|(x, _)| x.to_bytes()));
1943        builder.add_values(values.into_iter().map(|(x, _)| x));
1944        let mut builder = builder.into_phase2().await?;
1945        let mut triples = Vec::with_capacity(num_triples);
1946        triples.extend(layer.triples().map(move |t| {
1947            IdTriple::new(
1948                node_value_map[t.subject as usize],
1949                pred_map[t.predicate as usize],
1950                node_value_map[t.object as usize],
1951            )
1952        }));
1953        triples.sort();
1954
1955        builder.add_id_triples(triples.into_iter()).await?;
1956        builder.finalize().await?;
1957
1958        self.finalize_layer(layer_name).await?;
1959
1960        Ok(layer_name)
1961    }
1962
1963    async fn squash_upto(&self, layer: Arc<InternalLayer>, upto: [u32; 5]) -> io::Result<[u32; 5]> {
1964        let mut base_node_count = 0;
1965        let mut base_pred_count = 0;
1966        let mut base_value_count = 0;
1967        walk_backwards_from_disk!(self, upto, current, {
1968            base_node_count += self.get_node_count(current).await?.unwrap_or(0) as u64;
1969            base_value_count += self.get_value_count(current).await?.unwrap_or(0) as u64;
1970            base_pred_count += self.get_predicate_count(current).await?.unwrap_or(0) as u64;
1971        });
1972        let base_node_value_count = base_node_count + base_value_count;
1973        let mut node_value_count = base_node_value_count;
1974        let mut pred_count = base_pred_count;
1975        let stack_names = self
1976            .retrieve_layer_stack_names_upto(layer.name(), upto)
1977            .await?;
1978
1979        let mut structures = Vec::with_capacity(stack_names.len());
1980        for layer_name in stack_names {
1981            let node_dict = self.get_node_dictionary(layer_name).await?;
1982            let predicate_dict = self.get_predicate_dictionary(layer_name).await?;
1983            let value_dict = self.get_value_dictionary(layer_name).await?;
1984            let node_value_id_map = self.get_node_value_idmap(layer_name).await?;
1985            let predicate_id_map = self.get_predicate_idmap(layer_name).await?;
1986
1987            structures.push((
1988                node_dict,
1989                predicate_dict,
1990                value_dict,
1991                node_value_id_map,
1992                predicate_id_map,
1993            ));
1994        }
1995
1996        let stack_node_count = structures
1997            .iter()
1998            .map(|d| d.0.as_ref().map(|d| d.num_entries()).unwrap_or(0))
1999            .sum();
2000        let stack_pred_count = structures
2001            .iter()
2002            .map(|d| d.1.as_ref().map(|d| d.num_entries()).unwrap_or(0))
2003            .sum();
2004        let stack_value_count = structures
2005            .iter()
2006            .map(|d| d.2.as_ref().map(|d| d.num_entries()).unwrap_or(0))
2007            .sum();
2008        let stack_node_value_count = stack_node_count + stack_value_count;
2009
2010        // figure out what keys are actually used through a prepass over all triples
2011        let mut node_value_existences = bitvec![0;stack_node_value_count as usize+1];
2012        let mut predicate_existences = bitvec![0;stack_pred_count as usize+1];
2013        let mut num_triple_changes = 0;
2014        let layer_changes_upto = self.layer_changes_upto(layer.name(), upto).await?;
2015        for (change_type, triple) in layer_changes_upto.clone() {
2016            num_triple_changes += 1;
2017            if change_type == TripleChange::Removal {
2018                // removals can't possibly have dictionary entries
2019                // that are part of this layer stack, as they are
2020                // relative to upto, and therefore cannot contain
2021                // dictionary entries that didn't yet exist for the
2022                // upto layer.
2023                continue;
2024            }
2025            if triple.subject >= base_node_value_count {
2026                node_value_existences.set((triple.subject - base_node_value_count) as usize, true);
2027            }
2028            if triple.predicate >= base_pred_count {
2029                predicate_existences.set((triple.predicate - base_pred_count) as usize, true);
2030            }
2031            if triple.object >= base_node_value_count {
2032                node_value_existences.set((triple.object - base_node_value_count) as usize, true);
2033            }
2034        }
2035
2036        let mut nodes = Vec::with_capacity(stack_node_count);
2037        let mut predicates = Vec::with_capacity(stack_pred_count);
2038        let mut values = Vec::with_capacity(stack_value_count);
2039
2040        for (node_dict, predicate_dict, value_dict, node_value_id_map, predicate_id_map) in
2041            structures
2042        {
2043            let node_dict_len = node_dict.as_ref().map(|d| d.num_entries()).unwrap_or(0) as u64;
2044            let predicate_dict_len = predicate_dict
2045                .as_ref()
2046                .map(|d| d.num_entries())
2047                .unwrap_or(0) as u64;
2048            let value_dict_len = value_dict.as_ref().map(|d| d.num_entries()).unwrap_or(0) as u64;
2049            if let Some(node_dict) = node_dict.as_ref() {
2050                nodes.extend(node_dict.iter().enumerate().flat_map(|(i, x)| {
2051                    let mapped_node = node_value_id_map
2052                        .as_ref()
2053                        .map(|m| m.inner_to_outer(i as u64 + 1) + node_value_count)
2054                        .unwrap_or(i as u64 + node_value_count + 1);
2055                    if node_value_existences[(mapped_node - base_node_value_count) as usize] {
2056                        Some((x, mapped_node))
2057                    } else {
2058                        None
2059                    }
2060                }));
2061            }
2062            if let Some(predicate_dict) = predicate_dict.as_ref() {
2063                predicates.extend(predicate_dict.iter().enumerate().flat_map(|(i, x)| {
2064                    let mapped_predicate = predicate_id_map
2065                        .as_ref()
2066                        .map(|m| m.inner_to_outer(i as u64 + 1) + pred_count)
2067                        .unwrap_or(i as u64 + pred_count + 1);
2068                    if predicate_existences[(mapped_predicate - base_pred_count) as usize] {
2069                        Some((x, mapped_predicate))
2070                    } else {
2071                        None
2072                    }
2073                }));
2074            }
2075            if let Some(value_dict) = value_dict.as_ref() {
2076                values.extend(value_dict.iter().enumerate().flat_map(|(i, x)| {
2077                    let mapped_value = node_value_id_map
2078                        .as_ref()
2079                        .map(|m| m.inner_to_outer(i as u64 + node_dict_len + 1) + node_value_count)
2080                        .unwrap_or(i as u64 + node_value_count + node_dict_len + 1);
2081                    if node_value_existences[(mapped_value - base_node_value_count) as usize] {
2082                        Some((x, mapped_value))
2083                    } else {
2084                        None
2085                    }
2086                }));
2087            }
2088            node_value_count += node_dict_len + value_dict_len;
2089            pred_count += predicate_dict_len;
2090        }
2091        nodes.sort();
2092        predicates.sort();
2093        values.sort();
2094
2095        let mut node_value_map: Vec<u64> = vec![0; stack_node_count + stack_value_count + 1];
2096        for (remapped, (_, old)) in nodes.iter().enumerate() {
2097            node_value_map[(*old - base_node_value_count) as usize] =
2098                remapped as u64 + base_node_value_count + 1;
2099        }
2100        let node_count = nodes.len();
2101        for (remapped, (_, old)) in values.iter().enumerate() {
2102            node_value_map[(*old - base_node_value_count) as usize] =
2103                (remapped + node_count + 1) as u64 + base_node_value_count;
2104        }
2105
2106        let mut pred_map: Vec<u64> = vec![0; stack_pred_count + 1];
2107        for (remapped, (_, old)) in predicates.iter().enumerate() {
2108            pred_map[(*old - base_pred_count) as usize] = remapped as u64 + base_pred_count + 1;
2109        }
2110
2111        let predicate_count = predicates.len();
2112        let value_count = values.len();
2113
2114        let upto_layer = self.get_layer(upto).await?.expect("expected upto to exist");
2115        let layer_name = self.create_directory().await?;
2116        let child_layer_files = self.child_layer_files(layer_name).await?;
2117        let mut builder = DictionarySetFileBuilder::from_files(
2118            child_layer_files.node_dictionary_files.clone(),
2119            child_layer_files.predicate_dictionary_files.clone(),
2120            child_layer_files.value_dictionary_files.clone(),
2121        )
2122        .await?;
2123
2124        builder.add_nodes_bytes(nodes.into_iter().map(|(x, _)| x.to_bytes()));
2125        builder.add_predicates_bytes(predicates.into_iter().map(|(x, _)| x.to_bytes()));
2126        builder.add_values(values.into_iter().map(|(x, _)| x));
2127        builder.finalize().await?;
2128
2129        // TODO use more inner stuff to avoid parent checks as they are unnecessary here
2130        let mut builder = ChildLayerFileBuilderPhase2::new(
2131            upto_layer,
2132            child_layer_files,
2133            node_count,
2134            predicate_count,
2135            value_count,
2136        )
2137        .await?;
2138        let mut triple_changes = Vec::with_capacity(num_triple_changes);
2139        for (change_type, t) in layer_changes_upto {
2140            let mapped_subject = if t.subject <= base_node_value_count {
2141                t.subject
2142            } else {
2143                node_value_map[(t.subject - base_node_value_count) as usize]
2144            };
2145            let mapped_predicate = if t.predicate <= base_pred_count {
2146                t.predicate
2147            } else {
2148                pred_map[(t.predicate - base_pred_count) as usize]
2149            };
2150            let mapped_object = if t.object <= base_node_value_count {
2151                t.object
2152            } else {
2153                node_value_map[(t.object - base_node_value_count) as usize]
2154            };
2155            triple_changes.push((
2156                change_type,
2157                IdTriple::new(mapped_subject, mapped_predicate, mapped_object),
2158            ));
2159        }
2160        triple_changes.sort();
2161        for (change_type, mapped_triple) in triple_changes.into_iter() {
2162            match change_type {
2163                TripleChange::Addition => {
2164                    builder
2165                        .add_triple_unchecked(
2166                            mapped_triple.subject,
2167                            mapped_triple.predicate,
2168                            mapped_triple.object,
2169                        )
2170                        .await?
2171                }
2172                TripleChange::Removal => {
2173                    builder
2174                        .remove_triple_unchecked(
2175                            mapped_triple.subject,
2176                            mapped_triple.predicate,
2177                            mapped_triple.object,
2178                        )
2179                        .await?
2180                }
2181            }
2182        }
2183        builder.finalize().await?;
2184        self.write_parent_file(layer_name, upto).await?;
2185        self.finalize_layer(layer_name).await?;
2186
2187        Ok(layer_name)
2188    }
2189
2190    async fn merge_base_layer(
2191        &self,
2192        layers: &[[u32; 5]],
2193        temp_path: &Path,
2194    ) -> io::Result<[u32; 5]> {
2195        let mut layer_files = Vec::with_capacity(layers.len());
2196        for layer in layers {
2197            if self.layer_has_parent(*layer).await? {
2198                return Err(io::Error::new(
2199                    io::ErrorKind::Other,
2200                    format!(
2201                        "given layer is not a base layer: {}",
2202                        name_to_string(*layer)
2203                    ),
2204                ));
2205            }
2206            layer_files.push(self.base_layer_files(*layer).await?);
2207        }
2208
2209        let output_name = self.create_directory().await?;
2210        let output_layer_files = self.base_layer_files(output_name).await?;
2211
2212        merge_base_layers(&layer_files, output_layer_files, temp_path).await?;
2213
2214        self.finalize(output_name).await?;
2215
2216        Ok(output_name)
2217    }
2218
2219    async fn layer_is_ancestor_of(
2220        &self,
2221        mut descendant: [u32; 5],
2222        ancestor: [u32; 5],
2223    ) -> io::Result<bool> {
2224        loop {
2225            if ancestor == descendant {
2226                return Ok(true);
2227            }
2228
2229            if self.layer_has_parent(descendant).await? {
2230                let parent = self.read_parent_file(descendant).await?;
2231                descendant = parent;
2232            } else {
2233                return Ok(false);
2234            }
2235        }
2236    }
2237
2238    async fn triple_addition_exists(
2239        &self,
2240        layer: [u32; 5],
2241        subject: u64,
2242        predicate: u64,
2243        object: u64,
2244    ) -> io::Result<bool> {
2245        let (subjects_file, s_p_aj_files, sp_o_aj_files) =
2246            self.triple_addition_files(layer).await?;
2247
2248        file_triple_exists(
2249            subjects_file,
2250            s_p_aj_files,
2251            sp_o_aj_files,
2252            subject,
2253            predicate,
2254            object,
2255        )
2256        .await
2257    }
2258
2259    async fn triple_removal_exists(
2260        &self,
2261        layer: [u32; 5],
2262        subject: u64,
2263        predicate: u64,
2264        object: u64,
2265    ) -> io::Result<bool> {
2266        if let Some((subjects_file, s_p_aj_files, sp_o_aj_files)) =
2267            self.triple_removal_files(layer).await?
2268        {
2269            file_triple_exists(
2270                subjects_file,
2271                s_p_aj_files,
2272                sp_o_aj_files,
2273                subject,
2274                predicate,
2275                object,
2276            )
2277            .await
2278        } else {
2279            Ok(false)
2280        }
2281    }
2282
2283    async fn triple_additions(
2284        &self,
2285        layer: [u32; 5],
2286    ) -> io::Result<OptInternalLayerTripleSubjectIterator> {
2287        let (subjects_file, s_p_aj_files, sp_o_aj_files) =
2288            self.triple_addition_files(layer).await?;
2289
2290        Ok(OptInternalLayerTripleSubjectIterator(Some(
2291            file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files).await?,
2292        )))
2293    }
2294
2295    async fn triple_removals(
2296        &self,
2297        layer: [u32; 5],
2298    ) -> io::Result<OptInternalLayerTripleSubjectIterator> {
2299        if let Some((subjects_file, s_p_aj_files, sp_o_aj_files)) =
2300            self.triple_removal_files(layer).await?
2301        {
2302            Ok(OptInternalLayerTripleSubjectIterator(Some(
2303                file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files).await?,
2304            )))
2305        } else {
2306            Ok(OptInternalLayerTripleSubjectIterator(None))
2307        }
2308    }
2309
2310    async fn triple_additions_s(
2311        &self,
2312        layer: [u32; 5],
2313        subject: u64,
2314    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2315        let (subjects_file, s_p_aj_files, sp_o_aj_files) =
2316            self.triple_addition_files(layer).await?;
2317
2318        Ok(Box::new(
2319            file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files)
2320                .await?
2321                .seek_subject(subject)
2322                .take_while(move |t| t.subject == subject),
2323        ) as Box<dyn Iterator<Item = _> + Send>)
2324    }
2325
2326    async fn triple_removals_s(
2327        &self,
2328        layer: [u32; 5],
2329        subject: u64,
2330    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2331        if let Some((subjects_file, s_p_aj_files, sp_o_aj_files)) =
2332            self.triple_removal_files(layer).await?
2333        {
2334            Ok(Box::new(
2335                file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files)
2336                    .await?
2337                    .seek_subject(subject)
2338                    .take_while(move |t| t.subject == subject),
2339            ) as Box<dyn Iterator<Item = _> + Send>)
2340        } else {
2341            Ok(Box::new(std::iter::empty()) as Box<dyn Iterator<Item = _> + Send>)
2342        }
2343    }
2344
2345    async fn triple_additions_sp(
2346        &self,
2347        layer: [u32; 5],
2348        subject: u64,
2349        predicate: u64,
2350    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2351        let (subjects_file, s_p_aj_files, sp_o_aj_files) =
2352            self.triple_addition_files(layer).await?;
2353
2354        Ok(Box::new(
2355            file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files)
2356                .await?
2357                .seek_subject_predicate(subject, predicate)
2358                .take_while(move |t| t.predicate == predicate && t.subject == subject),
2359        ) as Box<dyn Iterator<Item = _> + Send>)
2360    }
2361
2362    async fn triple_removals_sp(
2363        &self,
2364        layer: [u32; 5],
2365        subject: u64,
2366        predicate: u64,
2367    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2368        if let Some((subjects_file, s_p_aj_files, sp_o_aj_files)) =
2369            self.triple_removal_files(layer).await?
2370        {
2371            Ok(Box::new(
2372                file_triple_iterator(subjects_file, s_p_aj_files, sp_o_aj_files)
2373                    .await?
2374                    .seek_subject_predicate(subject, predicate)
2375                    .take_while(move |t| t.predicate == predicate && t.subject == subject),
2376            ) as Box<dyn Iterator<Item = _> + Send>)
2377        } else {
2378            Ok(Box::new(std::iter::empty()) as Box<dyn Iterator<Item = _> + Send>)
2379        }
2380    }
2381
2382    async fn triple_additions_p(
2383        &self,
2384        layer: [u32; 5],
2385        predicate: u64,
2386    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2387        let (subjects_file, s_p_aj_files, sp_o_aj_files) =
2388            self.triple_addition_files(layer).await?;
2389        let predicate_wavelet_files = self.predicate_wavelet_addition_files(layer).await?;
2390
2391        Ok(Box::new(
2392            file_triple_iterator_by_predicate(
2393                subjects_file,
2394                s_p_aj_files,
2395                sp_o_aj_files,
2396                predicate_wavelet_files,
2397                predicate,
2398            )
2399            .await?,
2400        ) as Box<dyn Iterator<Item = _> + Send>)
2401    }
2402
2403    async fn triple_removals_p(
2404        &self,
2405        layer: [u32; 5],
2406        predicate: u64,
2407    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2408        if let (Some((subjects_file, s_p_aj_files, sp_o_aj_files)), Some(predicate_wavelet_files)) = (
2409            self.triple_removal_files(layer).await?,
2410            self.predicate_wavelet_removal_files(layer).await?,
2411        ) {
2412            Ok(Box::new(
2413                file_triple_iterator_by_predicate(
2414                    subjects_file,
2415                    s_p_aj_files,
2416                    sp_o_aj_files,
2417                    predicate_wavelet_files,
2418                    predicate,
2419                )
2420                .await?,
2421            ) as Box<dyn Iterator<Item = _> + Send>)
2422        } else {
2423            Ok(Box::new(std::iter::empty()) as Box<dyn Iterator<Item = _> + Send>)
2424        }
2425    }
2426
2427    async fn triple_additions_o(
2428        &self,
2429        layer: [u32; 5],
2430        object: u64,
2431    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2432        let (subjects_file, objects_file, o_ps_aj_files, s_p_aj_files) =
2433            self.triple_addition_files_by_object(layer).await?;
2434
2435        Ok(Box::new(
2436            file_triple_iterator_by_object(
2437                subjects_file,
2438                objects_file,
2439                o_ps_aj_files,
2440                s_p_aj_files,
2441                object,
2442            )
2443            .await?
2444            .take_while(move |t| t.object == object),
2445        ) as Box<dyn Iterator<Item = _> + Send>)
2446    }
2447
2448    async fn triple_removals_o(
2449        &self,
2450        layer: [u32; 5],
2451        object: u64,
2452    ) -> io::Result<Box<dyn Iterator<Item = IdTriple> + Send>> {
2453        if let Some((subjects_file, objects_file, o_ps_aj_files, s_p_aj_files)) =
2454            self.triple_removal_files_by_object(layer).await?
2455        {
2456            Ok(Box::new(
2457                file_triple_iterator_by_object(
2458                    subjects_file,
2459                    objects_file,
2460                    o_ps_aj_files,
2461                    s_p_aj_files,
2462                    object,
2463                )
2464                .await?
2465                .take_while(move |t| t.object == object),
2466            ) as Box<dyn Iterator<Item = _> + Send>)
2467        } else {
2468            Ok(Box::new(std::iter::empty()) as Box<dyn Iterator<Item = _> + Send>)
2469        }
2470    }
2471
2472    async fn triple_layer_addition_count(&self, layer: [u32; 5]) -> io::Result<usize> {
2473        let (s_p_nums_file, sp_o_bits_file, predicate_wavelet_files) =
2474            self.triple_layer_addition_count_files(layer).await?;
2475        file_triple_layer_count(s_p_nums_file, sp_o_bits_file, predicate_wavelet_files).await
2476    }
2477
2478    async fn triple_layer_removal_count(&self, layer: [u32; 5]) -> io::Result<usize> {
2479        if let Some((s_p_nums_file, sp_o_bits_file, predicate_wavelet_files)) =
2480            self.triple_layer_removal_count_files(layer).await?
2481        {
2482            file_triple_layer_count(s_p_nums_file, sp_o_bits_file, predicate_wavelet_files).await
2483        } else {
2484            Ok(0)
2485        }
2486    }
2487
2488    async fn retrieve_layer_stack_names(&self, name: [u32; 5]) -> io::Result<Vec<[u32; 5]>> {
2489        let mut result = vec![name];
2490
2491        loop {
2492            if self.layer_has_parent(*result.last().unwrap()).await? {
2493                let parent = self.read_parent_file(*result.last().unwrap()).await?;
2494                result.push(parent);
2495            } else {
2496                result.reverse();
2497
2498                return Ok(result);
2499            }
2500        }
2501    }
2502
2503    async fn retrieve_layer_stack_names_upto(
2504        &self,
2505        name: [u32; 5],
2506        upto: [u32; 5],
2507    ) -> io::Result<Vec<[u32; 5]>> {
2508        let mut result = vec![name];
2509
2510        loop {
2511            if self.layer_has_parent(*result.last().unwrap()).await? {
2512                let parent = self.read_parent_file(*result.last().unwrap()).await?;
2513                if parent == upto {
2514                    break;
2515                }
2516                result.push(parent);
2517            } else {
2518                return Err(io::Error::new(
2519                    io::ErrorKind::NotFound,
2520                    "parent layer not found while retrieving names of layer stack",
2521                ));
2522            }
2523        }
2524
2525        result.reverse();
2526        Ok(result)
2527    }
2528}
2529
2530pub(crate) async fn file_triple_exists<F: FileLoad + FileStore>(
2531    subjects_file: F,
2532    s_p_adjacency_list_files: AdjacencyListFiles<F>,
2533    sp_o_adjacency_list_files: AdjacencyListFiles<F>,
2534    subject: u64,
2535    predicate: u64,
2536    object: u64,
2537) -> io::Result<bool> {
2538    let s_p_maps = s_p_adjacency_list_files.map_all().await?;
2539    let sp_o_maps = sp_o_adjacency_list_files.map_all().await?;
2540
2541    let subjects: Option<MonotonicLogArray> = subjects_file
2542        .map_if_exists()
2543        .await?
2544        .map(|l| LogArray::parse(l).unwrap().into());
2545    let s_p_aj = s_p_maps.into();
2546    let sp_o_aj = sp_o_maps.into();
2547
2548    Ok(layer_triple_exists(
2549        subjects.as_ref(),
2550        &s_p_aj,
2551        &sp_o_aj,
2552        subject,
2553        predicate,
2554        object,
2555    ))
2556}
2557
2558pub(crate) async fn file_triple_iterator<F: FileLoad + FileStore>(
2559    subjects_file: F,
2560    s_p_adjacency_list_files: AdjacencyListFiles<F>,
2561    sp_o_adjacency_list_files: AdjacencyListFiles<F>,
2562) -> io::Result<InternalLayerTripleSubjectIterator> {
2563    let s_p_maps = s_p_adjacency_list_files.map_all().await?;
2564    let sp_o_maps = sp_o_adjacency_list_files.map_all().await?;
2565
2566    let subjects: Option<MonotonicLogArray> = subjects_file
2567        .map_if_exists()
2568        .await?
2569        .map(|l| LogArray::parse(l).unwrap().into());
2570    let s_p_aj = s_p_maps.into();
2571    let sp_o_aj = sp_o_maps.into();
2572
2573    Ok(InternalLayerTripleSubjectIterator::new(
2574        subjects, s_p_aj, sp_o_aj,
2575    ))
2576}
2577
2578pub(crate) async fn file_triple_iterator_by_predicate<F: FileLoad + FileStore>(
2579    subjects_file: F,
2580    s_p_adjacency_list_files: AdjacencyListFiles<F>,
2581    sp_o_adjacency_list_files: AdjacencyListFiles<F>,
2582    predicate_wavelet_files: BitIndexFiles<F>,
2583    predicate: u64,
2584) -> io::Result<impl Iterator<Item = IdTriple> + Send> {
2585    let s_p_maps = s_p_adjacency_list_files.map_all().await?;
2586    let sp_o_maps = sp_o_adjacency_list_files.map_all().await?;
2587    let predicate_wavelet_maps = predicate_wavelet_files.map_all().await?;
2588
2589    let subjects: Option<MonotonicLogArray> = subjects_file
2590        .map_if_exists()
2591        .await?
2592        .map(|l| LogArray::parse(l).unwrap().into());
2593    let s_p_aj: AdjacencyList = s_p_maps.into();
2594    let sp_o_aj: AdjacencyList = sp_o_maps.into();
2595
2596    let width = s_p_aj.nums().width();
2597    let wavelet_bits = predicate_wavelet_maps.into();
2598    let wtree = WaveletTree::from_parts(wavelet_bits, width);
2599    Ok(match wtree.lookup(predicate) {
2600        Some(lookup) => OptInternalLayerTriplePredicateIterator(Some(
2601            InternalLayerTriplePredicateIterator::new(lookup, subjects, s_p_aj, sp_o_aj),
2602        )),
2603        None => OptInternalLayerTriplePredicateIterator(None),
2604    })
2605}
2606
2607pub(crate) async fn file_triple_iterator_by_object<F: FileLoad + FileStore>(
2608    subjects_file: F,
2609    objects_file: F,
2610    o_ps_adjacency_list_files: AdjacencyListFiles<F>,
2611    s_p_adjacency_list_files: AdjacencyListFiles<F>,
2612    object: u64,
2613) -> io::Result<impl Iterator<Item = IdTriple> + Send> {
2614    let subjects: Option<MonotonicLogArray> = subjects_file
2615        .map_if_exists()
2616        .await?
2617        .map(|l| LogArray::parse(l).unwrap().into());
2618    let objects: Option<MonotonicLogArray> = objects_file
2619        .map_if_exists()
2620        .await?
2621        .map(|l| LogArray::parse(l).unwrap().into());
2622
2623    let o_ps_maps = o_ps_adjacency_list_files.map_all().await?;
2624    let s_p_maps = s_p_adjacency_list_files.map_all().await?;
2625    let o_ps_aj: AdjacencyList = o_ps_maps.into();
2626    let s_p_aj: AdjacencyList = s_p_maps.into();
2627
2628    Ok(
2629        InternalLayerTripleObjectIterator::new(subjects, objects, o_ps_aj, s_p_aj, true)
2630            .seek_object(object),
2631    )
2632}
2633
2634pub(crate) async fn file_triple_layer_count<F: FileLoad + FileStore>(
2635    s_p_nums_file: F,
2636    sp_o_bits_file: F,
2637    predicate_wavelet_files: BitIndexFiles<F>,
2638) -> io::Result<usize> {
2639    let (_, width) = logarray_file_get_length_and_width(s_p_nums_file).await?;
2640    let bits_len: usize = bitarray_len_from_file(sp_o_bits_file)
2641        .await?
2642        .try_into()
2643        .unwrap();
2644    let predicate_wavelet_maps = predicate_wavelet_files.map_all().await?;
2645    let wavelet_bits = predicate_wavelet_maps.into();
2646    let wtree = WaveletTree::from_parts(wavelet_bits, width);
2647
2648    Ok(bits_len - wtree.lookup(0).map(|l| l.len()).unwrap_or(0))
2649}
2650
2651#[cfg(test)]
2652mod tests {
2653    use super::*;
2654    use crate::layer::{ObjectType, ValueTriple};
2655    use crate::storage::directory::DirectoryLayerStore;
2656    use crate::storage::memory::MemoryLayerStore;
2657    use std::collections::HashMap;
2658    use tempfile::{tempdir, TempDir};
2659    // these tests are for both the memory store and the directory store
2660    // They test functionality that should really work for both
2661
2662    lazy_static! {
2663        static ref BASE_TRIPLES: Vec<ValueTriple> = vec![
2664            ValueTriple::new_string_value("cow", "says", "moo"),
2665            ValueTriple::new_string_value("cow", "says", "mooo"),
2666            ValueTriple::new_node("cow", "likes", "duck"),
2667            ValueTriple::new_node("cow", "likes", "pig"),
2668            ValueTriple::new_string_value("cow", "name", "clarabelle"),
2669            ValueTriple::new_string_value("pig", "says", "oink"),
2670            ValueTriple::new_node("pig", "hates", "cow"),
2671            ValueTriple::new_string_value("duck", "says", "quack"),
2672            ValueTriple::new_node("duck", "hates", "cow"),
2673            ValueTriple::new_node("duck", "hates", "pig"),
2674            ValueTriple::new_string_value("duck", "name", "donald"),
2675        ];
2676        static ref CHILD_ADDITION_TRIPLES: Vec<ValueTriple> = vec![
2677            ValueTriple::new_string_value("cow", "says", "moooo"),
2678            ValueTriple::new_string_value("cow", "says", "mooooo"),
2679            ValueTriple::new_node("cow", "likes", "horse"),
2680            ValueTriple::new_node("pig", "likes", "platypus"),
2681            ValueTriple::new_node("duck", "hates", "platypus"),
2682        ];
2683        static ref CHILD_REMOVAL_TRIPLES: Vec<ValueTriple> = vec![
2684            ValueTriple::new_string_value("cow", "says", "mooo"),
2685            ValueTriple::new_string_value("cow", "name", "clarabelle"),
2686            ValueTriple::new_node("pig", "hates", "cow"),
2687            ValueTriple::new_node("duck", "hates", "cow"),
2688            ValueTriple::new_node("duck", "hates", "pig"),
2689            ValueTriple::new_string_value("duck", "name", "donald"),
2690        ];
2691    }
2692
2693    async fn example_base_layer<S: LayerStore>(
2694        store: &S,
2695        invalidate: bool,
2696    ) -> io::Result<(
2697        [u32; 5],
2698        Option<Arc<InternalLayer>>,
2699        HashMap<ValueTriple, IdTriple>,
2700    )> {
2701        let mut builder = store.create_base_layer().await?;
2702        let name = builder.name();
2703        for t in BASE_TRIPLES.iter() {
2704            builder.add_value_triple(t.clone());
2705        }
2706        builder.commit_boxed().await?;
2707        let layer = store.get_layer(name).await?.unwrap();
2708
2709        let mut contents = HashMap::with_capacity(BASE_TRIPLES.len());
2710        for t in BASE_TRIPLES.iter() {
2711            let t_id = layer.value_triple_to_id(t).unwrap();
2712            contents.insert(t.clone(), t_id);
2713        }
2714
2715        let layer_opt = match invalidate {
2716            true => None,
2717            false => Some(layer),
2718        };
2719
2720        Ok((name, layer_opt, contents))
2721    }
2722
2723    async fn example_child_layer<S: LayerStore>(
2724        store: &S,
2725        invalidate: bool,
2726    ) -> io::Result<(
2727        [u32; 5],
2728        Option<Arc<InternalLayer>>,
2729        HashMap<ValueTriple, IdTriple>,
2730        HashMap<ValueTriple, IdTriple>,
2731    )> {
2732        let (base_name, _base_layer, _) = example_base_layer(store, false).await?;
2733        let mut builder = store.create_child_layer(base_name).await?;
2734        let name = builder.name();
2735        for t in CHILD_ADDITION_TRIPLES.iter() {
2736            builder.add_value_triple(t.clone());
2737        }
2738        for t in CHILD_REMOVAL_TRIPLES.iter() {
2739            builder.remove_value_triple(t.clone());
2740        }
2741        builder.commit_boxed().await?;
2742        let layer = store.get_layer(name).await?.unwrap();
2743
2744        let mut add_contents = HashMap::with_capacity(BASE_TRIPLES.len());
2745        for t in CHILD_ADDITION_TRIPLES.iter() {
2746            let t_id = layer.value_triple_to_id(t).unwrap();
2747            add_contents.insert(t.clone(), t_id);
2748        }
2749
2750        let mut remove_contents = HashMap::with_capacity(BASE_TRIPLES.len());
2751        for t in CHILD_REMOVAL_TRIPLES.iter() {
2752            let t_id = layer.value_triple_to_id(t).unwrap();
2753            remove_contents.insert(t.clone(), t_id);
2754        }
2755
2756        let layer_opt = match invalidate {
2757            true => None,
2758            false => Some(layer),
2759        };
2760
2761        Ok((name, layer_opt, add_contents, remove_contents))
2762    }
2763
2764    async fn base_layer_counts<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2765        let (name, _layer, _) = example_base_layer(store, invalidate).await?;
2766        assert_eq!(11, store.triple_layer_addition_count(name).await?);
2767        assert_eq!(0, store.triple_layer_removal_count(name).await?);
2768
2769        Ok(())
2770    }
2771
2772    #[tokio::test]
2773    async fn memory_base_layer_counts() {
2774        let store = MemoryLayerStore::new();
2775        base_layer_counts(&store, false).await.unwrap();
2776    }
2777
2778    #[tokio::test]
2779    async fn directory_base_layer_counts() {
2780        let dir = tempdir().unwrap();
2781        let store = DirectoryLayerStore::new(dir.path());
2782        base_layer_counts(&store, false).await.unwrap();
2783    }
2784
2785    async fn child_layer_counts<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2786        let (name, _layer, _, _) = example_child_layer(store, invalidate).await?;
2787        assert_eq!(5, store.triple_layer_addition_count(name).await?);
2788        assert_eq!(6, store.triple_layer_removal_count(name).await?);
2789
2790        Ok(())
2791    }
2792
2793    #[tokio::test]
2794    async fn memory_child_layer_counts() {
2795        let store = MemoryLayerStore::new();
2796        child_layer_counts(&store, false).await.unwrap();
2797    }
2798
2799    #[tokio::test]
2800    async fn directory_child_layer_counts() {
2801        let dir = tempdir().unwrap();
2802        let store = DirectoryLayerStore::new(dir.path());
2803        child_layer_counts(&store, false).await.unwrap();
2804    }
2805
2806    async fn base_layer_addition_exists<S: LayerStore>(
2807        store: &S,
2808        invalidate: bool,
2809    ) -> io::Result<()> {
2810        let (name, _layer, triples) = example_base_layer(store, invalidate).await?;
2811
2812        for t in triples.values() {
2813            assert!(
2814                store
2815                    .triple_addition_exists(name, t.subject, t.predicate, t.object)
2816                    .await?
2817            );
2818        }
2819
2820        assert!(!store.triple_addition_exists(name, 42, 42, 42).await?);
2821        Ok(())
2822    }
2823
2824    #[tokio::test]
2825    async fn memory_base_layer_addition_exists() {
2826        let store = MemoryLayerStore::new();
2827        base_layer_addition_exists(&store, false).await.unwrap();
2828    }
2829
2830    #[tokio::test]
2831    async fn directory_base_layer_addition_exists() {
2832        let dir = tempdir().unwrap();
2833        let store = DirectoryLayerStore::new(dir.path());
2834        base_layer_addition_exists(&store, false).await.unwrap();
2835    }
2836
2837    async fn base_layer_additions<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2838        let (name, _layer, triples) = example_base_layer(store, invalidate).await?;
2839        let mut values: Vec<_> = triples.values().cloned().collect();
2840        values.sort();
2841
2842        let additions: Vec<_> = store.triple_additions(name).await?.collect();
2843        assert_eq!(values, additions);
2844
2845        Ok(())
2846    }
2847
2848    #[tokio::test]
2849    async fn memory_base_layer_additions() {
2850        let store = MemoryLayerStore::new();
2851        base_layer_additions(&store, false).await.unwrap();
2852    }
2853
2854    #[tokio::test]
2855    async fn directory_base_layer_additions() {
2856        let dir = tempdir().unwrap();
2857        let store = DirectoryLayerStore::new(dir.path());
2858        base_layer_additions(&store, false).await.unwrap();
2859    }
2860
2861    async fn base_layer_additions_s<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2862        let (name, _layer, contents) = example_base_layer(store, invalidate).await?;
2863
2864        let mut triples: Vec<_> = contents
2865            .iter()
2866            .filter(|(t, _)| t.subject == "cow")
2867            .map(|(_, t)| t)
2868            .cloned()
2869            .collect();
2870        triples.sort();
2871        assert_eq!(5, triples.len());
2872
2873        let result: Vec<_> = store
2874            .triple_additions_s(name, triples[0].subject)
2875            .await?
2876            .collect();
2877        assert_eq!(triples, result);
2878
2879        assert!(store.triple_additions_s(name, 42).await?.next().is_none());
2880
2881        Ok(())
2882    }
2883
2884    #[tokio::test]
2885    async fn memory_base_layer_additions_s() {
2886        let store = MemoryLayerStore::new();
2887        base_layer_additions_s(&store, false).await.unwrap();
2888    }
2889
2890    #[tokio::test]
2891    async fn directory_base_layer_additions_s() {
2892        let dir = tempdir().unwrap();
2893        let store = DirectoryLayerStore::new(dir.path());
2894        base_layer_additions_s(&store, false).await.unwrap();
2895    }
2896
2897    async fn base_layer_additions_sp<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2898        let (name, _layer, contents) = example_base_layer(store, invalidate).await?;
2899
2900        let mut triples: Vec<_> = contents
2901            .iter()
2902            .filter(|(t, _)| t.subject == "cow" && t.predicate == "says")
2903            .map(|(_, t)| t)
2904            .cloned()
2905            .collect();
2906        triples.sort();
2907        assert_eq!(2, triples.len());
2908
2909        let result: Vec<_> = store
2910            .triple_additions_sp(name, triples[0].subject, triples[0].predicate)
2911            .await?
2912            .collect();
2913        assert_eq!(triples, result);
2914
2915        assert!(store
2916            .triple_additions_sp(name, 42, 42)
2917            .await?
2918            .next()
2919            .is_none());
2920
2921        Ok(())
2922    }
2923
2924    #[tokio::test]
2925    async fn memory_base_layer_additions_sp() {
2926        let store = MemoryLayerStore::new();
2927        base_layer_additions_sp(&store, false).await.unwrap();
2928    }
2929
2930    #[tokio::test]
2931    async fn directory_base_layer_additions_sp() {
2932        let dir = tempdir().unwrap();
2933        let store = DirectoryLayerStore::new(dir.path());
2934        base_layer_additions_sp(&store, false).await.unwrap();
2935    }
2936
2937    async fn base_layer_additions_p<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2938        let (name, _layer, contents) = example_base_layer(store, invalidate).await?;
2939
2940        let mut triples: Vec<_> = contents
2941            .iter()
2942            .filter(|(t, _)| t.predicate == "says")
2943            .map(|(_, t)| t)
2944            .cloned()
2945            .collect();
2946        triples.sort();
2947        assert_eq!(4, triples.len());
2948
2949        let result: Vec<_> = store
2950            .triple_additions_p(name, triples[0].predicate)
2951            .await?
2952            .collect();
2953        assert_eq!(triples, result);
2954
2955        assert!(store.triple_additions_p(name, 42).await?.next().is_none());
2956
2957        Ok(())
2958    }
2959
2960    #[tokio::test]
2961    async fn memory_base_layer_additions_p() {
2962        let store = MemoryLayerStore::new();
2963        base_layer_additions_p(&store, false).await.unwrap();
2964    }
2965
2966    #[tokio::test]
2967    async fn directory_base_layer_additions_p() {
2968        let dir = tempdir().unwrap();
2969        let store = DirectoryLayerStore::new(dir.path());
2970        base_layer_additions_p(&store, false).await.unwrap();
2971    }
2972
2973    async fn base_layer_additions_o<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
2974        let (name, _layer, contents) = example_base_layer(store, invalidate).await?;
2975
2976        let mut triples: Vec<_> = contents
2977            .iter()
2978            .filter(|(t, _)| t.object == ObjectType::Node("cow".to_string()))
2979            .map(|(_, t)| t)
2980            .cloned()
2981            .collect();
2982        triples.sort();
2983        assert_eq!(2, triples.len());
2984
2985        let result: Vec<_> = store
2986            .triple_additions_o(name, triples[0].object)
2987            .await?
2988            .collect();
2989        assert_eq!(triples, result);
2990
2991        assert!(store.triple_additions_o(name, 42).await?.next().is_none());
2992
2993        Ok(())
2994    }
2995
2996    #[tokio::test]
2997    async fn memory_base_layer_additions_o() {
2998        let store = MemoryLayerStore::new();
2999        base_layer_additions_o(&store, false).await.unwrap();
3000    }
3001
3002    #[tokio::test]
3003    async fn directory_base_layer_additions_o() {
3004        let dir = tempdir().unwrap();
3005        let store = DirectoryLayerStore::new(dir.path());
3006        base_layer_additions_o(&store, false).await.unwrap();
3007    }
3008
3009    async fn base_layer_removals<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3010        let (name, _layer, _) = example_base_layer(store, invalidate).await?;
3011
3012        assert!(!store.triple_removal_exists(name, 42, 42, 42).await?);
3013        assert!(store.triple_removals(name).await?.next().is_none());
3014        assert!(store.triple_removals_s(name, 42).await?.next().is_none());
3015        assert!(store
3016            .triple_removals_sp(name, 42, 42)
3017            .await?
3018            .next()
3019            .is_none());
3020        assert!(store.triple_removals_p(name, 42).await?.next().is_none());
3021        assert!(store.triple_removals_o(name, 42).await?.next().is_none());
3022
3023        Ok(())
3024    }
3025
3026    #[tokio::test]
3027    async fn memory_base_layer_removals() {
3028        let store = MemoryLayerStore::new();
3029        base_layer_removals(&store, false).await.unwrap();
3030    }
3031
3032    #[tokio::test]
3033    async fn directory_base_layer_removals() {
3034        let dir = tempdir().unwrap();
3035        let store = DirectoryLayerStore::new(dir.path());
3036        base_layer_removals(&store, false).await.unwrap();
3037    }
3038
3039    async fn child_layer_addition_exists<S: LayerStore>(
3040        store: &S,
3041        invalidate: bool,
3042    ) -> io::Result<()> {
3043        let (name, _layer, triples, _) = example_child_layer(store, invalidate).await?;
3044
3045        for t in triples.values() {
3046            assert!(
3047                store
3048                    .triple_addition_exists(name, t.subject, t.predicate, t.object)
3049                    .await?
3050            );
3051        }
3052
3053        assert!(!store.triple_addition_exists(name, 42, 42, 42).await?);
3054        Ok(())
3055    }
3056
3057    #[tokio::test]
3058    async fn memory_child_layer_addition_exists() {
3059        let store = MemoryLayerStore::new();
3060        child_layer_addition_exists(&store, false).await.unwrap();
3061    }
3062
3063    #[tokio::test]
3064    async fn directory_child_layer_addition_exists() {
3065        let dir = tempdir().unwrap();
3066        let store = DirectoryLayerStore::new(dir.path());
3067        child_layer_addition_exists(&store, false).await.unwrap();
3068    }
3069
3070    async fn child_layer_additions<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3071        let (name, _layer, triples, _) = example_child_layer(store, invalidate).await?;
3072        let mut values: Vec<_> = triples.values().cloned().collect();
3073        values.sort();
3074
3075        let additions: Vec<_> = store.triple_additions(name).await?.collect();
3076        assert_eq!(values, additions);
3077
3078        Ok(())
3079    }
3080
3081    #[tokio::test]
3082    async fn memory_child_layer_additions() {
3083        let store = MemoryLayerStore::new();
3084        child_layer_additions(&store, false).await.unwrap();
3085    }
3086
3087    #[tokio::test]
3088    async fn directory_child_layer_additions() {
3089        let dir = tempdir().unwrap();
3090        let store = DirectoryLayerStore::new(dir.path());
3091        child_layer_additions(&store, false).await.unwrap();
3092    }
3093
3094    async fn child_layer_additions_s<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3095        let (name, _layer, contents, _removals) = example_child_layer(store, invalidate).await?;
3096
3097        let mut triples: Vec<_> = contents
3098            .iter()
3099            .filter(|(t, _)| t.subject == "cow")
3100            .map(|(_, t)| t)
3101            .cloned()
3102            .collect();
3103        triples.sort();
3104        assert_eq!(3, triples.len());
3105
3106        let result: Vec<_> = store
3107            .triple_additions_s(name, triples[0].subject)
3108            .await?
3109            .collect();
3110        assert_eq!(triples, result);
3111
3112        assert!(store.triple_additions_s(name, 42).await?.next().is_none());
3113
3114        Ok(())
3115    }
3116
3117    #[tokio::test]
3118    async fn memory_child_layer_additions_s() {
3119        let store = MemoryLayerStore::new();
3120        child_layer_additions_s(&store, false).await.unwrap();
3121    }
3122
3123    #[tokio::test]
3124    async fn directory_child_layer_additions_s() {
3125        let dir = tempdir().unwrap();
3126        let store = DirectoryLayerStore::new(dir.path());
3127        child_layer_additions_s(&store, false).await.unwrap();
3128    }
3129
3130    async fn child_layer_additions_sp<S: LayerStore>(
3131        store: &S,
3132        invalidate: bool,
3133    ) -> io::Result<()> {
3134        let (name, _layer, contents, _removals) = example_child_layer(store, invalidate).await?;
3135
3136        let mut triples: Vec<_> = contents
3137            .iter()
3138            .filter(|(t, _)| t.subject == "cow" && t.predicate == "says")
3139            .map(|(_, t)| t)
3140            .cloned()
3141            .collect();
3142        triples.sort();
3143        assert_eq!(2, triples.len());
3144
3145        let result: Vec<_> = store
3146            .triple_additions_sp(name, triples[0].subject, triples[0].predicate)
3147            .await?
3148            .collect();
3149        assert_eq!(triples, result);
3150
3151        assert!(store
3152            .triple_additions_sp(name, 42, 42)
3153            .await?
3154            .next()
3155            .is_none());
3156
3157        Ok(())
3158    }
3159
3160    #[tokio::test]
3161    async fn memory_child_layer_additions_sp() {
3162        let store = MemoryLayerStore::new();
3163        child_layer_additions_sp(&store, false).await.unwrap();
3164    }
3165
3166    #[tokio::test]
3167    async fn directory_child_layer_additions_sp() {
3168        let dir = tempdir().unwrap();
3169        let store = DirectoryLayerStore::new(dir.path());
3170        child_layer_additions_sp(&store, false).await.unwrap();
3171    }
3172
3173    async fn child_layer_additions_p<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3174        let (name, _layer, contents, _removals) = example_child_layer(store, invalidate).await?;
3175
3176        let mut triples: Vec<_> = contents
3177            .iter()
3178            .filter(|(t, _)| t.predicate == "likes")
3179            .map(|(_, t)| t)
3180            .cloned()
3181            .collect();
3182        triples.sort();
3183        assert_eq!(2, triples.len());
3184
3185        let result: Vec<_> = store
3186            .triple_additions_p(name, triples[0].predicate)
3187            .await?
3188            .collect();
3189        assert_eq!(triples, result);
3190
3191        assert!(store.triple_additions_p(name, 42).await?.next().is_none());
3192
3193        Ok(())
3194    }
3195
3196    #[tokio::test]
3197    async fn memory_child_layer_additions_p() {
3198        let store = MemoryLayerStore::new();
3199        child_layer_additions_p(&store, false).await.unwrap();
3200    }
3201
3202    #[tokio::test]
3203    async fn directory_child_layer_additions_p() {
3204        let dir = tempdir().unwrap();
3205        let store = DirectoryLayerStore::new(dir.path());
3206        child_layer_additions_p(&store, false).await.unwrap();
3207    }
3208
3209    async fn child_layer_additions_o<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3210        let (name, _layer, contents, _removals) = example_child_layer(store, invalidate).await?;
3211
3212        let mut triples: Vec<_> = contents
3213            .iter()
3214            .filter(|(t, _)| t.object == ObjectType::Node("platypus".to_string()))
3215            .map(|(_, t)| t)
3216            .cloned()
3217            .collect();
3218        triples.sort();
3219        assert_eq!(2, triples.len());
3220
3221        let result: Vec<_> = store
3222            .triple_additions_o(name, triples[0].object)
3223            .await?
3224            .collect();
3225        assert_eq!(triples, result);
3226
3227        assert!(store.triple_additions_o(name, 42).await?.next().is_none());
3228
3229        Ok(())
3230    }
3231
3232    #[tokio::test]
3233    async fn memory_child_layer_additions_o() {
3234        let store = MemoryLayerStore::new();
3235        child_layer_additions_o(&store, false).await.unwrap();
3236    }
3237
3238    #[tokio::test]
3239    async fn directory_child_layer_additions_o() {
3240        let dir = tempdir().unwrap();
3241        let store = DirectoryLayerStore::new(dir.path());
3242        child_layer_additions_o(&store, false).await.unwrap();
3243    }
3244
3245    async fn child_layer_removal_exists<S: LayerStore>(
3246        store: &S,
3247        invalidate: bool,
3248    ) -> io::Result<()> {
3249        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3250
3251        for t in removals.values() {
3252            assert!(
3253                store
3254                    .triple_removal_exists(name, t.subject, t.predicate, t.object)
3255                    .await?
3256            );
3257        }
3258
3259        assert!(!store.triple_removal_exists(name, 42, 42, 42).await?);
3260        Ok(())
3261    }
3262
3263    #[tokio::test]
3264    async fn memory_child_layer_removal_exists() {
3265        let store = MemoryLayerStore::new();
3266        child_layer_removal_exists(&store, false).await.unwrap();
3267    }
3268
3269    #[tokio::test]
3270    async fn directory_child_layer_removal_exists() {
3271        let dir = tempdir().unwrap();
3272        let store = DirectoryLayerStore::new(dir.path());
3273        child_layer_removal_exists(&store, false).await.unwrap();
3274    }
3275
3276    async fn child_layer_removals<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3277        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3278        let mut values: Vec<_> = removals.values().cloned().collect();
3279        values.sort();
3280
3281        let removals: Vec<_> = store.triple_removals(name).await?.collect();
3282        assert_eq!(values, removals);
3283
3284        Ok(())
3285    }
3286
3287    #[tokio::test]
3288    async fn memory_child_layer_removals() {
3289        let store = MemoryLayerStore::new();
3290        child_layer_removals(&store, false).await.unwrap();
3291    }
3292
3293    #[tokio::test]
3294    async fn directory_child_layer_removals() {
3295        let dir = tempdir().unwrap();
3296        let store = DirectoryLayerStore::new(dir.path());
3297        child_layer_removals(&store, false).await.unwrap();
3298    }
3299
3300    async fn child_layer_removals_s<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3301        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3302
3303        let mut triples: Vec<_> = removals
3304            .iter()
3305            .filter(|(t, _)| t.subject == "duck")
3306            .map(|(_, t)| t)
3307            .cloned()
3308            .collect();
3309        triples.sort();
3310        assert_eq!(3, triples.len());
3311
3312        let result: Vec<_> = store
3313            .triple_removals_s(name, triples[0].subject)
3314            .await?
3315            .collect();
3316        assert_eq!(triples, result);
3317
3318        assert!(store.triple_removals_s(name, 42).await?.next().is_none());
3319
3320        Ok(())
3321    }
3322
3323    #[tokio::test]
3324    async fn memory_child_layer_removals_s() {
3325        let store = MemoryLayerStore::new();
3326        child_layer_removals_s(&store, false).await.unwrap();
3327    }
3328
3329    #[tokio::test]
3330    async fn directory_child_layer_removals_s() {
3331        let dir = tempdir().unwrap();
3332        let store = DirectoryLayerStore::new(dir.path());
3333        child_layer_removals_s(&store, false).await.unwrap();
3334    }
3335
3336    async fn child_layer_removals_sp<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3337        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3338
3339        let mut triples: Vec<_> = removals
3340            .iter()
3341            .filter(|(t, _)| t.subject == "duck" && t.predicate == "hates")
3342            .map(|(_, t)| t)
3343            .cloned()
3344            .collect();
3345        triples.sort();
3346        assert_eq!(2, triples.len());
3347
3348        let result: Vec<_> = store
3349            .triple_removals_sp(name, triples[0].subject, triples[0].predicate)
3350            .await?
3351            .collect();
3352        assert_eq!(triples, result);
3353
3354        assert!(store
3355            .triple_removals_sp(name, 42, 42)
3356            .await?
3357            .next()
3358            .is_none());
3359
3360        Ok(())
3361    }
3362
3363    #[tokio::test]
3364    async fn memory_child_layer_removals_sp() {
3365        let store = MemoryLayerStore::new();
3366        child_layer_removals_sp(&store, false).await.unwrap();
3367    }
3368
3369    #[tokio::test]
3370    async fn directory_child_layer_removals_sp() {
3371        let dir = tempdir().unwrap();
3372        let store = DirectoryLayerStore::new(dir.path());
3373        child_layer_removals_sp(&store, false).await.unwrap();
3374    }
3375
3376    async fn child_layer_removals_p<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3377        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3378
3379        let mut triples: Vec<_> = removals
3380            .iter()
3381            .filter(|(t, _)| t.predicate == "hates")
3382            .map(|(_, t)| t)
3383            .cloned()
3384            .collect();
3385        triples.sort();
3386        assert_eq!(3, triples.len());
3387
3388        let result: Vec<_> = store
3389            .triple_removals_p(name, triples[0].predicate)
3390            .await?
3391            .collect();
3392        assert_eq!(triples, result);
3393
3394        assert!(store.triple_removals_p(name, 42).await?.next().is_none());
3395
3396        Ok(())
3397    }
3398
3399    #[tokio::test]
3400    async fn memory_child_layer_removals_p() {
3401        let store = MemoryLayerStore::new();
3402        child_layer_removals_p(&store, false).await.unwrap();
3403    }
3404
3405    #[tokio::test]
3406    async fn directory_child_layer_removals_p() {
3407        let dir = tempdir().unwrap();
3408        let store = DirectoryLayerStore::new(dir.path());
3409        child_layer_removals_p(&store, false).await.unwrap();
3410    }
3411
3412    async fn child_layer_removals_o<S: LayerStore>(store: &S, invalidate: bool) -> io::Result<()> {
3413        let (name, _layer, _additions, removals) = example_child_layer(store, invalidate).await?;
3414
3415        let mut triples: Vec<_> = removals
3416            .iter()
3417            .filter(|(t, _)| t.object == ObjectType::Node("cow".to_string()))
3418            .map(|(_, t)| t)
3419            .cloned()
3420            .collect();
3421        triples.sort();
3422        assert_eq!(2, triples.len());
3423
3424        let result: Vec<_> = store
3425            .triple_removals_o(name, triples[0].object)
3426            .await?
3427            .collect();
3428        assert_eq!(triples, result);
3429
3430        assert!(store.triple_removals_o(name, 42).await?.next().is_none());
3431
3432        Ok(())
3433    }
3434
3435    #[tokio::test]
3436    async fn memory_child_layer_removals_o() {
3437        let store = MemoryLayerStore::new();
3438        child_layer_removals_o(&store, false).await.unwrap();
3439    }
3440
3441    #[tokio::test]
3442    async fn directory_child_layer_removals_o() {
3443        let dir = tempdir().unwrap();
3444        let store = DirectoryLayerStore::new(dir.path());
3445        child_layer_removals_o(&store, false).await.unwrap();
3446    }
3447
3448    fn make_cached_store() -> (TempDir, CachedLayerStore) {
3449        let dir = tempdir().unwrap();
3450        let store = DirectoryLayerStore::new(dir.path());
3451        let cache = CachedLayerStore::new(store, LockingHashMapLayerCache::new());
3452
3453        (dir, cache)
3454    }
3455
3456    #[tokio::test]
3457    async fn cached_base_layer_counts() {
3458        let (_dir, store) = make_cached_store();
3459        base_layer_counts(&store, false).await.unwrap();
3460    }
3461
3462    #[tokio::test]
3463    async fn uncached_base_layer_counts() {
3464        let (_dir, store) = make_cached_store();
3465        base_layer_counts(&store, true).await.unwrap();
3466    }
3467
3468    #[tokio::test]
3469    async fn cached_child_layer_counts() {
3470        let (_dir, store) = make_cached_store();
3471        child_layer_counts(&store, false).await.unwrap();
3472    }
3473
3474    #[tokio::test]
3475    async fn uncached_child_layer_counts() {
3476        let (_dir, store) = make_cached_store();
3477        child_layer_counts(&store, true).await.unwrap();
3478    }
3479
3480    #[tokio::test]
3481    async fn cached_base_layer_addition_exists() {
3482        let (_dir, store) = make_cached_store();
3483        base_layer_addition_exists(&store, false).await.unwrap();
3484    }
3485
3486    #[tokio::test]
3487    async fn uncached_base_layer_addition_exists() {
3488        let (_dir, store) = make_cached_store();
3489        base_layer_addition_exists(&store, true).await.unwrap();
3490    }
3491
3492    #[tokio::test]
3493    async fn cached_base_layer_additions() {
3494        let (_dir, store) = make_cached_store();
3495        base_layer_additions(&store, false).await.unwrap();
3496    }
3497
3498    #[tokio::test]
3499    async fn uncached_base_layer_additions() {
3500        let (_dir, store) = make_cached_store();
3501        base_layer_additions(&store, true).await.unwrap();
3502    }
3503
3504    #[tokio::test]
3505    async fn cached_base_layer_additions_s() {
3506        let (_dir, store) = make_cached_store();
3507        base_layer_additions_s(&store, false).await.unwrap();
3508    }
3509
3510    #[tokio::test]
3511    async fn uncached_base_layer_additions_s() {
3512        let (_dir, store) = make_cached_store();
3513        base_layer_additions_s(&store, true).await.unwrap();
3514    }
3515
3516    #[tokio::test]
3517    async fn cached_base_layer_additions_sp() {
3518        let (_dir, store) = make_cached_store();
3519        base_layer_additions_sp(&store, false).await.unwrap();
3520    }
3521
3522    #[tokio::test]
3523    async fn uncached_base_layer_additions_sp() {
3524        let (_dir, store) = make_cached_store();
3525        base_layer_additions_sp(&store, true).await.unwrap();
3526    }
3527
3528    #[tokio::test]
3529    async fn cached_base_layer_additions_p() {
3530        let (_dir, store) = make_cached_store();
3531        base_layer_additions_p(&store, false).await.unwrap();
3532    }
3533
3534    #[tokio::test]
3535    async fn uncached_base_layer_additions_p() {
3536        let (_dir, store) = make_cached_store();
3537        base_layer_additions_p(&store, true).await.unwrap();
3538    }
3539
3540    #[tokio::test]
3541    async fn cached_base_layer_additions_o() {
3542        let (_dir, store) = make_cached_store();
3543        base_layer_additions_o(&store, false).await.unwrap();
3544    }
3545
3546    #[tokio::test]
3547    async fn uncached_base_layer_additions_o() {
3548        let (_dir, store) = make_cached_store();
3549        base_layer_additions_o(&store, true).await.unwrap();
3550    }
3551
3552    #[tokio::test]
3553    async fn cached_base_layer_removals() {
3554        let (_dir, store) = make_cached_store();
3555        base_layer_removals(&store, false).await.unwrap();
3556    }
3557
3558    #[tokio::test]
3559    async fn uncached_base_layer_removals() {
3560        let (_dir, store) = make_cached_store();
3561        base_layer_removals(&store, true).await.unwrap();
3562    }
3563
3564    #[tokio::test]
3565    async fn cached_child_layer_addition_exists() {
3566        let (_dir, store) = make_cached_store();
3567        child_layer_addition_exists(&store, false).await.unwrap();
3568    }
3569
3570    #[tokio::test]
3571    async fn uncached_child_layer_addition_exists() {
3572        let (_dir, store) = make_cached_store();
3573        child_layer_addition_exists(&store, true).await.unwrap();
3574    }
3575
3576    #[tokio::test]
3577    async fn cached_child_layer_additions() {
3578        let (_dir, store) = make_cached_store();
3579        child_layer_additions(&store, false).await.unwrap();
3580    }
3581
3582    #[tokio::test]
3583    async fn uncached_child_layer_additions() {
3584        let (_dir, store) = make_cached_store();
3585        child_layer_additions(&store, true).await.unwrap();
3586    }
3587
3588    #[tokio::test]
3589    async fn cached_child_layer_additions_s() {
3590        let (_dir, store) = make_cached_store();
3591        child_layer_additions_s(&store, false).await.unwrap();
3592    }
3593
3594    #[tokio::test]
3595    async fn uncached_child_layer_additions_s() {
3596        let (_dir, store) = make_cached_store();
3597        child_layer_additions_s(&store, true).await.unwrap();
3598    }
3599
3600    #[tokio::test]
3601    async fn cached_child_layer_additions_sp() {
3602        let (_dir, store) = make_cached_store();
3603        child_layer_additions_sp(&store, false).await.unwrap();
3604    }
3605
3606    #[tokio::test]
3607    async fn uncached_child_layer_additions_sp() {
3608        let (_dir, store) = make_cached_store();
3609        child_layer_additions_sp(&store, true).await.unwrap();
3610    }
3611
3612    #[tokio::test]
3613    async fn cached_child_layer_additions_p() {
3614        let (_dir, store) = make_cached_store();
3615        child_layer_additions_p(&store, false).await.unwrap();
3616    }
3617
3618    #[tokio::test]
3619    async fn uncached_child_layer_additions_p() {
3620        let (_dir, store) = make_cached_store();
3621        child_layer_additions_p(&store, true).await.unwrap();
3622    }
3623
3624    #[tokio::test]
3625    async fn cached_child_layer_additions_o() {
3626        let (_dir, store) = make_cached_store();
3627        child_layer_additions_o(&store, false).await.unwrap();
3628    }
3629
3630    #[tokio::test]
3631    async fn uncached_child_layer_additions_o() {
3632        let (_dir, store) = make_cached_store();
3633        child_layer_additions_o(&store, true).await.unwrap();
3634    }
3635
3636    #[tokio::test]
3637    async fn cached_child_layer_removal_exists() {
3638        let (_dir, store) = make_cached_store();
3639        child_layer_removal_exists(&store, false).await.unwrap();
3640    }
3641
3642    #[tokio::test]
3643    async fn uncached_child_layer_removal_exists() {
3644        let (_dir, store) = make_cached_store();
3645        child_layer_removal_exists(&store, true).await.unwrap();
3646    }
3647
3648    #[tokio::test]
3649    async fn cached_child_layer_removals() {
3650        let (_dir, store) = make_cached_store();
3651        child_layer_removals(&store, false).await.unwrap();
3652    }
3653
3654    #[tokio::test]
3655    async fn uncached_child_layer_removals() {
3656        let (_dir, store) = make_cached_store();
3657        child_layer_removals(&store, true).await.unwrap();
3658    }
3659
3660    #[tokio::test]
3661    async fn cached_child_layer_removals_s() {
3662        let (_dir, store) = make_cached_store();
3663        child_layer_removals_s(&store, false).await.unwrap();
3664    }
3665
3666    #[tokio::test]
3667    async fn uncached_child_layer_removals_s() {
3668        let (_dir, store) = make_cached_store();
3669        child_layer_removals_s(&store, true).await.unwrap();
3670    }
3671
3672    #[tokio::test]
3673    async fn cached_child_layer_removals_sp() {
3674        let (_dir, store) = make_cached_store();
3675        child_layer_removals_sp(&store, false).await.unwrap();
3676    }
3677
3678    #[tokio::test]
3679    async fn uncached_child_layer_removals_sp() {
3680        let (_dir, store) = make_cached_store();
3681        child_layer_removals_sp(&store, true).await.unwrap();
3682    }
3683
3684    #[tokio::test]
3685    async fn cached_child_layer_removals_p() {
3686        let (_dir, store) = make_cached_store();
3687        child_layer_removals_p(&store, false).await.unwrap();
3688    }
3689
3690    #[tokio::test]
3691    async fn uncached_child_layer_removals_p() {
3692        let (_dir, store) = make_cached_store();
3693        child_layer_removals_p(&store, true).await.unwrap();
3694    }
3695
3696    #[tokio::test]
3697    async fn cached_child_layer_removals_o() {
3698        let (_dir, store) = make_cached_store();
3699        child_layer_removals_o(&store, false).await.unwrap();
3700    }
3701
3702    #[tokio::test]
3703    async fn uncached_child_layer_removals_o() {
3704        let (_dir, store) = make_cached_store();
3705        child_layer_removals_o(&store, true).await.unwrap();
3706    }
3707}