1use async_stream::try_stream;
2use futures::StreamExt;
3use futures::{Stream, stream::BoxStream};
4use std::collections::{HashMap, HashSet, hash_map};
5use tracing::{debug, trace};
6
7use super::Directory;
8use crate::{B3Digest, Node};
9
10#[derive(thiserror::Error, Debug, Eq, PartialEq)]
12pub enum OrderingError {
13 #[error("wrong size for digest {digest}, referenced with {referenced}, but got {actual}")]
14 WrongSize {
15 digest: B3Digest,
16 referenced: u64,
17 actual: u64,
18 },
19
20 #[error("unknown digest {digest} referenced for {path_component} in parent {parent_digest}")]
21 UnknownLTR {
22 digest: B3Digest,
23 parent_digest: B3Digest,
24 path_component: crate::PathComponent,
25 },
26
27 #[error("unexpected Directory with digest {0} encountered", directory.digest())]
28 Unexpected { directory: Directory },
29
30 #[error("some directories missing")]
31 DirectoriesMissing(HashSet<B3Digest>),
32
33 #[error("no directories received")]
34 EmptySet,
35}
36
37impl From<OrderingError> for crate::directoryservice::Error {
38 fn from(value: OrderingError) -> Self {
39 Self(Box::new(value))
40 }
41}
42
43pub trait OrderValidator {
45 fn try_accept(&mut self, directory: &Directory) -> Result<(), OrderingError>;
47 fn finalize(self) -> Result<(), OrderingError>;
49}
50
51pub struct RootToLeaves {
68 root_digest: B3Digest,
70
71 referenced_directories: HashMap<B3Digest, u64>,
73
74 pending_directories: HashSet<B3Digest>,
78
79 poison: bool,
81}
82
83impl RootToLeaves {
84 pub fn new_with_root_digest(root_digest: B3Digest) -> Self {
87 Self {
88 root_digest,
89 referenced_directories: HashMap::default(),
90 pending_directories: HashSet::from_iter([root_digest]),
91 poison: false,
92 }
93 }
94
95 pub fn would_accept(&self, digest: &B3Digest) -> bool {
103 assert!(!self.poison, "Snix bug: RootToLeavesValidator poisoned");
104 digest == &self.root_digest || self.referenced_directories.contains_key(digest)
105 }
106
107 fn introduce_children_of(&mut self, directory: &Directory) {
109 for (_name, node) in directory.nodes() {
110 if let Node::Directory { digest, size } = node {
111 if self
113 .referenced_directories
114 .insert(digest.to_owned(), *size)
115 .is_none()
116 {
117 self.pending_directories.insert(digest.to_owned());
118 }
119 }
120 }
121 }
122
123 pub fn validate_stream<'s, S>(
128 root_digest: B3Digest,
129 directories: S,
130 ) -> BoxStream<'s, Result<Directory, OrderingError>>
131 where
132 S: Stream<Item = Directory> + Send + 's,
133 {
134 let mut validator = RootToLeaves::new_with_root_digest(root_digest);
135 let mut directories = directories.boxed();
136
137 Box::pin(try_stream! {
138 while let Some(directory) = directories.next().await {
139 validator.try_accept(&directory)?;
140 yield directory;
141 }
142 validator.finalize()?;
143 })
144 }
145}
146
147impl OrderValidator for RootToLeaves {
148 #[tracing::instrument(level = "trace", skip_all, fields(directory.digest = %directory.digest(), directory.size = directory.size()), err)]
150 fn try_accept(&mut self, directory: &Directory) -> Result<(), OrderingError> {
151 assert!(!self.poison, "Snix bug: RootToLeavesValidator poisoned");
152
153 let size = directory.size();
154 let digest = directory.digest();
155
156 match self.referenced_directories.get(&digest) {
158 #[cfg(feature = "compat-accept-bigger-sizes")]
159 Some(size_referenced) if (size..=directory.size_max()).contains(size_referenced) => {
160 if !self.pending_directories.remove(&digest) {
161 debug!("directory received multiple times");
162 };
163
164 if *size_referenced != size {
165 debug!(directory.size_referenced=%size_referenced, "directory was referenced with a larger size (legacy size calculation)");
166 }
167
168 self.introduce_children_of(directory);
170 Ok(())
171 }
172 #[cfg(not(feature = "compat-accept-bigger-sizes"))]
173 Some(size_referenced) if size == *size_referenced => {
174 if !self.pending_directories.remove(&digest) {
175 debug!("directory received multiple times");
176 };
177
178 self.introduce_children_of(directory);
180 Ok(())
181 }
182 Some(size_referenced) => {
183 self.poison = true;
184 Err(OrderingError::WrongSize {
185 digest,
186 referenced: *size_referenced,
187 actual: size,
188 })
189 }
190 None if digest == self.root_digest => {
192 self.introduce_children_of(directory);
194 self.pending_directories.remove(&self.root_digest);
195 Ok(())
196 }
197 None => {
198 self.poison = true;
199 Err(OrderingError::Unexpected {
200 directory: directory.clone(),
201 })
202 }
203 }
204 }
205
206 #[tracing::instrument(level = "trace", skip_all, err)]
209 fn finalize(self) -> Result<(), OrderingError> {
210 match self.pending_directories.len() {
211 0 => Ok(()),
212 1 if self.pending_directories.iter().next().unwrap() == &self.root_digest => {
213 Err(OrderingError::EmptySet)
214 }
215 _ => Err(OrderingError::DirectoriesMissing(self.pending_directories)),
216 }
217 }
218}
219
220#[derive(Default)]
221pub struct LeavesToRoot {
229 #[cfg(feature = "compat-accept-bigger-sizes")]
230 accepted_directories: HashMap<B3Digest, std::ops::RangeInclusive<u64>>,
233
234 #[cfg(not(feature = "compat-accept-bigger-sizes"))]
235 accepted_directories: HashMap<B3Digest, u64>,
237
238 pending_directories: HashSet<B3Digest>,
241
242 #[cfg(debug_assertions)]
244 last_inserted_digest: Option<B3Digest>,
245
246 poison: bool,
249}
250
251impl LeavesToRoot {
252 pub fn new() -> Self {
253 Self {
254 accepted_directories: Default::default(),
255 pending_directories: Default::default(),
256 #[cfg(debug_assertions)]
257 last_inserted_digest: None,
258 poison: false,
259 }
260 }
261
262 pub fn validate_stream<'s, S>(directories: S) -> BoxStream<'s, Result<Directory, OrderingError>>
266 where
267 S: Stream<Item = Directory> + Send + 's,
268 {
269 let mut directories = directories.boxed();
270 let mut validator = Self::new();
271
272 Box::pin(try_stream! {
273 while let Some(directory) = directories.next().await {
274 validator.try_accept(&directory)?;
275 yield directory;
276 }
277
278 validator.finalize()?;
279 })
280 }
281}
282
283impl OrderValidator for LeavesToRoot {
284 #[tracing::instrument(level = "trace", skip_all, fields(directory.digest = %directory.digest(), directory.size = directory.size()), err)]
286 fn try_accept(&mut self, directory: &Directory) -> Result<(), OrderingError> {
287 assert!(!self.poison, "Snix bug: LeavesToRootValidator poisoned");
288
289 for (name, node) in directory.nodes() {
292 trace!(%name, ?node, "at node");
293 if let Node::Directory {
294 digest,
295 size: referenced_size,
296 } = node
297 {
298 match self.accepted_directories.get(digest) {
299 #[cfg(feature = "compat-accept-bigger-sizes")]
300 Some(size_range) if size_range.contains(referenced_size) => {
301 let minimal_size = size_range.start();
302 if referenced_size != minimal_size {
303 debug!(
304 directory.size_referenced=%referenced_size,
305 directory.size_referenced_minimal=%minimal_size,
306 "directory was referenced with a larger size (legacy size calculation)"
307 );
308 }
309 self.pending_directories.remove(digest);
310 }
311 #[cfg(not(feature = "compat-accept-bigger-sizes"))]
312 Some(size) if size == referenced_size => {
313 self.pending_directories.remove(digest);
314 }
315 Some(s) => {
316 self.poison = true;
317 Err(OrderingError::WrongSize {
318 digest: digest.to_owned(),
319 referenced: *referenced_size,
320 #[cfg(feature = "compat-accept-bigger-sizes")]
321 actual: *s.start(),
322 #[cfg(not(feature = "compat-accept-bigger-sizes"))]
323 actual: *s,
324 })?
325 }
326 None => {
327 self.poison = true;
328 Err(OrderingError::UnknownLTR {
329 digest: digest.to_owned(),
330 parent_digest: directory.digest(),
331 path_component: name.to_owned(),
332 })?
333 }
334 }
335 }
336 }
337
338 let directory_digest = directory.digest();
341 match self.accepted_directories.entry(directory_digest) {
342 hash_map::Entry::Occupied(_) => {
343 debug!("directory received multiple times");
344 }
345 hash_map::Entry::Vacant(entry) => {
346 #[cfg(feature = "compat-accept-bigger-sizes")]
347 entry.insert(directory.size()..=directory.size_max());
348 #[cfg(not(feature = "compat-accept-bigger-sizes"))]
349 entry.insert(directory.size());
350
351 #[cfg(debug_assertions)]
352 {
353 self.last_inserted_digest = Some(directory_digest)
354 }
355 self.pending_directories.insert(directory_digest);
356 }
357 }
358
359 Ok(())
360 }
361
362 #[tracing::instrument(level = "trace", skip_all, err)]
365 #[allow(unused_mut)]
366 fn finalize(mut self) -> Result<(), OrderingError> {
367 assert!(!self.poison, "Snix bug: LeavesToRootValidator poisoned");
368
369 if self.accepted_directories.is_empty() {
370 return Err(OrderingError::EmptySet);
371 }
372
373 if self.pending_directories.len() != 1 {
376 Err(OrderingError::DirectoriesMissing(
377 self.pending_directories.clone(),
378 ))?
379 }
380 #[cfg(debug_assertions)]
381 {
382 let last_inserted_digest = self
383 .last_inserted_digest
384 .expect("Snix bug: have dangling_directories, but no last_inserted_digest");
385 self.pending_directories
386 .get(&last_inserted_digest)
387 .expect("Snix bug: dangling directory is not last inserted one");
388 self.poison = true;
389 }
390
391 Ok(())
392 }
393}
394
395#[cfg(test)]
396mod tests {
397 use super::{LeavesToRoot, OrderValidator, RootToLeaves};
398 use crate::Node;
399 use crate::directoryservice::Directory;
400 use crate::fixtures::{
401 DIRECTORY_A, DIRECTORY_B, DIRECTORY_C, DIRECTORY_D, DIRECTORY_E, DIRECTORY_WITH_KEEP,
402 };
403 use futures::TryStreamExt;
404 use rstest::rstest;
405 use tracing_test::traced_test;
406
407 #[rstest]
408 #[case::empty_directory(&[&*DIRECTORY_A], false, false)]
410 #[case::simple_closure(&[&*DIRECTORY_A, &*DIRECTORY_B], false, false)]
412 #[case::same_child(&[&*DIRECTORY_A, &*DIRECTORY_A, &*DIRECTORY_C], false, false)]
415 #[case::same_child_dedup(&[&*DIRECTORY_A, &*DIRECTORY_C], false, false)]
417 #[case::unconnected_node(&[&*DIRECTORY_A, &*DIRECTORY_C, &*DIRECTORY_B], false, true)]
420 #[case::dangling_pointer(&[&*DIRECTORY_B], true, false)]
422 #[case::empty(&[], false, true)]
424 fn leaves_to_root(
425 #[case] directories_to_upload: &[&Directory],
426 #[case] exp_fail_upload_last: bool,
427 #[case] exp_fail_finalize: bool,
428 ) {
429 let mut validator = LeavesToRoot::default();
430 let mut it = directories_to_upload.iter().peekable();
431
432 while let Some(d) = it.next() {
433 if it.peek().is_none() && exp_fail_upload_last {
434 validator
435 .try_accept(d)
436 .expect_err("last try_accept to fail");
437 } else {
438 assert!(validator.try_accept(d).is_ok(), "try_accept to succeed");
439 }
440 }
441
442 if !exp_fail_upload_last {
443 if !exp_fail_finalize {
444 validator.finalize().expect("finalize to succeed");
445 } else {
446 let _ = validator.finalize();
447 }
448 }
449 }
450
451 #[rstest]
452 #[case::empty_directory(&[&*DIRECTORY_A], false)]
454 #[case::simple_closure(&[&*DIRECTORY_B, &*DIRECTORY_A], false)]
456 #[case::same_child_dedup(&[&*DIRECTORY_C, &*DIRECTORY_A], false)]
458 #[case::same_child_redundant(&[&*DIRECTORY_C, &*DIRECTORY_A, &*DIRECTORY_A], false)]
460 #[case::with_root_sent_twice(&[&*DIRECTORY_C, &*DIRECTORY_C, &*DIRECTORY_A], false)]
462 #[case::more_levels(&[&*DIRECTORY_E, &*DIRECTORY_D, &*DIRECTORY_A, &*DIRECTORY_B], false)]
464 #[case::unconnected_node(&[&*DIRECTORY_C, &*DIRECTORY_B], true)]
466 fn root_to_leaves(
467 #[case] directories_to_upload: &[&Directory],
468 #[case] exp_fail_upload_last: bool,
469 ) {
470 let root_digest = directories_to_upload[0].digest();
471 let mut validator = RootToLeaves::new_with_root_digest(root_digest);
472 let mut it = directories_to_upload.iter().peekable();
473
474 while let Some(d) = it.next() {
475 if it.peek().is_none() && exp_fail_upload_last {
476 assert!(
477 !validator.would_accept(&d.digest()),
478 "would_accept not expected to accept last failing element"
479 );
480
481 validator
482 .try_accept(d)
483 .expect_err("last try_accept to fail");
484 } else {
485 assert!(
486 validator.would_accept(&d.digest()),
487 "would_accept expected to accept directory"
488 );
489 assert!(validator.try_accept(d).is_ok(), "try_accept to succeed");
490 }
491 }
492
493 if !exp_fail_upload_last {
494 validator.finalize().expect("finalize to succeed");
495 }
496 }
497
498 fn legacy_size_dirs() -> Vec<Directory> {
501 let a = DIRECTORY_WITH_KEEP.to_owned();
502 assert_eq!(a.size(), 1);
503 #[cfg(feature = "compat-accept-bigger-sizes")]
504 assert_eq!(a.size_max(), 2);
505
506 let b = crate::Directory::try_from_iter([
508 (
509 "symlink".try_into().unwrap(),
510 Node::Symlink {
511 target: "somewhereelse".try_into().unwrap(),
512 },
513 ),
514 (
515 "dir".try_into().unwrap(),
516 Node::Directory {
517 digest: DIRECTORY_WITH_KEEP.digest(),
518 size: 1,
519 },
520 ),
521 ])
522 .unwrap();
523
524 assert_eq!(3, b.size());
525 #[cfg(feature = "compat-accept-bigger-sizes")]
526 assert_eq!(5, b.size_max());
527 let b_size = 1 + 2 + 1; let root = crate::Directory::try_from_iter([
530 (
531 "a".try_into().unwrap(),
532 Node::Directory {
533 digest: a.digest(),
534 size: 2,
536 },
537 ),
538 (
539 "b".try_into().unwrap(),
540 Node::Directory {
541 digest: b.digest(),
542 size: b_size,
543 },
544 ),
545 ])
546 .unwrap();
547
548 vec![root, b, a]
549 }
550
551 #[test]
552 #[traced_test]
553 fn root_to_leaves_legacy_size() {
555 let dirs = legacy_size_dirs();
556
557 let mut validator = RootToLeaves::new_with_root_digest(dirs[0].digest());
558 validator.try_accept(&dirs[0]).expect("to accept root");
559
560 if cfg!(feature = "compat-accept-bigger-sizes") {
561 validator.try_accept(&dirs[1]).expect("to accept b");
562 validator.try_accept(&dirs[2]).expect("to accept leaf a");
563 validator.finalize().expect("to finalize");
564
565 assert!(logs_contain("legacy size calculation"));
566 } else {
567 validator
568 .try_accept(&dirs[1])
569 .expect_err("to reject b due to wrong size used in root");
570 }
571 }
572
573 #[test]
574 #[traced_test]
575 fn leaves_to_root_legacy_size() {
577 let dirs = legacy_size_dirs();
578
579 let mut validator = LeavesToRoot::new();
580 validator.try_accept(&dirs[2]).expect("to accept leaf a");
581 validator.try_accept(&dirs[1]).expect("to accept b");
582
583 if cfg!(feature = "compat-accept-bigger-sizes") {
584 validator.try_accept(&dirs[0]).expect("to accept root");
585 validator.finalize().expect("to finalize");
586
587 assert!(logs_contain("legacy size calculation"));
588 } else {
589 validator
590 .try_accept(&dirs[0])
591 .expect_err("to reject root due to referring to a with wrong size");
592 }
593 }
594
595 #[test]
596 fn reject_too_small_size() {
599 assert_eq!(1, DIRECTORY_B.size());
600
601 let root = Directory::try_from_iter([(
603 "b".try_into().unwrap(),
604 Node::Directory {
605 digest: DIRECTORY_B.digest(),
606 size: 0,
607 },
608 )])
609 .unwrap();
610
611 let mut validator = RootToLeaves::new_with_root_digest(root.digest());
612 validator.try_accept(&root).expect("should accept root");
613 validator
614 .try_accept(&DIRECTORY_B)
615 .expect_err("should reject B due to wrong size");
616
617 let mut validator = LeavesToRoot::new();
618 validator.try_accept(&DIRECTORY_A).expect("should accept A");
619 validator.try_accept(&DIRECTORY_B).expect("should accept B");
620 validator
621 .try_accept(&root)
622 .expect_err("should reject root due to referring by wrong size");
623 }
624
625 #[test]
626 fn reject_too_big_size() {
629 #[cfg(feature = "compat-accept-bigger-sizes")]
630 assert_eq!(2, DIRECTORY_B.size_max());
631
632 let root = Directory::try_from_iter([(
634 "b".try_into().unwrap(),
635 Node::Directory {
636 digest: DIRECTORY_B.digest(),
637 size: 3,
638 },
639 )])
640 .unwrap();
641
642 let mut validator = RootToLeaves::new_with_root_digest(root.digest());
643 validator.try_accept(&root).expect("should accept root");
644 validator
645 .try_accept(&DIRECTORY_B)
646 .expect_err("should reject B due to wrong size");
647
648 let mut validator = LeavesToRoot::new();
649 validator.try_accept(&DIRECTORY_A).expect("should accept A");
650 validator.try_accept(&DIRECTORY_B).expect("should accept B");
651 validator
652 .try_accept(&root)
653 .expect_err("should reject root due to referring by wrong size");
654 }
655
656 #[test]
657 fn root_to_leaves_root_mismatch() {
659 let mut validator = RootToLeaves::new_with_root_digest(DIRECTORY_A.digest());
660
661 validator
662 .try_accept(&DIRECTORY_B)
663 .expect_err("shouldn't accept wrong first directory");
664 validator.finalize().expect_err("expect finalize to fail");
665 }
666
667 #[tokio::test]
668 async fn root_to_leaves_stream() {
669 let directories_to_upload = vec![
670 DIRECTORY_E.to_owned(),
671 DIRECTORY_D.to_owned(),
672 DIRECTORY_A.to_owned(),
673 DIRECTORY_B.to_owned(),
674 ];
675 let root_digest = directories_to_upload[0].digest();
676
677 let validated_stream = RootToLeaves::validate_stream(
678 root_digest,
679 futures::stream::iter(directories_to_upload.iter().map(|d| (*d).to_owned())),
680 );
681
682 let validated_directories: Vec<Directory> = validated_stream
683 .try_collect()
684 .await
685 .expect("stream to collect successfully");
686
687 assert_eq!(directories_to_upload, validated_directories);
688
689 RootToLeaves::validate_stream(root_digest, futures::stream::empty())
690 .try_collect::<Vec<_>>()
691 .await
692 .expect_err("an empty stream to fail");
693 }
694}