Skip to main content

xet_data/progress_tracking/
upload_tracking.rs

1use std::collections::BTreeSet;
2use std::collections::hash_map::Entry as HashMapEntry;
3use std::mem::take;
4use std::sync::atomic::Ordering;
5use std::sync::{Arc, Mutex};
6
7use more_asserts::{debug_assert_ge, debug_assert_le};
8#[cfg(target_family = "wasm")]
9use tokio_with_wasm::alias as tokio;
10use xet_core_structures::MerkleHashMap;
11use xet_core_structures::merklehash::MerkleHash;
12use xet_runtime::utils::UniqueId;
13
14use super::progress_types::{GroupProgress, ItemProgressUpdater};
15
16pub struct FileXorbDependency {
17    pub file_id: u64,
18    pub xorb_hash: MerkleHash,
19    pub n_bytes: u64,
20    pub is_external: bool,
21}
22
23/// A type with with to track a File ID; reporting is done by Arc<str>, but
24/// this ensures the bookkeeping is correct across duplicates and speeds up the
25/// updates.
26pub type CompletionTrackerFileId = u64;
27
28/// Keeps track of which files depend on a given xorb.
29#[derive(Default)]
30struct XorbDependency {
31    /// List of file indices that need this xorb.
32    file_indices: BTreeSet<usize>,
33
34    /// Number of bytes completed so far
35    completed_bytes: u64,
36
37    /// Number of bytes in that xorb.
38    xorb_size: u64,
39
40    /// True if the xorb has already been updated successfully.
41    is_completed: bool,
42}
43
44#[derive(Default, Debug)]
45struct XorbPartCompletionStats {
46    completed_bytes: u64,
47    n_bytes: u64,
48}
49
50/// Represents a file that depends on one or more xorbs.
51struct FileDependency {
52    tracking_id: UniqueId,
53    updater: Arc<ItemProgressUpdater>,
54    name: Arc<str>,
55    total_bytes: u64,
56    is_final_size_known: bool,
57    completed_bytes: u64,
58    remaining_xorbs_parts: MerkleHashMap<XorbPartCompletionStats>,
59}
60
61/// Tracks all files and all xorbs, allowing you to register file
62/// dependencies on xorbs and then mark xorbs as completed when they
63/// are fully uploaded.
64#[derive(Default)]
65struct CompletionTrackerImpl {
66    files: Vec<FileDependency>,
67    xorbs: MerkleHashMap<XorbDependency>,
68
69    total_upload_bytes: u64,
70    total_upload_bytes_completed: u64,
71
72    total_bytes: u64,
73    total_bytes_completed: u64,
74}
75
76pub struct CompletionTracker {
77    inner: Mutex<CompletionTrackerImpl>,
78    group: Arc<GroupProgress>,
79}
80
81impl CompletionTrackerImpl {
82    fn register_new_file(
83        &mut self,
84        updater: Arc<ItemProgressUpdater>,
85        n_bytes: Option<u64>,
86    ) -> CompletionTrackerFileId {
87        let (total_bytes, is_final_size_known) = match n_bytes {
88            Some(size) => (size, true),
89            None => (0, false),
90        };
91
92        updater.update_item_size(total_bytes, n_bytes.is_some());
93
94        let file_id = self.files.len() as CompletionTrackerFileId;
95        let tracking_id = updater.item().id;
96        let name = updater.item().name.clone();
97
98        let file_dependency = FileDependency {
99            tracking_id,
100            updater,
101            name,
102            total_bytes,
103            is_final_size_known,
104            completed_bytes: 0,
105            remaining_xorbs_parts: MerkleHashMap::new(),
106        };
107
108        self.files.push(file_dependency);
109        self.total_bytes += total_bytes;
110
111        file_id
112    }
113
114    fn increment_file_size(&mut self, file_id: CompletionTrackerFileId, size_increment: u64) {
115        let file_entry = &mut self.files[file_id as usize];
116
117        if file_entry.is_final_size_known {
118            return;
119        }
120
121        file_entry.total_bytes += size_increment;
122        self.total_bytes += size_increment;
123
124        debug_assert_ge!(file_entry.total_bytes, file_entry.completed_bytes);
125        debug_assert_ge!(self.total_bytes, self.total_bytes_completed);
126
127        file_entry.updater.update_item_size(file_entry.total_bytes, false);
128    }
129
130    fn register_dependencies(&mut self, dependencies: &[FileXorbDependency]) {
131        let mut file_bytes_processed = 0;
132
133        for dep in dependencies {
134            let file_entry = &mut self.files[dep.file_id as usize];
135
136            if dep.is_external {
137                file_entry.completed_bytes += dep.n_bytes;
138                debug_assert_le!(file_entry.completed_bytes, file_entry.total_bytes);
139
140                file_entry.updater.report_bytes_completed(dep.n_bytes);
141                file_bytes_processed += dep.n_bytes;
142            } else {
143                debug_assert_ne!(dep.xorb_hash, MerkleHash::marker());
144
145                let entry = self.xorbs.entry(dep.xorb_hash).or_default();
146
147                if entry.is_completed {
148                    file_entry.completed_bytes += dep.n_bytes;
149                    debug_assert_le!(file_entry.completed_bytes, file_entry.total_bytes);
150
151                    file_entry.updater.report_bytes_completed(dep.n_bytes);
152                    file_bytes_processed += dep.n_bytes;
153                } else {
154                    entry.file_indices.insert(dep.file_id as usize);
155                    file_entry.remaining_xorbs_parts.entry(dep.xorb_hash).or_default().n_bytes += dep.n_bytes;
156                }
157            }
158        }
159
160        self.total_bytes_completed += file_bytes_processed;
161        debug_assert_le!(self.total_bytes_completed, self.total_bytes);
162    }
163
164    fn register_new_xorb(&mut self, group: &Arc<GroupProgress>, xorb_hash: MerkleHash, xorb_size: u64) -> bool {
165        match self.xorbs.entry(xorb_hash) {
166            HashMapEntry::Occupied(mut occupied_entry) => {
167                let entry = occupied_entry.get_mut();
168                if entry.xorb_size == 0 {
169                    entry.xorb_size = xorb_size;
170                    self.total_upload_bytes += xorb_size;
171                    group.total_transfer_bytes.fetch_add(xorb_size, Ordering::Release);
172                    true
173                } else {
174                    debug_assert_eq!(entry.xorb_size, xorb_size);
175                    false
176                }
177            },
178            HashMapEntry::Vacant(vacant_entry) => {
179                vacant_entry.insert(XorbDependency {
180                    file_indices: Default::default(),
181                    xorb_size,
182                    completed_bytes: 0,
183                    is_completed: false,
184                });
185
186                self.total_upload_bytes += xorb_size;
187                group.total_transfer_bytes.fetch_add(xorb_size, Ordering::Release);
188                true
189            },
190        }
191    }
192
193    fn register_xorb_upload_completion(&mut self, group: &Arc<GroupProgress>, xorb_hash: MerkleHash) {
194        let (file_indices, byte_completion_increment) = {
195            let entry = self.xorbs.entry(xorb_hash).or_default();
196
197            if entry.is_completed {
198                return;
199            }
200
201            let new_byte_increment = entry.xorb_size - entry.completed_bytes;
202            entry.is_completed = true;
203
204            (take(&mut entry.file_indices), new_byte_increment)
205        };
206
207        let mut file_bytes_processed = 0;
208
209        for file_id in file_indices {
210            let file_entry = &mut self.files[file_id];
211
212            debug_assert!(file_entry.remaining_xorbs_parts.contains_key(&xorb_hash));
213
214            let xorb_part = file_entry.remaining_xorbs_parts.remove(&xorb_hash).unwrap_or_default();
215            debug_assert_le!(xorb_part.completed_bytes, xorb_part.n_bytes);
216
217            let n_bytes_remaining = xorb_part.n_bytes - xorb_part.completed_bytes;
218
219            if n_bytes_remaining > 0 {
220                file_entry.completed_bytes += n_bytes_remaining;
221                file_entry.updater.report_bytes_completed(n_bytes_remaining);
222                file_bytes_processed += n_bytes_remaining;
223            }
224        }
225
226        debug_assert_le!(self.total_upload_bytes_completed + byte_completion_increment, self.total_upload_bytes);
227        self.total_upload_bytes_completed += byte_completion_increment;
228        group
229            .total_transfer_bytes_completed
230            .fetch_add(byte_completion_increment, Ordering::Release);
231
232        self.total_bytes_completed += file_bytes_processed;
233        debug_assert_le!(self.total_bytes_completed, self.total_bytes);
234    }
235
236    fn register_xorb_upload_progress(
237        &mut self,
238        group: &Arc<GroupProgress>,
239        xorb_hash: MerkleHash,
240        new_byte_progress: u64,
241        check_ordering: bool,
242    ) {
243        debug_assert!(self.xorbs.contains_key(&xorb_hash));
244
245        let entry = self.xorbs.entry(xorb_hash).or_default();
246
247        if !check_ordering && entry.is_completed {
248            return;
249        }
250
251        debug_assert!(!entry.is_completed);
252        debug_assert_le!(entry.completed_bytes + new_byte_progress, entry.xorb_size);
253
254        entry.completed_bytes += new_byte_progress;
255
256        let new_completion_ratio = (entry.completed_bytes as f64) / (entry.xorb_size as f64);
257
258        let mut file_bytes_processed = 0;
259
260        for &file_id in entry.file_indices.iter() {
261            let file_entry = &mut self.files[file_id];
262
263            debug_assert!(file_entry.remaining_xorbs_parts.contains_key(&xorb_hash));
264
265            let incremental_update = 'update: {
266                let Some(xorb_part) = file_entry.remaining_xorbs_parts.get_mut(&xorb_hash) else {
267                    break 'update 0;
268                };
269                debug_assert_le!(xorb_part.completed_bytes, xorb_part.n_bytes);
270
271                let new_completion_bytes = ((xorb_part.n_bytes as f64) * new_completion_ratio).floor() as u64;
272
273                debug_assert_ge!(new_completion_bytes, xorb_part.completed_bytes);
274
275                let incremental_update = new_completion_bytes.saturating_sub(xorb_part.completed_bytes);
276                xorb_part.completed_bytes += incremental_update;
277
278                debug_assert_le!(xorb_part.completed_bytes, xorb_part.n_bytes);
279
280                incremental_update
281            };
282
283            if incremental_update != 0 {
284                file_entry.completed_bytes += incremental_update;
285                file_entry.updater.report_bytes_completed(incremental_update);
286                file_bytes_processed += incremental_update;
287            }
288        }
289
290        self.total_upload_bytes_completed += new_byte_progress;
291        debug_assert_le!(self.total_upload_bytes_completed, self.total_upload_bytes);
292
293        group
294            .total_transfer_bytes_completed
295            .fetch_add(new_byte_progress, Ordering::Release);
296
297        self.total_bytes_completed += file_bytes_processed;
298        debug_assert_le!(self.total_bytes_completed, self.total_bytes);
299    }
300
301    fn status(&self) -> (u64, u64) {
302        let (mut sum_completed, mut sum_total) = (0, 0);
303        for file in &self.files {
304            sum_completed += file.completed_bytes;
305            sum_total += file.total_bytes;
306        }
307        (sum_completed, sum_total)
308    }
309
310    fn is_complete(&self) -> bool {
311        let (done, total) = self.status();
312
313        #[cfg(debug_assertions)]
314        {
315            if done == total {
316                self.assert_complete();
317            }
318        }
319
320        done == total
321    }
322
323    fn assert_complete(&self) {
324        for (idx, file) in self.files.iter().enumerate() {
325            assert_eq!(
326                file.completed_bytes, file.total_bytes,
327                "File #{} ({}, {}) is not fully completed: {}/{} bytes",
328                idx, file.name, file.tracking_id, file.completed_bytes, file.total_bytes
329            );
330            assert!(
331                file.remaining_xorbs_parts.is_empty(),
332                "File #{} ({}) still has uncompleted xorb parts: {:?}",
333                idx,
334                file.name,
335                file.remaining_xorbs_parts
336            );
337        }
338
339        for (hash, xorb_dep) in self.xorbs.iter() {
340            assert!(xorb_dep.is_completed, "Xorb {hash:?} is not marked completed.");
341            assert!(
342                xorb_dep.file_indices.is_empty(),
343                "Xorb {:?} still has file references: {:?}",
344                hash,
345                xorb_dep.file_indices
346            );
347        }
348    }
349}
350
351impl CompletionTracker {
352    pub fn new(group: Arc<GroupProgress>) -> Self {
353        Self {
354            inner: Mutex::new(CompletionTrackerImpl::default()),
355            group,
356        }
357    }
358
359    pub fn register_new_file(
360        &self,
361        updater: Arc<ItemProgressUpdater>,
362        n_bytes: Option<u64>,
363    ) -> CompletionTrackerFileId {
364        let mut update_lock = self.inner.lock().unwrap();
365        update_lock.register_new_file(updater, n_bytes)
366    }
367
368    pub fn increment_file_size(&self, file_id: CompletionTrackerFileId, size_increment: u64) {
369        let mut update_lock = self.inner.lock().unwrap();
370        update_lock.increment_file_size(file_id, size_increment);
371    }
372
373    pub fn register_new_xorb(&self, xorb_hash: MerkleHash, xorb_size: u64) -> bool {
374        let mut update_lock = self.inner.lock().unwrap();
375        update_lock.register_new_xorb(&self.group, xorb_hash, xorb_size)
376    }
377
378    pub fn register_dependencies(&self, dependencies: &[FileXorbDependency]) {
379        let mut update_lock = self.inner.lock().unwrap();
380        update_lock.register_dependencies(dependencies);
381    }
382
383    pub fn register_xorb_upload_completion(&self, xorb_hash: MerkleHash) {
384        let mut update_lock = self.inner.lock().unwrap();
385        update_lock.register_xorb_upload_completion(&self.group, xorb_hash);
386    }
387
388    pub fn register_xorb_upload_progress(&self, xorb_hash: MerkleHash, new_byte_progress: u64) {
389        self.register_xorb_upload_progress_impl(xorb_hash, new_byte_progress, true);
390    }
391
392    pub fn register_xorb_upload_progress_background(self: Arc<Self>, xorb_hash: MerkleHash, new_byte_progress: u64) {
393        tokio::spawn(async move {
394            self.register_xorb_upload_progress_impl(xorb_hash, new_byte_progress, false);
395        });
396    }
397
398    fn register_xorb_upload_progress_impl(&self, xorb_hash: MerkleHash, new_byte_progress: u64, check_ordering: bool) {
399        let mut update_lock = self.inner.lock().unwrap();
400        update_lock.register_xorb_upload_progress(&self.group, xorb_hash, new_byte_progress, check_ordering);
401    }
402
403    pub fn status(&self) -> (u64, u64) {
404        self.inner.lock().unwrap().status()
405    }
406
407    pub fn is_complete(&self) -> bool {
408        self.inner.lock().unwrap().is_complete()
409    }
410
411    pub fn assert_complete(&self) {
412        self.inner.lock().unwrap().assert_complete();
413    }
414}
415
416#[cfg(test)]
417mod tests {
418    use xet_core_structures::merklehash::MerkleHash;
419
420    use super::*;
421
422    #[test]
423    fn test_status_and_is_complete() {
424        let group = GroupProgress::new();
425        let tracker = CompletionTracker::new(group.clone());
426
427        let updater_a = group.new_item(UniqueId::new(), "fileA");
428        let file_a = tracker.register_new_file(updater_a, Some(100));
429
430        let updater_b = group.new_item(UniqueId::new(), "fileB");
431        let file_b = tracker.register_new_file(updater_b, Some(50));
432
433        let (done, total) = tracker.status();
434        assert_eq!(done, 0);
435        assert_eq!(total, 150);
436        assert!(!tracker.is_complete());
437
438        let x = MerkleHash::random_from_seed(1);
439        tracker.register_dependencies(&[FileXorbDependency {
440            file_id: file_a,
441            xorb_hash: x,
442            n_bytes: 100,
443            is_external: true,
444        }]);
445
446        let (done, total) = tracker.status();
447        assert_eq!(done, 100);
448        assert_eq!(total, 150);
449        assert!(!tracker.is_complete());
450
451        let y = MerkleHash::random_from_seed(2);
452        tracker.register_dependencies(&[FileXorbDependency {
453            file_id: file_b,
454            xorb_hash: y,
455            n_bytes: 50,
456            is_external: false,
457        }]);
458
459        let (done, total) = tracker.status();
460        assert_eq!(done, 100);
461        assert_eq!(total, 150);
462
463        tracker.register_new_xorb(y, 50);
464        tracker.register_xorb_upload_completion(y);
465
466        let (done, total) = tracker.status();
467        assert_eq!(done, 150);
468        assert_eq!(total, 150);
469        assert!(tracker.is_complete());
470
471        tracker.assert_complete();
472        group.assert_complete();
473    }
474
475    #[test]
476    fn test_multiple_files_one_shared_xorb() {
477        let group = GroupProgress::new();
478        let tracker = CompletionTracker::new(group.clone());
479
480        let updater_a = group.new_item(UniqueId::new(), "fileA");
481        let file_a = tracker.register_new_file(updater_a, Some(200));
482
483        let updater_b = group.new_item(UniqueId::new(), "fileB");
484        let file_b = tracker.register_new_file(updater_b, Some(300));
485
486        let (done, total) = tracker.status();
487        assert_eq!(done, 0);
488        assert_eq!(total, 500);
489
490        let xhash = MerkleHash::random_from_seed(1);
491
492        tracker.register_new_xorb(xhash, 1000);
493
494        tracker.register_dependencies(&[
495            FileXorbDependency {
496                file_id: file_a,
497                xorb_hash: xhash,
498                n_bytes: 100,
499                is_external: false,
500            },
501            FileXorbDependency {
502                file_id: file_b,
503                xorb_hash: xhash,
504                n_bytes: 200,
505                is_external: true,
506            },
507        ]);
508
509        let (done, total) = tracker.status();
510        assert_eq!(done, 200);
511        assert_eq!(total, 500);
512        assert!(!tracker.is_complete());
513
514        tracker.register_xorb_upload_completion(xhash);
515
516        let (done, total) = tracker.status();
517        assert_eq!(done, 300);
518        assert_eq!(total, 500);
519
520        let x2 = MerkleHash::random_from_seed(2);
521
522        tracker.register_new_xorb(x2, 1000);
523
524        tracker.register_dependencies(&[FileXorbDependency {
525            file_id: file_a,
526            xorb_hash: x2,
527            n_bytes: 100,
528            is_external: true,
529        }]);
530
531        let (done, total) = tracker.status();
532        assert_eq!(done, 400);
533        assert_eq!(total, 500);
534
535        tracker.register_dependencies(&[FileXorbDependency {
536            file_id: file_b,
537            xorb_hash: x2,
538            n_bytes: 100,
539            is_external: false,
540        }]);
541
542        let (done, total) = tracker.status();
543        assert_eq!(done, 400);
544        assert_eq!(total, 500);
545        assert!(!tracker.is_complete());
546
547        tracker.register_xorb_upload_completion(x2);
548        let (done, total) = tracker.status();
549        assert_eq!(done, 500);
550        assert_eq!(total, 500);
551        assert!(tracker.is_complete());
552
553        tracker.assert_complete();
554        group.assert_complete();
555    }
556
557    #[test]
558    fn test_single_file_multiple_xorbs() {
559        let group = GroupProgress::new();
560        let tracker = CompletionTracker::new(group.clone());
561
562        let updater = group.new_item(UniqueId::new(), "bigFile");
563        let f = tracker.register_new_file(updater, Some(300));
564
565        let x1 = MerkleHash::random_from_seed(1);
566        let x2 = MerkleHash::random_from_seed(2);
567        let x3 = MerkleHash::random_from_seed(3);
568
569        tracker.register_new_xorb(x1, 100);
570        tracker.register_new_xorb(x3, 100);
571
572        tracker.register_dependencies(&[
573            FileXorbDependency {
574                file_id: f,
575                xorb_hash: x1,
576                n_bytes: 100,
577                is_external: false,
578            },
579            FileXorbDependency {
580                file_id: f,
581                xorb_hash: x2,
582                n_bytes: 100,
583                is_external: true,
584            },
585            FileXorbDependency {
586                file_id: f,
587                xorb_hash: x3,
588                n_bytes: 100,
589                is_external: false,
590            },
591        ]);
592
593        let (done, total) = tracker.status();
594        assert_eq!(done, 100);
595        assert_eq!(total, 300);
596        assert!(!tracker.is_complete());
597
598        tracker.register_xorb_upload_completion(x1);
599        let (done, total) = tracker.status();
600        assert_eq!(done, 200);
601        assert_eq!(total, 300);
602        assert!(!tracker.is_complete());
603
604        tracker.register_xorb_upload_completion(x3);
605        let (done, total) = tracker.status();
606        assert_eq!(done, 300);
607        assert_eq!(total, 300);
608        assert!(tracker.is_complete());
609
610        tracker.assert_complete();
611        group.assert_complete();
612    }
613
614    #[test]
615    fn test_xorb_completed_before_dependencies() {
616        let group = GroupProgress::new();
617        let tracker = CompletionTracker::new(group.clone());
618
619        let updater = group.new_item(UniqueId::new(), "lateFile");
620        let file_id = tracker.register_new_file(updater, Some(50));
621
622        let x = MerkleHash::random_from_seed(999);
623        tracker.register_new_xorb(x, 1000);
624
625        tracker.register_xorb_upload_completion(x);
626
627        tracker.register_dependencies(&[FileXorbDependency {
628            file_id,
629            xorb_hash: x,
630            n_bytes: 50,
631            is_external: false,
632        }]);
633
634        let (done, total) = tracker.status();
635        assert_eq!(done, 50);
636        assert_eq!(total, 50);
637        assert!(tracker.is_complete());
638
639        tracker.assert_complete();
640        group.assert_complete();
641    }
642
643    #[test]
644    fn test_contradictory_logic_with_completed_xorb() {
645        let group = GroupProgress::new();
646        let tracker = CompletionTracker::new(group.clone());
647
648        let updater = group.new_item(UniqueId::new(), "someFile");
649        let file_id = tracker.register_new_file(updater, Some(100));
650        let x = MerkleHash::random_from_seed(123);
651
652        tracker.register_new_xorb(x, 1000);
653
654        tracker.register_xorb_upload_completion(x);
655
656        tracker.register_dependencies(&[FileXorbDependency {
657            file_id,
658            xorb_hash: x,
659            n_bytes: 100,
660            is_external: false,
661        }]);
662
663        let (done, total) = tracker.status();
664        assert_eq!(done, 100);
665        assert_eq!(total, 100);
666        assert!(tracker.is_complete());
667
668        tracker.assert_complete();
669        group.assert_complete();
670    }
671
672    #[test]
673    fn test_increment_file_size_basic() {
674        let group = GroupProgress::new();
675        let tracker = CompletionTracker::new(group.clone());
676
677        let updater = group.new_item(UniqueId::new(), "growingFile");
678        let file_id = tracker.register_new_file(updater, None);
679
680        let (done, total) = tracker.status();
681        assert_eq!(done, 0);
682        assert_eq!(total, 0);
683
684        tracker.increment_file_size(file_id, 100);
685        let (done, total) = tracker.status();
686        assert_eq!(done, 0);
687        assert_eq!(total, 100);
688
689        tracker.increment_file_size(file_id, 150);
690        let (done, total) = tracker.status();
691        assert_eq!(done, 0);
692        assert_eq!(total, 250);
693
694        tracker.increment_file_size(file_id, 50);
695        let (done, total) = tracker.status();
696        assert_eq!(done, 0);
697        assert_eq!(total, 300);
698
699        let x = MerkleHash::random_from_seed(1);
700        tracker.register_dependencies(&[FileXorbDependency {
701            file_id,
702            xorb_hash: x,
703            n_bytes: 300,
704            is_external: true,
705        }]);
706
707        let (done, total) = tracker.status();
708        assert_eq!(done, 300);
709        assert_eq!(total, 300);
710        assert!(tracker.is_complete());
711
712        tracker.assert_complete();
713        group.assert_complete();
714    }
715
716    #[test]
717    fn test_increment_file_size_with_xorb_uploads() {
718        let group = GroupProgress::new();
719        let tracker = CompletionTracker::new(group.clone());
720
721        let updater = group.new_item(UniqueId::new(), "streamFile");
722        let file_id = tracker.register_new_file(updater, None);
723
724        let x1 = MerkleHash::random_from_seed(10);
725        let x2 = MerkleHash::random_from_seed(20);
726
727        tracker.register_new_xorb(x1, 500);
728        tracker.register_new_xorb(x2, 500);
729
730        tracker.increment_file_size(file_id, 200);
731        tracker.register_dependencies(&[FileXorbDependency {
732            file_id,
733            xorb_hash: x1,
734            n_bytes: 200,
735            is_external: false,
736        }]);
737
738        let (done, total) = tracker.status();
739        assert_eq!(done, 0);
740        assert_eq!(total, 200);
741
742        tracker.register_xorb_upload_completion(x1);
743        let (done, total) = tracker.status();
744        assert_eq!(done, 200);
745        assert_eq!(total, 200);
746
747        tracker.increment_file_size(file_id, 300);
748        tracker.register_dependencies(&[FileXorbDependency {
749            file_id,
750            xorb_hash: x2,
751            n_bytes: 300,
752            is_external: false,
753        }]);
754
755        let (done, total) = tracker.status();
756        assert_eq!(done, 200);
757        assert_eq!(total, 500);
758
759        tracker.register_xorb_upload_completion(x2);
760        let (done, total) = tracker.status();
761        assert_eq!(done, 500);
762        assert_eq!(total, 500);
763        assert!(tracker.is_complete());
764
765        tracker.assert_complete();
766        group.assert_complete();
767    }
768
769    #[test]
770    fn test_increment_file_size_mixed_known_unknown() {
771        let group = GroupProgress::new();
772        let tracker = CompletionTracker::new(group.clone());
773
774        let updater_a = group.new_item(UniqueId::new(), "fileA");
775        let file_a = tracker.register_new_file(updater_a, Some(100));
776
777        let updater_b = group.new_item(UniqueId::new(), "fileB");
778        let file_b = tracker.register_new_file(updater_b, None);
779
780        let (done, total) = tracker.status();
781        assert_eq!(done, 0);
782        assert_eq!(total, 100);
783
784        let xa = MerkleHash::random_from_seed(1);
785        tracker.register_dependencies(&[FileXorbDependency {
786            file_id: file_a,
787            xorb_hash: xa,
788            n_bytes: 100,
789            is_external: true,
790        }]);
791
792        let (done, total) = tracker.status();
793        assert_eq!(done, 100);
794        assert_eq!(total, 100);
795
796        tracker.increment_file_size(file_b, 200);
797        let (done, total) = tracker.status();
798        assert_eq!(done, 100);
799        assert_eq!(total, 300);
800
801        let xb = MerkleHash::random_from_seed(2);
802        tracker.register_new_xorb(xb, 200);
803        tracker.register_dependencies(&[FileXorbDependency {
804            file_id: file_b,
805            xorb_hash: xb,
806            n_bytes: 200,
807            is_external: false,
808        }]);
809
810        tracker.register_xorb_upload_completion(xb);
811
812        let (done, total) = tracker.status();
813        assert_eq!(done, 300);
814        assert_eq!(total, 300);
815        assert!(tracker.is_complete());
816
817        tracker.assert_complete();
818        group.assert_complete();
819    }
820
821    #[test]
822    fn test_increment_file_size_ignored_when_already_final() {
823        let group = GroupProgress::new();
824        let tracker = CompletionTracker::new(group.clone());
825
826        let updater = group.new_item(UniqueId::new(), "fixedFile");
827        let file_id = tracker.register_new_file(updater, Some(100));
828
829        tracker.increment_file_size(file_id, 999);
830        let (_, total) = tracker.status();
831        assert_eq!(total, 100);
832
833        let x = MerkleHash::random_from_seed(1);
834        tracker.register_dependencies(&[FileXorbDependency {
835            file_id,
836            xorb_hash: x,
837            n_bytes: 100,
838            is_external: true,
839        }]);
840
841        assert!(tracker.is_complete());
842        tracker.assert_complete();
843        group.assert_complete();
844    }
845
846    #[test]
847    fn test_increment_file_size_with_partial_xorb_progress() {
848        let group = GroupProgress::new();
849        let tracker = CompletionTracker::new(group.clone());
850
851        let updater = group.new_item(UniqueId::new(), "partialFile");
852        let file_id = tracker.register_new_file(updater, None);
853
854        let x = MerkleHash::random_from_seed(42);
855        tracker.register_new_xorb(x, 1000);
856
857        tracker.increment_file_size(file_id, 400);
858        tracker.register_dependencies(&[FileXorbDependency {
859            file_id,
860            xorb_hash: x,
861            n_bytes: 400,
862            is_external: false,
863        }]);
864
865        tracker.register_xorb_upload_progress(x, 500);
866        let (done, total) = tracker.status();
867        assert_eq!(total, 400);
868        assert!(done > 0);
869        assert!(done < 400);
870
871        tracker.increment_file_size(file_id, 200);
872        let (_, total) = tracker.status();
873        assert_eq!(total, 600);
874
875        tracker.register_dependencies(&[FileXorbDependency {
876            file_id,
877            xorb_hash: MerkleHash::random_from_seed(99),
878            n_bytes: 200,
879            is_external: true,
880        }]);
881
882        tracker.register_xorb_upload_completion(x);
883
884        let (done, total) = tracker.status();
885        assert_eq!(done, 600);
886        assert_eq!(total, 600);
887        assert!(tracker.is_complete());
888
889        tracker.assert_complete();
890        group.assert_complete();
891    }
892}