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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 if self.directory_exists(layer).await? {
1107 if self.layer_has_parent(layer).await? {
1108 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 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 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 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 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 if self.directory_exists(layer).await? {
1259 if self.layer_has_parent(layer).await? {
1260 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 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 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 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 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 if self.directory_exists(layer).await? {
1360 if self.layer_has_parent(layer).await? {
1361 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 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 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 layers_to_load.pop().unwrap();
1464
1465 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(); 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 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 return Ok(layer.name());
1732 }
1733
1734 if self.layer_has_rollup(layer.name()).await? {
1736 let rollup = self.read_rollup_file(layer.name()).await?;
1737 if !self.layer_has_parent(rollup).await? {
1739 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 return Ok(layer.name());
1761 }
1762
1763 if layer.parent_name() == Some(upto) {
1764 return Ok(layer.name());
1766 }
1767
1768 if self.layer_has_rollup(layer.name()).await? {
1770 let rollup = self.read_rollup_file(layer.name()).await?;
1771
1772 if upto == self.read_parent_file(rollup).await? {
1774 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 return Ok(layer.name());
1796 }
1797
1798 if layer.parent_name() == Some(upto) {
1799 return Ok(layer.name());
1801 }
1802
1803 if self.layer_has_rollup(layer.name()).await? {
1805 let rollup = self.read_rollup_file(layer.name()).await?;
1806
1807 if upto == self.read_parent_file(rollup).await? {
1809 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 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 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 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 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 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 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 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}