Skip to main content

asupersync/atp/swarm/
piece_tracker.rs

1//! ATP Swarm Piece Tracker - Tracks piece availability and download progress.
2//!
3//! Manages the state of which pieces are available from which peers,
4//! tracks download progress, and coordinates piece requests.
5
6use 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
12/// Upper bound for the in-memory piece tracker.
13///
14/// `PieceTracker` currently materializes one status and one redundancy entry
15/// per piece, so untrusted `total_pieces` values must be capped before
16/// allocation.
17const MAX_TRACKED_PIECES_PER_TRANSFER: u64 = 1_000_000;
18
19/// Tracks piece availability and download state across transfers.
20#[derive(Debug)]
21pub struct PieceTracker {
22    /// Per-transfer piece maps
23    transfer_maps: HashMap<MailboxTransferId, TransferPieceMap>,
24
25    /// Transfer-scoped piece availability across all peers.
26    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/// Piece availability and progress for a single transfer.
36#[derive(Debug, Clone)]
37struct TransferPieceMap {
38    /// Total number of pieces in transfer
39    total_pieces: u64,
40
41    /// Pieces and their current status
42    piece_status: HashMap<PieceId, PieceStatus>,
43
44    /// Pieces available from each peer
45    peer_pieces: HashMap<PeerId, BTreeSet<PieceId>>,
46
47    /// Redundancy count for each piece
48    redundancy: HashMap<PieceId, u32>,
49}
50
51/// Map of piece availability across the swarm.
52#[derive(Debug, Clone, Serialize, Deserialize)]
53pub struct PieceMap {
54    /// Total number of pieces
55    pub total_pieces: u64,
56
57    /// Size of each piece in bytes
58    pub piece_size: u32,
59
60    /// Pieces available from each peer
61    pub peer_availability: HashMap<PeerId, BTreeSet<PieceId>>,
62
63    /// Content hash for verification
64    pub content_hash: String,
65
66    /// Creation timestamp
67    pub created_at: Time,
68}
69
70/// Status of an individual piece.
71#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
72pub enum PieceStatus {
73    /// Piece is needed and not yet requested
74    #[default]
75    Needed,
76
77    /// Piece has been requested from a peer
78    Requested {
79        /// Time when request was sent
80        requested_at: Time,
81        /// Peer from which piece was requested
82        peer_id: PeerId,
83    },
84
85    /// Piece is currently being downloaded
86    Downloading {
87        /// Download start time
88        started_at: Time,
89        /// Peer providing the piece
90        peer_id: PeerId,
91        /// Progress percentage (0.0 to 1.0)
92        progress: f64,
93    },
94
95    /// Piece download completed successfully
96    Completed {
97        /// Completion time
98        completed_at: Time,
99        /// Peer that provided the piece
100        peer_id: PeerId,
101    },
102
103    /// Piece download failed
104    Failed {
105        /// Failure time
106        failed_at: Time,
107        /// Peer that failed to provide piece
108        peer_id: PeerId,
109        /// Failure reason
110        reason: String,
111    },
112
113    /// Piece is being verified
114    Verifying {
115        /// Verification start time
116        started_at: Time,
117        /// Peer that provided the piece
118        peer_id: PeerId,
119    },
120}
121
122/// Statistics about piece distribution and redundancy.
123#[derive(Debug, Clone, Serialize, Deserialize)]
124pub struct PieceDistributionStats {
125    /// Total unique pieces tracked
126    pub total_unique_pieces: u64,
127
128    /// Average redundancy factor
129    pub avg_redundancy: f64,
130
131    /// Minimum redundancy (rarest piece)
132    pub min_redundancy: u32,
133
134    /// Maximum redundancy
135    pub max_redundancy: u32,
136
137    /// Pieces with only one peer (rarest)
138    pub rarest_pieces: Vec<PieceId>,
139
140    /// Distribution of redundancy levels
141    pub redundancy_distribution: BTreeMap<u32, u32>,
142}
143
144impl PieceMap {
145    /// Create a new piece map.
146    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    /// Add piece availability for a peer.
157    pub fn add_peer_pieces(&mut self, peer_id: PeerId, pieces: BTreeSet<PieceId>) {
158        self.peer_availability.insert(peer_id, pieces);
159    }
160
161    /// Get all peers that have a specific piece.
162    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    /// Calculate redundancy for a piece.
176    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    /// Get statistics about piece distribution.
184    pub fn get_distribution_stats(&self) -> PieceDistributionStats {
185        let mut redundancy_counts = HashMap::new();
186
187        // Calculate redundancy for each piece
188        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        // Find rarest pieces
205        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        // Build redundancy distribution
212        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    /// Create a new piece tracker.
230    pub fn new() -> Self {
231        Self {
232            transfer_maps: HashMap::new(),
233            global_availability: HashMap::new(),
234        }
235    }
236
237    /// Initialize tracking for a new transfer.
238    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        // Create piece status map
272        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        // Build redundancy map
278        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        // Update global availability
294        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    /// Get pieces that still need to be downloaded for a transfer.
310    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    /// Get pieces sorted by rarity (rarest first).
331    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        // Sort by redundancy (ascending - rarest first)
353        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    /// Mark a piece as requested.
362    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    /// Mark a piece as downloading.
387    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    /// Update download progress for a piece.
413    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    /// Mark a piece as completed.
446    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        // Get peer ID from current status
459        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    /// Mark a piece as failed.
485    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        // Completed pieces are terminal: a late failure report must not put a
499        // verified piece back into the retry set or regress transfer progress.
500        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    /// Get status of a specific piece.
530    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    /// Get transfer progress statistics.
552    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    /// Clean up completed transfers.
587    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    /// Remove a peer from every tracked transfer.
613    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    /// Return the tracked redundancy for a piece in a transfer.
668    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    /// Get pieces available from a specific peer.
690    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/// Progress statistics for a transfer.
711#[derive(Debug, Clone, Default, Serialize, Deserialize)]
712pub struct TransferProgress {
713    /// Total number of pieces
714    pub total_pieces: u64,
715
716    /// Pieces still needed
717    pub needed: u64,
718
719    /// Pieces requested but not yet downloading
720    pub requested: u64,
721
722    /// Pieces currently downloading
723    pub downloading: u64,
724
725    /// Pieces being verified
726    pub verifying: u64,
727
728    /// Pieces completed successfully
729    pub completed: u64,
730
731    /// Pieces that failed
732    pub failed: u64,
733
734    /// Overall completion percentage
735    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); // All pieces initially needed
828    }
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        // Test requested -> downloading -> completed
843        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        // Complete one piece
1083        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        // Piece 0: only on peer1 (redundancy 1)
1104        assert_eq!(piece_map.get_piece_redundancy(&PieceId::new(0)), 1);
1105
1106        // Piece 4: on both peers (redundancy 2)
1107        assert_eq!(piece_map.get_piece_redundancy(&PieceId::new(4)), 2);
1108
1109        // Piece 8: only on peer2 (redundancy 1)
1110        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        // Should be sorted with rarest pieces (redundancy 1) first
1129        assert_eq!(by_rarity.len(), 10);
1130    }
1131}
1132
1133// Additional error types needed for piece tracker
1134impl 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
1151// Add missing error variants to the main SwarmError enum in mod.rs
1152use 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}