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
24pub type CompletionTrackerFileId = u64;
28
29#[derive(Default)]
31struct XorbDependency {
32 file_indices: BTreeSet<usize>,
34
35 completed_bytes: u64,
37
38 xorb_size: u64,
40
41 is_completed: bool,
43}
44
45#[derive(Default, Debug)]
46struct XorbPartCompletionStats {
47 completed_bytes: u64,
48 n_bytes: u64,
49}
50
51struct 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#[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}