1use super::{PeerId, PieceId, SwarmError, SwarmResult, swarm_time_now};
7use crate::atp::mailbox::MailboxTransferId;
8use crate::types::Time;
9use serde::{Deserialize, Serialize};
10use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
11
12const MAX_TRACKED_PIECES_PER_TRANSFER: u64 = 1_000_000;
18
19#[derive(Debug)]
21pub struct PieceTracker {
22 transfer_maps: HashMap<MailboxTransferId, TransferPieceMap>,
24
25 global_availability: HashMap<AvailabilityKey, HashSet<PeerId>>,
27}
28
29#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
30struct AvailabilityKey {
31 transfer_id: MailboxTransferId,
32 piece_id: PieceId,
33}
34
35#[derive(Debug, Clone)]
37struct TransferPieceMap {
38 total_pieces: u64,
40
41 piece_status: HashMap<PieceId, PieceStatus>,
43
44 peer_pieces: HashMap<PeerId, BTreeSet<PieceId>>,
46
47 redundancy: HashMap<PieceId, u32>,
49}
50
51#[derive(Debug, Clone, Serialize, Deserialize)]
53pub struct PieceMap {
54 pub total_pieces: u64,
56
57 pub piece_size: u32,
59
60 pub peer_availability: HashMap<PeerId, BTreeSet<PieceId>>,
62
63 pub content_hash: String,
65
66 pub created_at: Time,
68}
69
70#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
72pub enum PieceStatus {
73 #[default]
75 Needed,
76
77 Requested {
79 requested_at: Time,
81 peer_id: PeerId,
83 },
84
85 Downloading {
87 started_at: Time,
89 peer_id: PeerId,
91 progress: f64,
93 },
94
95 Completed {
97 completed_at: Time,
99 peer_id: PeerId,
101 },
102
103 Failed {
105 failed_at: Time,
107 peer_id: PeerId,
109 reason: String,
111 },
112
113 Verifying {
115 started_at: Time,
117 peer_id: PeerId,
119 },
120}
121
122#[derive(Debug, Clone, Serialize, Deserialize)]
124pub struct PieceDistributionStats {
125 pub total_unique_pieces: u64,
127
128 pub avg_redundancy: f64,
130
131 pub min_redundancy: u32,
133
134 pub max_redundancy: u32,
136
137 pub rarest_pieces: Vec<PieceId>,
139
140 pub redundancy_distribution: BTreeMap<u32, u32>,
142}
143
144impl PieceMap {
145 pub fn new(total_pieces: u64, piece_size: u32, content_hash: String) -> Self {
147 Self {
148 total_pieces,
149 piece_size,
150 peer_availability: HashMap::new(),
151 content_hash,
152 created_at: swarm_time_now(),
153 }
154 }
155
156 pub fn add_peer_pieces(&mut self, peer_id: PeerId, pieces: BTreeSet<PieceId>) {
158 self.peer_availability.insert(peer_id, pieces);
159 }
160
161 pub fn get_peers_for_piece(&self, piece_id: &PieceId) -> Vec<PeerId> {
163 self.peer_availability
164 .iter()
165 .filter_map(|(peer_id, pieces): (&PeerId, &BTreeSet<PieceId>)| {
166 if pieces.contains(piece_id) {
167 Some(peer_id.clone())
168 } else {
169 None
170 }
171 })
172 .collect()
173 }
174
175 pub fn get_piece_redundancy(&self, piece_id: &PieceId) -> u32 {
177 self.peer_availability
178 .values()
179 .filter(|pieces: &&BTreeSet<PieceId>| pieces.contains(piece_id))
180 .count() as u32
181 }
182
183 pub fn get_distribution_stats(&self) -> PieceDistributionStats {
185 let mut redundancy_counts = HashMap::new();
186
187 for piece_id in 0..self.total_pieces {
189 let piece_id = PieceId::new(piece_id);
190 let redundancy = self.get_piece_redundancy(&piece_id);
191 redundancy_counts.insert(piece_id, redundancy);
192 }
193
194 let redundancy_values: Vec<u32> = redundancy_counts.values().copied().collect();
195 let avg_redundancy = if redundancy_values.is_empty() {
196 0.0
197 } else {
198 redundancy_values.iter().sum::<u32>() as f64 / redundancy_values.len() as f64
199 };
200
201 let min_redundancy = redundancy_values.iter().min().copied().unwrap_or(0);
202 let max_redundancy = redundancy_values.iter().max().copied().unwrap_or(0);
203
204 let rarest_pieces: Vec<PieceId> = redundancy_counts
206 .iter()
207 .filter(|&(_, &redundancy)| redundancy == min_redundancy)
208 .map(|(piece_id, _)| *piece_id)
209 .collect();
210
211 let mut redundancy_distribution = BTreeMap::new();
213 for &redundancy in &redundancy_values {
214 *redundancy_distribution.entry(redundancy).or_insert(0) += 1;
215 }
216
217 PieceDistributionStats {
218 total_unique_pieces: self.total_pieces,
219 avg_redundancy,
220 min_redundancy,
221 max_redundancy,
222 rarest_pieces,
223 redundancy_distribution,
224 }
225 }
226}
227
228impl PieceTracker {
229 pub fn new() -> Self {
231 Self {
232 transfer_maps: HashMap::new(),
233 global_availability: HashMap::new(),
234 }
235 }
236
237 pub fn initialize_transfer(
239 &mut self,
240 transfer_id: &MailboxTransferId,
241 piece_map: &PieceMap,
242 ) -> SwarmResult<()> {
243 let total_pieces = piece_map.total_pieces;
244
245 if total_pieces > MAX_TRACKED_PIECES_PER_TRANSFER {
246 return Err(SwarmError::ConfigurationError {
247 details: format!(
248 "piece map declares {total_pieces} pieces, exceeding tracker limit {MAX_TRACKED_PIECES_PER_TRANSFER}"
249 ),
250 });
251 }
252
253 for (peer_id, pieces) in &piece_map.peer_availability {
254 if let Some(piece_id) = pieces
255 .iter()
256 .find(|piece_id| piece_id.as_u64() >= total_pieces)
257 {
258 return Err(SwarmError::ConfigurationError {
259 details: format!(
260 "peer {peer_id:?} advertises piece {piece_id} outside declared range 0..{total_pieces}"
261 ),
262 });
263 }
264 }
265
266 let piece_capacity =
267 usize::try_from(total_pieces).map_err(|_| SwarmError::ConfigurationError {
268 details: format!("piece map declares {total_pieces} pieces, exceeding usize"),
269 })?;
270
271 let mut piece_status = HashMap::with_capacity(piece_capacity);
273 for piece_id in 0..total_pieces {
274 piece_status.insert(PieceId::new(piece_id), PieceStatus::Needed);
275 }
276
277 let mut redundancy = HashMap::with_capacity(piece_capacity);
279 for piece_id in 0..total_pieces {
280 let piece_id = PieceId::new(piece_id);
281 redundancy.insert(piece_id, piece_map.get_piece_redundancy(&piece_id));
282 }
283
284 let transfer_map = TransferPieceMap {
285 total_pieces,
286 piece_status,
287 peer_pieces: piece_map.peer_availability.clone(),
288 redundancy,
289 };
290
291 self.transfer_maps.insert(*transfer_id, transfer_map);
292
293 for (peer_id, pieces) in &piece_map.peer_availability {
295 for piece_id in pieces {
296 self.global_availability
297 .entry(AvailabilityKey {
298 transfer_id: *transfer_id,
299 piece_id: *piece_id,
300 })
301 .or_default()
302 .insert(peer_id.clone());
303 }
304 }
305
306 Ok(())
307 }
308
309 pub fn get_needed_pieces(&self, transfer_id: &MailboxTransferId) -> SwarmResult<Vec<PieceId>> {
311 let transfer_map =
312 self.transfer_maps
313 .get(transfer_id)
314 .ok_or(SwarmError::TransferNotFound {
315 transfer_id: *transfer_id,
316 })?;
317
318 let needed_pieces: Vec<PieceId> = transfer_map
319 .piece_status
320 .iter()
321 .filter_map(|(piece_id, status)| match status {
322 PieceStatus::Needed | PieceStatus::Failed { .. } => Some(*piece_id),
323 _ => None,
324 })
325 .collect();
326
327 Ok(needed_pieces)
328 }
329
330 pub fn get_pieces_by_rarity(
332 &self,
333 transfer_id: &MailboxTransferId,
334 ) -> SwarmResult<Vec<PieceId>> {
335 let transfer_map =
336 self.transfer_maps
337 .get(transfer_id)
338 .ok_or(SwarmError::TransferNotFound {
339 transfer_id: *transfer_id,
340 })?;
341
342 let needed_pieces = self.get_needed_pieces(transfer_id)?;
343
344 let mut rarity_sorted: Vec<(PieceId, u32)> = needed_pieces
345 .into_iter()
346 .map(|piece_id| {
347 let redundancy = transfer_map.redundancy.get(&piece_id).copied().unwrap_or(0);
348 (piece_id, redundancy)
349 })
350 .collect();
351
352 rarity_sorted.sort_by_key(|(_, redundancy)| *redundancy);
354
355 Ok(rarity_sorted
356 .into_iter()
357 .map(|(piece_id, _)| piece_id)
358 .collect())
359 }
360
361 pub fn mark_piece_requested(
363 &mut self,
364 transfer_id: &MailboxTransferId,
365 piece_id: PieceId,
366 peer_id: PeerId,
367 ) -> SwarmResult<()> {
368 let transfer_map =
369 self.transfer_maps
370 .get_mut(transfer_id)
371 .ok_or(SwarmError::TransferNotFound {
372 transfer_id: *transfer_id,
373 })?;
374
375 transfer_map.piece_status.insert(
376 piece_id,
377 PieceStatus::Requested {
378 requested_at: swarm_time_now(),
379 peer_id,
380 },
381 );
382
383 Ok(())
384 }
385
386 pub fn mark_piece_downloading(
388 &mut self,
389 transfer_id: &MailboxTransferId,
390 piece_id: PieceId,
391 peer_id: PeerId,
392 ) -> SwarmResult<()> {
393 let transfer_map =
394 self.transfer_maps
395 .get_mut(transfer_id)
396 .ok_or(SwarmError::TransferNotFound {
397 transfer_id: *transfer_id,
398 })?;
399
400 transfer_map.piece_status.insert(
401 piece_id,
402 PieceStatus::Downloading {
403 started_at: swarm_time_now(),
404 peer_id,
405 progress: 0.0,
406 },
407 );
408
409 Ok(())
410 }
411
412 pub fn update_piece_progress(
414 &mut self,
415 transfer_id: &MailboxTransferId,
416 piece_id: PieceId,
417 progress: f64,
418 ) -> SwarmResult<()> {
419 let transfer_map =
420 self.transfer_maps
421 .get_mut(transfer_id)
422 .ok_or(SwarmError::TransferNotFound {
423 transfer_id: *transfer_id,
424 })?;
425
426 if let Some(PieceStatus::Downloading {
427 started_at,
428 peer_id,
429 ..
430 }) = transfer_map.piece_status.get(&piece_id)
431 {
432 transfer_map.piece_status.insert(
433 piece_id,
434 PieceStatus::Downloading {
435 started_at: *started_at,
436 peer_id: peer_id.clone(),
437 progress: progress.clamp(0.0, 1.0),
438 },
439 );
440 }
441
442 Ok(())
443 }
444
445 pub fn mark_piece_completed(
447 &mut self,
448 transfer_id: &MailboxTransferId,
449 piece_id: PieceId,
450 ) -> SwarmResult<()> {
451 let transfer_map =
452 self.transfer_maps
453 .get_mut(transfer_id)
454 .ok_or(SwarmError::TransferNotFound {
455 transfer_id: *transfer_id,
456 })?;
457
458 let peer_id = match transfer_map.piece_status.get(&piece_id) {
460 Some(
461 PieceStatus::Requested { peer_id, .. }
462 | PieceStatus::Downloading { peer_id, .. }
463 | PieceStatus::Verifying { peer_id, .. },
464 ) => peer_id.clone(),
465 _ => {
466 return Err(SwarmError::InvalidPieceState {
467 piece_id,
468 current_state: "not requested, downloading, or verifying".to_string(),
469 });
470 }
471 };
472
473 transfer_map.piece_status.insert(
474 piece_id,
475 PieceStatus::Completed {
476 completed_at: swarm_time_now(),
477 peer_id,
478 },
479 );
480
481 Ok(())
482 }
483
484 pub fn mark_piece_failed(
486 &mut self,
487 transfer_id: &MailboxTransferId,
488 piece_id: PieceId,
489 reason: String,
490 ) -> SwarmResult<()> {
491 let transfer_map =
492 self.transfer_maps
493 .get_mut(transfer_id)
494 .ok_or(SwarmError::TransferNotFound {
495 transfer_id: *transfer_id,
496 })?;
497
498 let peer_id = match transfer_map.piece_status.get(&piece_id) {
501 Some(PieceStatus::Completed { .. }) => return Ok(()),
502 Some(
503 PieceStatus::Downloading { peer_id, .. }
504 | PieceStatus::Verifying { peer_id, .. }
505 | PieceStatus::Requested { peer_id, .. },
506 ) => peer_id.clone(),
507 Some(PieceStatus::Failed { peer_id, .. }) => peer_id.clone(),
508 Some(PieceStatus::Needed) => {
509 return Err(SwarmError::InvalidPieceState {
510 piece_id,
511 current_state: "needed".to_string(),
512 });
513 }
514 None => return Err(SwarmError::PieceNotFound { piece_id }),
515 };
516
517 transfer_map.piece_status.insert(
518 piece_id,
519 PieceStatus::Failed {
520 failed_at: swarm_time_now(),
521 peer_id,
522 reason,
523 },
524 );
525
526 Ok(())
527 }
528
529 pub fn get_piece_status(
531 &self,
532 transfer_id: &MailboxTransferId,
533 piece_id: &PieceId,
534 ) -> SwarmResult<PieceStatus> {
535 let transfer_map =
536 self.transfer_maps
537 .get(transfer_id)
538 .ok_or(SwarmError::TransferNotFound {
539 transfer_id: *transfer_id,
540 })?;
541
542 transfer_map
543 .piece_status
544 .get(piece_id)
545 .cloned()
546 .ok_or(SwarmError::PieceNotFound {
547 piece_id: *piece_id,
548 })
549 }
550
551 pub fn get_transfer_progress(
553 &self,
554 transfer_id: &MailboxTransferId,
555 ) -> SwarmResult<TransferProgress> {
556 let transfer_map =
557 self.transfer_maps
558 .get(transfer_id)
559 .ok_or(SwarmError::TransferNotFound {
560 transfer_id: *transfer_id,
561 })?;
562
563 let mut progress = TransferProgress::default();
564 progress.total_pieces = transfer_map.total_pieces;
565
566 for status in transfer_map.piece_status.values() {
567 match status {
568 PieceStatus::Needed => progress.needed += 1,
569 PieceStatus::Requested { .. } => progress.requested += 1,
570 PieceStatus::Downloading { .. } => progress.downloading += 1,
571 PieceStatus::Completed { .. } => progress.completed += 1,
572 PieceStatus::Failed { .. } => progress.failed += 1,
573 PieceStatus::Verifying { .. } => progress.verifying += 1,
574 }
575 }
576
577 progress.completion_percentage = if progress.total_pieces > 0 {
578 (progress.completed as f64 / progress.total_pieces as f64) * 100.0
579 } else {
580 0.0
581 };
582
583 Ok(progress)
584 }
585
586 pub fn cleanup_transfer(&mut self, transfer_id: &MailboxTransferId) {
588 if let Some(transfer_map) = self.transfer_maps.remove(transfer_id) {
589 for (peer_id, pieces) in transfer_map.peer_pieces {
590 for piece_id in pieces {
591 let availability_key = AvailabilityKey {
592 transfer_id: *transfer_id,
593 piece_id,
594 };
595 let entry_is_empty = if let Some(peer_set) =
596 self.global_availability.get_mut(&availability_key)
597 {
598 peer_set.remove(&peer_id);
599 peer_set.is_empty()
600 } else {
601 false
602 };
603
604 if entry_is_empty {
605 self.global_availability.remove(&availability_key);
606 }
607 }
608 }
609 }
610 }
611
612 pub fn remove_peer(&mut self, peer_id: &PeerId) {
614 let failed_at = swarm_time_now();
615
616 for (transfer_id, transfer_map) in &mut self.transfer_maps {
617 if let Some(removed_pieces) = transfer_map.peer_pieces.remove(peer_id) {
618 for piece_id in removed_pieces {
619 let availability_key = AvailabilityKey {
620 transfer_id: *transfer_id,
621 piece_id,
622 };
623 let entry_is_empty = if let Some(peer_set) =
624 self.global_availability.get_mut(&availability_key)
625 {
626 peer_set.remove(peer_id);
627 peer_set.is_empty()
628 } else {
629 false
630 };
631
632 if entry_is_empty {
633 self.global_availability.remove(&availability_key);
634 }
635
636 let remaining_redundancy = transfer_map
637 .peer_pieces
638 .values()
639 .filter(|pieces| pieces.contains(&piece_id))
640 .count() as u32;
641 transfer_map
642 .redundancy
643 .insert(piece_id, remaining_redundancy);
644 }
645 }
646
647 for status in transfer_map.piece_status.values_mut() {
648 let was_assigned_to_removed_peer = matches!(
649 status,
650 PieceStatus::Requested { peer_id: assigned_peer, .. }
651 | PieceStatus::Downloading { peer_id: assigned_peer, .. }
652 | PieceStatus::Verifying { peer_id: assigned_peer, .. }
653 if assigned_peer == peer_id
654 );
655
656 if was_assigned_to_removed_peer {
657 *status = PieceStatus::Failed {
658 failed_at,
659 peer_id: peer_id.clone(),
660 reason: "peer removed".to_string(),
661 };
662 }
663 }
664 }
665 }
666
667 pub fn get_piece_redundancy(
669 &self,
670 transfer_id: &MailboxTransferId,
671 piece_id: &PieceId,
672 ) -> SwarmResult<u32> {
673 let transfer_map =
674 self.transfer_maps
675 .get(transfer_id)
676 .ok_or(SwarmError::TransferNotFound {
677 transfer_id: *transfer_id,
678 })?;
679
680 transfer_map
681 .redundancy
682 .get(piece_id)
683 .copied()
684 .ok_or(SwarmError::PieceNotFound {
685 piece_id: *piece_id,
686 })
687 }
688
689 pub fn get_peer_pieces(
691 &self,
692 transfer_id: &MailboxTransferId,
693 peer_id: &PeerId,
694 ) -> SwarmResult<BTreeSet<PieceId>> {
695 let transfer_map =
696 self.transfer_maps
697 .get(transfer_id)
698 .ok_or(SwarmError::TransferNotFound {
699 transfer_id: *transfer_id,
700 })?;
701
702 Ok(transfer_map
703 .peer_pieces
704 .get(peer_id)
705 .cloned()
706 .unwrap_or_default())
707 }
708}
709
710#[derive(Debug, Clone, Default, Serialize, Deserialize)]
712pub struct TransferProgress {
713 pub total_pieces: u64,
715
716 pub needed: u64,
718
719 pub requested: u64,
721
722 pub downloading: u64,
724
725 pub verifying: u64,
727
728 pub completed: u64,
730
731 pub failed: u64,
733
734 pub completion_percentage: f64,
736}
737
738#[cfg(test)]
739mod tests {
740 use super::*;
741
742 fn create_test_piece_map() -> PieceMap {
743 let mut piece_map = PieceMap::new(10, 1024, "test-hash".to_string());
744
745 let peer1 = PeerId::new("peer1");
746 let peer2 = PeerId::new("peer2");
747
748 piece_map.add_peer_pieces(peer1, (0..5).map(PieceId::new).collect());
749 piece_map.add_peer_pieces(peer2, (3..10).map(PieceId::new).collect());
750
751 piece_map
752 }
753
754 #[test]
755 fn test_piece_tracker_creation() {
756 let tracker = PieceTracker::new();
757 assert_eq!(tracker.transfer_maps.len(), 0);
758 assert_eq!(tracker.global_availability.len(), 0);
759 }
760
761 #[test]
762 fn test_initialize_transfer() {
763 let mut tracker = PieceTracker::new();
764 let piece_map = create_test_piece_map();
765 let transfer_id = MailboxTransferId::new();
766
767 let result = tracker.initialize_transfer(&transfer_id, &piece_map);
768 assert!(result.is_ok());
769 assert!(tracker.transfer_maps.contains_key(&transfer_id));
770 }
771
772 #[test]
773 fn initialize_transfer_rejects_excessive_piece_count_before_allocation() {
774 let mut tracker = PieceTracker::new();
775 let piece_map = PieceMap::new(
776 MAX_TRACKED_PIECES_PER_TRANSFER + 1,
777 1024,
778 "test-hash".to_string(),
779 );
780 let transfer_id = MailboxTransferId::new();
781
782 let err = tracker
783 .initialize_transfer(&transfer_id, &piece_map)
784 .expect_err("excessive piece count must be rejected");
785
786 assert!(
787 matches!(&err, SwarmError::ConfigurationError { details } if details.contains("exceeding tracker limit")),
788 "unexpected error: {err}"
789 );
790 assert!(tracker.transfer_maps.is_empty());
791 assert!(tracker.global_availability.is_empty());
792 }
793
794 #[test]
795 fn initialize_transfer_rejects_out_of_range_peer_piece() {
796 let mut tracker = PieceTracker::new();
797 let mut piece_map = PieceMap::new(2, 1024, "test-hash".to_string());
798 piece_map.add_peer_pieces(
799 PeerId::new("peer1"),
800 [PieceId::new(0), PieceId::new(2)].into_iter().collect(),
801 );
802 let transfer_id = MailboxTransferId::new();
803
804 let err = tracker
805 .initialize_transfer(&transfer_id, &piece_map)
806 .expect_err("out-of-range advertised piece must be rejected");
807
808 assert!(
809 matches!(&err, SwarmError::ConfigurationError { details } if details.contains("outside declared range 0..2")),
810 "unexpected error: {err}"
811 );
812 assert!(tracker.transfer_maps.is_empty());
813 assert!(tracker.global_availability.is_empty());
814 }
815
816 #[test]
817 fn test_get_needed_pieces() {
818 let mut tracker = PieceTracker::new();
819 let piece_map = create_test_piece_map();
820 let transfer_id = MailboxTransferId::new();
821
822 tracker
823 .initialize_transfer(&transfer_id, &piece_map)
824 .unwrap();
825 let needed = tracker.get_needed_pieces(&transfer_id).unwrap();
826
827 assert_eq!(needed.len(), 10); }
829
830 #[test]
831 fn test_piece_status_transitions() {
832 let mut tracker = PieceTracker::new();
833 let piece_map = create_test_piece_map();
834 let transfer_id = MailboxTransferId::new();
835 let piece_id = PieceId::new(0);
836 let peer_id = PeerId::new("peer1");
837
838 tracker
839 .initialize_transfer(&transfer_id, &piece_map)
840 .unwrap();
841
842 tracker
844 .mark_piece_requested(&transfer_id, piece_id, peer_id.clone())
845 .unwrap();
846 let status = tracker.get_piece_status(&transfer_id, &piece_id).unwrap();
847 assert!(matches!(status, PieceStatus::Requested { .. }));
848
849 tracker
850 .mark_piece_downloading(&transfer_id, piece_id, peer_id.clone())
851 .unwrap();
852 let status = tracker.get_piece_status(&transfer_id, &piece_id).unwrap();
853 assert!(matches!(status, PieceStatus::Downloading { .. }));
854
855 tracker
856 .mark_piece_completed(&transfer_id, piece_id)
857 .unwrap();
858 let status = tracker.get_piece_status(&transfer_id, &piece_id).unwrap();
859 assert!(matches!(status, PieceStatus::Completed { .. }));
860 }
861
862 #[test]
863 fn completed_piece_cannot_be_regressed_to_failed() {
864 let mut tracker = PieceTracker::new();
865 let piece_map = create_test_piece_map();
866 let transfer_id = MailboxTransferId::new();
867 let piece_id = PieceId::new(0);
868 let peer_id = PeerId::new("peer1");
869
870 tracker
871 .initialize_transfer(&transfer_id, &piece_map)
872 .unwrap();
873 tracker
874 .mark_piece_downloading(&transfer_id, piece_id, peer_id.clone())
875 .unwrap();
876 tracker
877 .mark_piece_completed(&transfer_id, piece_id)
878 .unwrap();
879
880 tracker
881 .mark_piece_failed(
882 &transfer_id,
883 piece_id,
884 "late verification failure".to_string(),
885 )
886 .unwrap();
887
888 let status = tracker.get_piece_status(&transfer_id, &piece_id).unwrap();
889 assert!(
890 matches!(status, PieceStatus::Completed { peer_id: completed_by, .. } if completed_by == peer_id)
891 );
892 let needed = tracker.get_needed_pieces(&transfer_id).unwrap();
893 assert!(
894 !needed.contains(&piece_id),
895 "completed piece must not re-enter the retry set"
896 );
897 let progress = tracker.get_transfer_progress(&transfer_id).unwrap();
898 assert_eq!(progress.completed, 1);
899 assert_eq!(progress.failed, 0);
900 assert_eq!(progress.needed, 9);
901 assert_eq!(progress.completion_percentage, 10.0);
902 }
903
904 #[test]
905 fn failed_piece_does_not_synthesize_unknown_peer() {
906 let mut tracker = PieceTracker::new();
907 let piece_map = create_test_piece_map();
908 let transfer_id = MailboxTransferId::new();
909 let unassigned_piece = PieceId::new(1);
910 let requested_piece = PieceId::new(0);
911 let real_unknown_peer = PeerId::new("unknown");
912
913 tracker
914 .initialize_transfer(&transfer_id, &piece_map)
915 .unwrap();
916
917 let err = tracker
918 .mark_piece_failed(
919 &transfer_id,
920 unassigned_piece,
921 "failure before request".to_string(),
922 )
923 .expect_err("unassigned failures must not invent a peer id");
924 assert!(
925 matches!(&err, SwarmError::InvalidPieceState { piece_id, current_state } if *piece_id == unassigned_piece && current_state == "needed"),
926 "unexpected error: {err}"
927 );
928 assert!(matches!(
929 tracker
930 .get_piece_status(&transfer_id, &unassigned_piece)
931 .unwrap(),
932 PieceStatus::Needed
933 ));
934
935 tracker
936 .mark_piece_requested(&transfer_id, requested_piece, real_unknown_peer.clone())
937 .unwrap();
938 tracker
939 .mark_piece_failed(
940 &transfer_id,
941 requested_piece,
942 "real peer named unknown failed".to_string(),
943 )
944 .unwrap();
945
946 let status = tracker
947 .get_piece_status(&transfer_id, &requested_piece)
948 .unwrap();
949 assert!(
950 matches!(status, PieceStatus::Failed { peer_id, .. } if peer_id == real_unknown_peer)
951 );
952 }
953
954 #[test]
955 fn cleanup_transfer_preserves_other_transfer_availability_for_same_piece() {
956 let mut tracker = PieceTracker::new();
957 let transfer_a = MailboxTransferId::new();
958 let transfer_b = MailboxTransferId::new();
959 let piece_id = PieceId::new(0);
960 let peer_a = PeerId::new("peer-a");
961 let peer_b = PeerId::new("peer-b");
962
963 let mut map_a = PieceMap::new(1, 1024, "hash-a".to_string());
964 map_a.add_peer_pieces(peer_a.clone(), std::iter::once(piece_id).collect());
965 let mut map_b = PieceMap::new(1, 1024, "hash-b".to_string());
966 map_b.add_peer_pieces(peer_b.clone(), std::iter::once(piece_id).collect());
967
968 tracker.initialize_transfer(&transfer_a, &map_a).unwrap();
969 tracker.initialize_transfer(&transfer_b, &map_b).unwrap();
970 assert_eq!(tracker.global_availability.len(), 2);
971
972 tracker.cleanup_transfer(&transfer_a);
973
974 assert!(!tracker.transfer_maps.contains_key(&transfer_a));
975 assert!(tracker.transfer_maps.contains_key(&transfer_b));
976 assert!(!tracker.global_availability.contains_key(&AvailabilityKey {
977 transfer_id: transfer_a,
978 piece_id,
979 }));
980
981 let remaining_peers = tracker
982 .global_availability
983 .get(&AvailabilityKey {
984 transfer_id: transfer_b,
985 piece_id,
986 })
987 .expect("second transfer availability must remain");
988 assert_eq!(remaining_peers.len(), 1);
989 assert!(remaining_peers.contains(&peer_b));
990 assert!(!remaining_peers.contains(&peer_a));
991 }
992
993 #[test]
994 fn remove_peer_drops_availability_and_retries_inflight_pieces() {
995 let mut tracker = PieceTracker::new();
996 let transfer_id = MailboxTransferId::new();
997 let peer_a = PeerId::new("peer-a");
998 let peer_b = PeerId::new("peer-b");
999 let piece_a = PieceId::new(0);
1000 let shared_piece = PieceId::new(1);
1001 let mut piece_map = PieceMap::new(2, 1024, "test-hash".to_string());
1002 piece_map.add_peer_pieces(
1003 peer_a.clone(),
1004 [piece_a, shared_piece].into_iter().collect(),
1005 );
1006 piece_map.add_peer_pieces(peer_b.clone(), std::iter::once(shared_piece).collect());
1007
1008 tracker
1009 .initialize_transfer(&transfer_id, &piece_map)
1010 .unwrap();
1011 tracker
1012 .mark_piece_requested(&transfer_id, piece_a, peer_a.clone())
1013 .unwrap();
1014 tracker
1015 .mark_piece_downloading(&transfer_id, shared_piece, peer_b.clone())
1016 .unwrap();
1017
1018 tracker.remove_peer(&peer_a);
1019
1020 assert!(
1021 tracker
1022 .get_peer_pieces(&transfer_id, &peer_a)
1023 .unwrap()
1024 .is_empty()
1025 );
1026 assert_eq!(
1027 tracker.get_peer_pieces(&transfer_id, &peer_b).unwrap(),
1028 std::iter::once(shared_piece).collect()
1029 );
1030 assert_eq!(
1031 tracker
1032 .get_piece_redundancy(&transfer_id, &piece_a)
1033 .unwrap(),
1034 0
1035 );
1036 assert_eq!(
1037 tracker
1038 .get_piece_redundancy(&transfer_id, &shared_piece)
1039 .unwrap(),
1040 1
1041 );
1042
1043 let failed_status = tracker.get_piece_status(&transfer_id, &piece_a).unwrap();
1044 assert!(
1045 matches!(failed_status, PieceStatus::Failed { peer_id, reason, .. } if peer_id == peer_a && reason == "peer removed")
1046 );
1047 let shared_status = tracker
1048 .get_piece_status(&transfer_id, &shared_piece)
1049 .unwrap();
1050 assert!(
1051 matches!(shared_status, PieceStatus::Downloading { peer_id, .. } if peer_id == peer_b)
1052 );
1053
1054 let needed = tracker.get_needed_pieces(&transfer_id).unwrap();
1055 assert!(needed.contains(&piece_a));
1056 assert!(!needed.contains(&shared_piece));
1057 assert!(!tracker.global_availability.contains_key(&AvailabilityKey {
1058 transfer_id,
1059 piece_id: piece_a,
1060 }));
1061 let remaining_peers = tracker
1062 .global_availability
1063 .get(&AvailabilityKey {
1064 transfer_id,
1065 piece_id: shared_piece,
1066 })
1067 .expect("shared piece availability must remain");
1068 assert_eq!(remaining_peers.len(), 1);
1069 assert!(remaining_peers.contains(&peer_b));
1070 }
1071
1072 #[test]
1073 fn test_transfer_progress() {
1074 let mut tracker = PieceTracker::new();
1075 let piece_map = create_test_piece_map();
1076 let transfer_id = MailboxTransferId::new();
1077
1078 tracker
1079 .initialize_transfer(&transfer_id, &piece_map)
1080 .unwrap();
1081
1082 let piece_id = PieceId::new(0);
1084 let peer_id = PeerId::new("peer1");
1085 tracker
1086 .mark_piece_downloading(&transfer_id, piece_id, peer_id)
1087 .unwrap();
1088 tracker
1089 .mark_piece_completed(&transfer_id, piece_id)
1090 .unwrap();
1091
1092 let progress = tracker.get_transfer_progress(&transfer_id).unwrap();
1093 assert_eq!(progress.total_pieces, 10);
1094 assert_eq!(progress.completed, 1);
1095 assert_eq!(progress.needed, 9);
1096 assert_eq!(progress.completion_percentage, 10.0);
1097 }
1098
1099 #[test]
1100 fn test_piece_map_redundancy() {
1101 let piece_map = create_test_piece_map();
1102
1103 assert_eq!(piece_map.get_piece_redundancy(&PieceId::new(0)), 1);
1105
1106 assert_eq!(piece_map.get_piece_redundancy(&PieceId::new(4)), 2);
1108
1109 assert_eq!(piece_map.get_piece_redundancy(&PieceId::new(8)), 1);
1111
1112 let stats = piece_map.get_distribution_stats();
1113 assert_eq!(stats.min_redundancy, 1);
1114 assert_eq!(stats.max_redundancy, 2);
1115 }
1116
1117 #[test]
1118 fn test_pieces_by_rarity() {
1119 let mut tracker = PieceTracker::new();
1120 let piece_map = create_test_piece_map();
1121 let transfer_id = MailboxTransferId::new();
1122
1123 tracker
1124 .initialize_transfer(&transfer_id, &piece_map)
1125 .unwrap();
1126 let by_rarity = tracker.get_pieces_by_rarity(&transfer_id).unwrap();
1127
1128 assert_eq!(by_rarity.len(), 10);
1130 }
1131}
1132
1133impl SwarmError {
1135 pub fn transfer_not_found(transfer_id: MailboxTransferId) -> Self {
1136 SwarmError::TransferNotFound { transfer_id }
1137 }
1138
1139 pub fn piece_not_found(piece_id: PieceId) -> Self {
1140 SwarmError::PieceNotFound { piece_id }
1141 }
1142
1143 pub fn invalid_piece_state(piece_id: PieceId, current_state: String) -> Self {
1144 SwarmError::InvalidPieceState {
1145 piece_id,
1146 current_state,
1147 }
1148 }
1149}
1150
1151use thiserror::Error;
1153
1154#[derive(Debug, Error)]
1155pub enum PieceTrackerError {
1156 #[error("Transfer not found: {transfer_id}")]
1157 TransferNotFound { transfer_id: MailboxTransferId },
1158
1159 #[error("Piece not found: {piece_id:?}")]
1160 PieceNotFound { piece_id: PieceId },
1161
1162 #[error("Invalid piece state for {piece_id:?}: {current_state}")]
1163 InvalidPieceState {
1164 piece_id: PieceId,
1165 current_state: String,
1166 },
1167}