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