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
23pub type CompletionTrackerFileId = u64;
27
28#[derive(Default)]
30struct XorbDependency {
31 file_indices: BTreeSet<usize>,
33
34 completed_bytes: u64,
36
37 xorb_size: u64,
39
40 is_completed: bool,
42}
43
44#[derive(Default, Debug)]
45struct XorbPartCompletionStats {
46 completed_bytes: u64,
47 n_bytes: u64,
48}
49
50struct 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#[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}