Skip to main content

aria2_core/engine/
bt_download_execute.rs

1use async_trait::async_trait;
2use std::collections::{HashMap, HashSet};
3use std::time::{Duration, Instant};
4use tracing::{debug, info, warn};
5
6use crate::engine::bt_download_command::{
7    BLOCK_SIZE, BtDownloadCommand, MAX_PUBLIC_TRACKERS_TO_TRY, MAX_RETRIES,
8    PEER_CONNECTION_DELAY_MS, PUBLIC_TRACKER_PEER_THRESHOLD,
9};
10use crate::engine::bt_message_handler::BtMessageHandler;
11use crate::engine::bt_peer_connection::BtPeerConn;
12use crate::engine::bt_peer_interaction::BtPeerInteraction;
13use crate::engine::bt_piece_downloader::write_piece_to_multi_files_coalesced;
14use crate::engine::bt_piece_selector::BtPieceSelector;
15use crate::engine::bt_post_download_handler::{DownloadStatus, HookContext};
16use crate::engine::bt_progress_info_file::{BtProgress, DownloadStats as ProgressDownloadStats};
17use crate::engine::bt_tracker_comm::{announce_to_public_tracker, perform_http_tracker_announce};
18use crate::engine::choking_algorithm::{ChokingAlgorithm, ChokingConfig};
19use crate::engine::command::{Command, CommandStatus};
20use crate::engine::peer_stats::PeerStats;
21use crate::engine::udp_tracker_client::UdpTrackerClient;
22use crate::engine::udp_tracker_manager::UdpTrackerManager;
23use crate::error::{Aria2Error, FatalError, RecoverableError, Result};
24use crate::filesystem::disk_writer::{DefaultDiskWriter, DiskWriter};
25use crate::rate_limiter::{RateLimiter, RateLimiterConfig, ThrottledWriter};
26use aria2_protocol::bittorrent::extension::pex::PexHandler;
27use aria2_protocol::bittorrent::message::serializer;
28use aria2_protocol::bittorrent::piece::peer_tracker::PeerBitfieldTracker;
29
30/// Tracks duplicate requests during endgame mode.
31///
32/// In endgame mode (when <=5 pieces remain incomplete), we request the same block
33/// from multiple peers simultaneously to speed up completion. When any peer responds
34/// with the block data, we send Cancel messages to the other peers that also received
35/// the request for that block.
36///
37/// This struct maintains the mapping from block identifiers to the list of peers
38/// that were sent duplicate requests, enabling efficient cancellation on arrival.
39pub struct EndgameState {
40    /// Map from (piece_index, offset, length) -> list of peer indices that received this request
41    active_duplicate_requests: HashMap<(u32, u32, u32), Vec<usize>>,
42    /// Whether we're currently in endgame mode
43    active: bool,
44}
45
46impl EndgameState {
47    /// Create a new EndgameState in inactive state
48    pub fn new() -> Self {
49        Self {
50            active_duplicate_requests: HashMap::new(),
51            active: false,
52        }
53    }
54
55    /// Enter endgame mode - enables duplicate request tracking
56    pub fn enter_endgame(&mut self) {
57        if !self.active {
58            self.active = true;
59            info!("[BT] === Entering endgame mode ===");
60        }
61    }
62
63    /// Exit endgame mode and clear all tracked requests
64    pub fn exit_endgame(&mut self) {
65        if self.active {
66            self.active = false;
67            self.active_duplicate_requests.clear();
68            debug!(
69                "[BT] Exiting endgame mode, cleared {} tracked requests",
70                self.active_duplicate_requests.len()
71            );
72        }
73    }
74
75    /// Register that a request was sent to a peer during endgame
76    ///
77    /// This tracks which peers have pending requests for each block so we can
78    /// cancel redundant requests when the first response arrives.
79    pub fn track_request(&mut self, piece: u32, offset: u32, len: u32, peer_id: usize) {
80        let key = (piece, offset, len);
81        self.active_duplicate_requests
82            .entry(key)
83            .or_default()
84            .push(peer_id);
85    }
86
87    /// When a block arrives, find other peers that have pending requests for the same block
88    ///
89    /// Returns the list of peer indices that should receive Cancel messages.
90    /// Does NOT remove the entry (call remove_request after sending cancels).
91    pub fn get_cancel_targets(&self, piece: u32, offset: u32, len: u32) -> Vec<usize> {
92        let key = (piece, offset, len);
93        self.active_duplicate_requests
94            .get(&key)
95            .map(|peers| peers.to_vec())
96            .unwrap_or_default()
97    }
98
99    /// Remove a tracked request after cancel or completion
100    ///
101    /// Called after Cancel messages have been sent and the block is fully processed.
102    pub fn remove_request(&mut self, piece: u32, offset: u32, len: u32) {
103        let key = (piece, offset, len);
104        self.active_duplicate_requests.remove(&key);
105    }
106
107    /// Check if endgame mode is currently active
108    pub fn is_endgame_active(&self) -> bool {
109        self.active
110    }
111
112    /// Get the number of actively tracked duplicate requests (for debugging/metrics)
113    #[allow(dead_code)] // Debugging metric; used in tests only
114    pub fn tracked_count(&self) -> usize {
115        self.active_duplicate_requests.len()
116    }
117}
118
119impl Default for EndgameState {
120    fn default() -> Self {
121        Self::new()
122    }
123}
124
125#[async_trait]
126impl Command for BtDownloadCommand {
127    async fn execute(&mut self) -> Result<()> {
128        if !self.started {
129            self.group.write().await.start().await?;
130            self.started = true;
131        }
132
133        let (meta, piece_length, total_size, num_pieces) = self.prepare_environment().await?;
134
135        // P1 集成: 尝试从 .aria2 文件恢复已保存的进度
136        if let Some(ref mgr) = self.progress_manager {
137            match mgr.load_progress(&meta.info_hash.bytes) {
138                Ok(saved) => {
139                    info!(
140                        pieces_done = saved.num_pieces,
141                        ratio = saved.completion_ratio(),
142                        "Resuming from saved progress"
143                    );
144                }
145                Err(e) => {
146                    debug!(
147                        error = %e,
148                        "No saved progress found, starting fresh download"
149                    );
150                }
151            }
152        }
153
154        let peer_addrs = self
155            .discover_peers(&meta, total_size, &meta.info_hash.bytes)
156            .await?;
157
158        if peer_addrs.is_empty() {
159            return Err(Aria2Error::Recoverable(
160                RecoverableError::TemporaryNetworkFailure {
161                    message: "No peers from tracker or DHT".into(),
162                },
163            ));
164        }
165
166        let mut active_connections = self
167            .connect_to_peers(&peer_addrs, &meta.info_hash.bytes, num_pieces)
168            .await?;
169
170        // Initialize PEX known peers list from discovered peers for BEP 11 exchange.
171        // BEP 0027 (Private Torrent): PEX must be disabled for private torrents
172        // because it exchanges peer lists with connected peers, which would leak
173        // the swarm membership beyond the tracker-controlled peer set.
174        if self.is_private {
175            info!("[BT] Private torrent: PEX disabled (BEP 0027)");
176        } else {
177            let pex_peers: Vec<aria2_protocol::bittorrent::peer::connection::PeerAddr> = peer_addrs
178                .iter()
179                .map(|pa| {
180                    aria2_protocol::bittorrent::peer::connection::PeerAddr::new(&pa.ip, pa.port)
181                })
182                .collect();
183            self.set_pex_known_peers(pex_peers);
184            info!(
185                "[PEX] Initialized with {} known peers from tracker/DHT",
186                self.pex_known_peers.len()
187            );
188        }
189
190        // Initialize web seed manager if web seeds are available (BEP 19)
191        let web_seed_manager = if !self.web_seed_urls.is_empty() {
192            info!(
193                "[BT] Initializing web seed manager with {} URL(s)",
194                self.web_seed_urls.len()
195            );
196            Some(crate::engine::bt_web_seed::WebSeedManager::new(
197                self.web_seed_urls.clone(),
198                piece_length,
199                total_size,
200            ))
201        } else {
202            None
203        };
204
205        // PEX Integration: Initialize PEX state tracking for active connections
206        // Each peer may support ut_pex extension (BEP 11) for peer discovery.
207        // BEP 0027 (Private Torrent): leave pex_enabled_peers empty so no PEX
208        // messages are ever sent.
209        let mut pex_enabled_peers: HashSet<usize> = HashSet::new();
210        let mut last_pex_send = Instant::now();
211        const PEX_SEND_INTERVAL_SECS: u64 = 60;
212
213        if !self.is_private {
214            // Check PEX support for each connection (simplified - assume all peers support PEX)
215            // In a full implementation, this would check extension handshake results
216            for (idx, _conn) in active_connections.iter().enumerate() {
217                pex_enabled_peers.insert(idx);
218            }
219            info!(
220                "[PEX] Initialized PEX tracking for {} peers (assuming ut_pex support)",
221                pex_enabled_peers.len()
222            );
223        }
224
225        self.download_pieces_loop(
226            &mut active_connections,
227            &meta,
228            piece_length,
229            total_size,
230            num_pieces,
231            web_seed_manager.as_ref(),
232            &mut pex_enabled_peers,
233            &mut last_pex_send,
234            PEX_SEND_INTERVAL_SECS,
235        )
236        .await?;
237
238        if self.seed_enabled && !active_connections.is_empty() {
239            info!(
240                "Starting seeding phase with {} peers...",
241                active_connections.len()
242            );
243            self.run_seeding_phase(active_connections, piece_length, num_pieces)
244                .await?;
245        } else {
246            info!(
247                "Skipping seeding (enabled={}, connections={})",
248                self.seed_enabled,
249                active_connections.len()
250            );
251            for conn in &mut active_connections {
252                let _ = conn;
253            }
254        }
255
256        self.finalize_download(Instant::now(), &meta).await?;
257
258        Ok(())
259    }
260
261    fn status(&self) -> CommandStatus {
262        if self.completed_bytes > 0 {
263            CommandStatus::Running
264        } else {
265            CommandStatus::Pending
266        }
267    }
268
269    fn timeout(&self) -> Option<Duration> {
270        Some(Duration::from_secs(600))
271    }
272}
273
274impl BtDownloadCommand {
275    async fn prepare_environment(
276        &mut self,
277    ) -> Result<(
278        aria2_protocol::bittorrent::torrent::parser::TorrentMeta,
279        u32,
280        u64,
281        u32,
282    )> {
283        if let Some(parent) = self.output_path.parent()
284            && !parent.exists()
285        {
286            std::fs::create_dir_all(parent).map_err(|e| {
287                Aria2Error::Fatal(FatalError::Config(format!("mkdir failed: {}", e)))
288            })?;
289        }
290
291        if let Some(ref layout) = self.multi_file_layout {
292            layout.create_directories().map_err(|e| {
293                Aria2Error::Fatal(FatalError::Config(format!(
294                    "create_directories failed: {}",
295                    e
296                )))
297            })?;
298            info!(
299                "[BT] Multi-file mode: {} files under {}",
300                layout.num_files(),
301                self.output_path.display()
302            );
303        }
304
305        let meta =
306            aria2_protocol::bittorrent::torrent::parser::TorrentMeta::parse(&self.torrent_data)
307                .map_err(|e| {
308                    Aria2Error::Fatal(FatalError::Config(format!("Torrent parse error: {}", e)))
309                })?;
310
311        {
312            let mut g = self.group.write().await;
313            g.set_total_length(meta.total_size()).await;
314            g.set_total_length_atomic(meta.total_size());
315        }
316
317        let piece_length = meta.info.piece_length;
318        let total_size = meta.total_size();
319        let num_pieces = meta.num_pieces() as u32;
320
321        Ok((meta, piece_length, total_size, num_pieces))
322    }
323
324    async fn discover_peers(
325        &mut self,
326        meta: &aria2_protocol::bittorrent::torrent::parser::TorrentMeta,
327        total_size: u64,
328        info_hash_raw: &[u8; 20],
329    ) -> Result<Vec<aria2_protocol::bittorrent::peer::connection::PeerAddr>> {
330        let my_peer_id = aria2_protocol::bittorrent::peer::id::generate_peer_id();
331        let mut peer_addrs =
332            perform_http_tracker_announce(&meta.announce, info_hash_raw, &my_peer_id, total_size)
333                .await?;
334
335        if let Ok(udp) = UdpTrackerClient::new(0).await {
336            self.udp_client = Some(std::sync::Arc::new(tokio::sync::Mutex::new(udp)));
337            if let Some(ref shared_client) = self.udp_client {
338                let mut mgr = UdpTrackerManager::new(std::sync::Arc::clone(shared_client)).await;
339                let urls: Vec<String> = meta.announce_list.iter().flatten().cloned().collect();
340                mgr.parse_tracker_urls(&urls);
341
342                if mgr.endpoint_count() > 0 {
343                    debug!("Trying {} UDP tracker endpoints", mgr.endpoint_count());
344
345                    match mgr.announce(
346                        info_hash_raw, &my_peer_id,
347                        0, total_size as i64, 0,
348                        aria2_protocol::bittorrent::tracker::udp_tracker_protocol::UdpEvent::Started,
349                        50,
350                    ).await {
351                        udp_responses if !udp_responses.is_empty() => {
352                            let udp_peers = UdpTrackerManager::collect_all_peers(&udp_responses);
353                            debug!("UDP trackers returned {} additional peers", udp_peers.len());
354                            for (ip, port) in udp_peers {
355                                peer_addrs.push(aria2_protocol::bittorrent::peer::connection::PeerAddr::new(&ip, port));
356                            }
357                        }
358                        _ => { debug!("No response from UDP trackers"); }
359                    }
360                }
361            }
362        }
363
364        if peer_addrs.is_empty() {
365            tracing::error!("[BT] ERROR: No peers from tracker");
366        }
367
368        // BEP 0027 (Private Torrent): DHT must be disabled for private torrents
369        // to prevent leaking the info_hash to the public DHT network.
370        let enable_dht = { self.group.read().await.options().enable_dht } && !self.is_private;
371        if self.is_private {
372            info!("[BT] Private torrent: DHT disabled (BEP 0027)");
373        }
374        if enable_dht && self.dht_engine.is_none() {
375            let dht_port = { self.group.read().await.options().dht_listen_port };
376            let dht_file_path = { self.group.read().await.options().dht_file_path.clone() };
377            let dht_entry_points = { self.group.read().await.options().dht_entry_point.clone() };
378
379            // Parse custom bootstrap nodes if provided
380            let bootstrap_nodes: Vec<std::net::SocketAddr> =
381                if let Some(ref entry_points) = dht_entry_points {
382                    entry_points
383                        .iter()
384                        .filter_map(|ep| ep.parse::<std::net::SocketAddr>().ok())
385                        .collect()
386                } else {
387                    vec![]
388                };
389
390            let dht_config = aria2_protocol::bittorrent::dht::engine::DhtEngineConfig {
391                port: dht_port.unwrap_or(0),
392                dht_file_path,
393                ..Default::default()
394            };
395
396            match aria2_protocol::bittorrent::dht::engine::DhtEngine::start(dht_config).await {
397                Ok(engine) => {
398                    // Add custom bootstrap nodes to routing table
399                    if !bootstrap_nodes.is_empty() {
400                        for addr in &bootstrap_nodes {
401                            use aria2_protocol::bittorrent::dht::node::DhtNode;
402                            use rand::RngCore;
403                            let mut id = [0u8; 20];
404                            rand::thread_rng().fill_bytes(&mut id);
405                            let node = DhtNode::new(id, *addr);
406                            engine.add_node(node).await;
407                        }
408                        tracing::info!(
409                            "[BT] Added {} custom DHT bootstrap nodes",
410                            bootstrap_nodes.len()
411                        );
412                    }
413
414                    self.dht_engine = Some(engine);
415                    tracing::info!("[BT] DHT engine started");
416                    if let Some(dht) = self.dht_engine.as_ref() {
417                        dht.start_maintenance_loop();
418                    }
419                }
420                Err(e) => {
421                    warn!("[BT] DHT engine start failed: {}", e);
422                }
423            }
424        }
425
426        if let Some(ref engine) = self.dht_engine {
427            let result = engine.find_peers(info_hash_raw).await;
428            if !result.peers.is_empty() {
429                let before = peer_addrs.len();
430                for addr in &result.peers {
431                    let ip_str = addr.ip().to_string();
432                    let paddr = aria2_protocol::bittorrent::peer::connection::PeerAddr::new(
433                        &ip_str,
434                        addr.port(),
435                    );
436                    if !peer_addrs
437                        .iter()
438                        .any(|p| p.ip == paddr.ip && p.port == paddr.port)
439                    {
440                        peer_addrs.push(paddr);
441                    }
442                }
443                tracing::info!(
444                    "[BT] DHT discovered {} extra peers (total: {}, contacted {} DHT nodes)",
445                    peer_addrs.len() - before,
446                    peer_addrs.len(),
447                    result.nodes_contacted
448                );
449            } else {
450                debug!("[BT] DHT find_peers returned no peers");
451            }
452        }
453
454        // BEP 0027 (Private Torrent): public tracker announcement is forbidden
455        // for private torrents because it would leak the info_hash to trackers
456        // not explicitly listed in the torrent's announce list.
457        let enable_public_trackers =
458            { self.group.read().await.options().enable_public_trackers } && !self.is_private;
459        if self.is_private {
460            info!("[BT] Private torrent: public trackers disabled (BEP 0027)");
461        }
462        if enable_public_trackers
463            && self.public_trackers.is_none()
464            && peer_addrs.len() < PUBLIC_TRACKER_PEER_THRESHOLD
465        {
466            let ptl = std::sync::Arc::new(
467                aria2_protocol::bittorrent::tracker::public_list::PublicTrackerList::new(),
468            );
469            ptl.start_auto_update(
470                "https://cf.trackerslist.com/best.txt".to_string(),
471                std::time::Duration::from_secs(86400),
472            );
473            self.public_trackers = Some(ptl);
474        }
475
476        if let Some(ref pt) = self.public_trackers {
477            let http_urls = pt.get_http_trackers().await;
478            let mut extra_peers: Vec<(String, u16)> = Vec::new();
479            let mut announced = 0usize;
480
481            for url in http_urls.iter().take(MAX_PUBLIC_TRACKERS_TO_TRY) {
482                match announce_to_public_tracker(url, info_hash_raw, &my_peer_id, total_size).await
483                {
484                    Ok(peers) => {
485                        announced += 1;
486                        extra_peers.extend(peers);
487                    }
488                    Err(e) => {
489                        debug!("[BT] Public tracker {} failed: {}", url, e);
490                    }
491                }
492            }
493
494            if !extra_peers.is_empty() {
495                let before = peer_addrs.len();
496                for (ip, port) in extra_peers {
497                    let paddr =
498                        aria2_protocol::bittorrent::peer::connection::PeerAddr::new(&ip, port);
499                    if !peer_addrs
500                        .iter()
501                        .any(|p| p.ip == paddr.ip && p.port == paddr.port)
502                    {
503                        peer_addrs.push(paddr);
504                    }
505                }
506                tracing::info!(
507                    "[BT] Public trackers discovered {} extra peers (announced to {} of {})",
508                    peer_addrs.len() - before,
509                    announced,
510                    http_urls.len()
511                );
512            } else if announced > 0 {
513                debug!("[BT] Public trackers responded but no peers found");
514            }
515        }
516
517        // P2: Integrate LPD-discovered LAN peers
518        // BEP 0027 (Private Torrent): LPD (Local Peer Discovery) uses UDP
519        // multicast which would leak the info_hash to the local network, so it
520        // must be disabled for private torrents.
521        if self.is_private {
522            if self.lpd_manager.is_some() {
523                info!("[BT] Private torrent: LPD disabled (BEP 0027)");
524            }
525        } else if let Some(ref lpd) = self.lpd_manager {
526            // Convert raw 20-byte info_hash to 40-char hex string for LPD
527            let info_hash_hex = hex::encode(*info_hash_raw);
528            let lpd_peers = lpd.get_peers_for(&info_hash_hex).await;
529            if !lpd_peers.is_empty() {
530                let before = peer_addrs.len();
531                for lpd_peer in &lpd_peers {
532                    // LpdPeer.addr is IpAddr, LpdPeer.port is u16
533                    let ip_str = lpd_peer.addr.to_string();
534                    let paddr = aria2_protocol::bittorrent::peer::connection::PeerAddr::new(
535                        &ip_str,
536                        lpd_peer.port,
537                    );
538                    if !peer_addrs
539                        .iter()
540                        .any(|p| p.ip == paddr.ip && p.port == paddr.port)
541                    {
542                        peer_addrs.push(paddr);
543                    }
544                }
545
546                info!(
547                    lpd_count = lpd_peers.len(),
548                    total_added = peer_addrs.len() - before,
549                    "LPD discovered local peers"
550                );
551
552                // Register current download for LPD announcement
553                let _ = lpd.register_torrent(&info_hash_hex).await;
554            } else {
555                debug!("LPD no local peers found for this torrent");
556            }
557        }
558
559        Ok(peer_addrs)
560    }
561
562    async fn connect_to_peers(
563        &mut self,
564        peer_addrs: &[aria2_protocol::bittorrent::peer::connection::PeerAddr],
565        info_hash_raw: &[u8; 20],
566        num_pieces: u32,
567    ) -> Result<Vec<BtPeerConn>> {
568        let require_crypto = { self.group.read().await.options().bt_require_crypto };
569        let force_encrypt = { self.group.read().await.options().bt_force_encrypt };
570
571        let conn_result = BtPeerInteraction::connect_to_peers(
572            peer_addrs,
573            info_hash_raw,
574            num_pieces,
575            require_crypto,
576            force_encrypt,
577        )
578        .await?;
579
580        let active_connections = conn_result.connections;
581
582        tracing::info!("[BT] Active connections: {}", active_connections.len());
583        if active_connections.is_empty() {
584            return Err(Aria2Error::Recoverable(
585                RecoverableError::TemporaryNetworkFailure {
586                    message: "All peer connections failed".into(),
587                },
588            ));
589        }
590
591        {
592            let options = self.group.read().await.options().clone();
593            let config = ChokingConfig {
594                max_upload_slots: options.bt_max_upload_slots.unwrap_or(4) as usize,
595                optimistic_unchoke_interval_secs: options
596                    .bt_optimistic_unchoke_interval
597                    .unwrap_or(30),
598                snubbed_timeout_secs: options.bt_snubbed_timeout.unwrap_or(60),
599                choke_rotation_interval_secs: 10,
600            };
601
602            let mut algo = ChokingAlgorithm::new(config);
603
604            for addr in peer_addrs {
605                let socket_addr = std::net::SocketAddr::new(
606                    addr.ip.parse().unwrap_or_else(|_| {
607                        std::net::IpAddr::V4(std::net::Ipv4Addr::new(127, 0, 0, 1))
608                    }),
609                    addr.port,
610                );
611                let peer_stats = PeerStats::new([0u8; 20], socket_addr);
612                algo.add_peer(peer_stats);
613            }
614
615            self.choking_algo = Some(algo);
616            tracing::info!(
617                "[BT] Choking algorithm initialized with {} peers",
618                self.choking_algo.as_ref().unwrap().len()
619            );
620        }
621
622        Ok(active_connections)
623    }
624
625    // Parameters are individually meaningful; grouping into a struct would
626    // reduce clarity for this inner download loop.
627    #[allow(clippy::too_many_arguments)]
628    async fn download_pieces_loop(
629        &mut self,
630        active_connections: &mut [BtPeerConn],
631        meta: &aria2_protocol::bittorrent::torrent::parser::TorrentMeta,
632        piece_length: u32,
633        total_size: u64,
634        num_pieces: u32,
635        web_seed_manager: Option<&crate::engine::bt_web_seed::WebSeedManager>,
636        pex_enabled_peers: &mut HashSet<usize>,
637        last_pex_send: &mut Instant,
638        pex_send_interval_secs: u64,
639    ) -> Result<()> {
640        let raw_writer = DefaultDiskWriter::new(&self.output_path);
641        let rate_limit = {
642            let g = self.group.read().await;
643            g.options().max_download_limit
644        };
645        let mut writer: Box<dyn DiskWriter> = match rate_limit {
646            Some(rate) if rate > 0 => Box::new(ThrottledWriter::new(
647                raw_writer,
648                RateLimiter::new(&RateLimiterConfig::new(Some(rate), None)),
649            )),
650            _ => Box::new(raw_writer),
651        };
652        let start_time = Instant::now();
653        let mut last_speed_update = Instant::now();
654        let mut last_completed = 0u64;
655
656        // P1 集成: 进度保存时间追踪
657        let mut last_progress_save = Instant::now();
658
659        let piece_selector = BtPieceSelector::new(num_pieces);
660
661        let mut piece_manager = aria2_protocol::bittorrent::piece::manager::PieceManager::new(
662            num_pieces,
663            piece_length,
664            total_size,
665            meta.info.pieces.clone(),
666        );
667
668        let mut piece_picker =
669            aria2_protocol::bittorrent::piece::picker::PiecePicker::new(num_pieces);
670        piece_picker.set_strategy(
671            aria2_protocol::bittorrent::piece::picker::PieceSelectionStrategy::Sequential,
672        );
673
674        // G2: Set piece priority mode from config option (--bt-prioritize-piece)
675        let prioritize_piece_mode = {
676            let g = self.group.read().await;
677            g.options().bt_prioritize_piece.clone()
678        };
679        match prioritize_piece_mode.as_str() {
680            "head" => {
681                piece_picker.set_priority_mode(
682                    aria2_protocol::bittorrent::piece::picker::PiecePriorityMode::SequentialHead,
683                );
684                info!("[BT] Piece priority mode: SequentialHead (from start)");
685            }
686            "tail" => {
687                piece_picker.set_priority_mode(
688                    aria2_protocol::bittorrent::piece::picker::PiecePriorityMode::SequentialTail,
689                );
690                info!("[BT] Piece priority mode: SequentialTail (from end)");
691            }
692            _ => {
693                // Default: RarestFirst (also handles "rarest" and empty string)
694                piece_picker.set_priority_mode(
695                    aria2_protocol::bittorrent::piece::picker::PiecePriorityMode::RarestFirst,
696                );
697                info!("[BT] Piece priority mode: RarestFirst (default)");
698            }
699        }
700
701        let mut peer_tracker = PeerBitfieldTracker::new(num_pieces);
702        BtPeerInteraction::initialize_peer_tracking(
703            active_connections,
704            num_pieces,
705            &mut peer_tracker,
706        );
707
708        piece_selector.initialize_frequencies(&mut piece_picker, &peer_tracker);
709
710        tracing::info!(
711            "[BT] Piece selection strategy: {:?}, {} pieces total, {} peers tracked",
712            piece_picker.priority_mode(),
713            num_pieces,
714            peer_tracker.peer_count()
715        );
716
717        // Phase 14 - B1: Initialize endgame state for this download session
718        let mut endgame_state = EndgameState::new();
719
720        // G1: Snub detection state - track last data received time per peer index
721        const SNUB_TIMEOUT_SECS: u64 = 30;
722        let mut peer_last_data_time: std::collections::HashMap<usize, Instant> =
723            std::collections::HashMap::new();
724        let mut last_snub_check = Instant::now();
725        const SNUB_CHECK_INTERVAL_SECS: u64 = 10;
726
727        // Initialize last-data-time tracking for all active peers
728        for (idx, _conn) in active_connections.iter().enumerate() {
729            peer_last_data_time.insert(idx, Instant::now());
730        }
731
732        loop {
733            if BtPieceSelector::is_complete(&piece_picker) {
734                // Exit endgame mode when download is complete
735                if endgame_state.is_endgame_active() {
736                    endgame_state.exit_endgame();
737                }
738                break;
739            }
740
741            // Phase 14 - B1: Check if we should enter endgame mode
742            // Endgame activates when <=5 pieces remain incomplete (matching picker threshold)
743            let endgame_candidates = piece_picker.endgame_candidates();
744            if !endgame_candidates.is_empty() && !endgame_state.is_endgame_active() {
745                endgame_state.enter_endgame();
746                info!(
747                    "[BT] Endgame mode activated: {}/{} pieces remaining",
748                    endgame_candidates.len(),
749                    num_pieces
750                );
751            } else if endgame_candidates.is_empty() && endgame_state.is_endgame_active() {
752                // This shouldn't normally happen (candidates empty means >5 remaining or all done)
753                // but handle gracefully
754                endgame_state.exit_endgame();
755            }
756
757            // G1: Periodic snub detection - check all peers for inactivity
758            if last_snub_check.elapsed().as_secs() >= SNUB_CHECK_INTERVAL_SECS {
759                last_snub_check = Instant::now();
760                let mut newly_snubbed = Vec::new();
761                for (&peer_id, &last_time) in &peer_last_data_time {
762                    if last_time.elapsed().as_secs() > SNUB_TIMEOUT_SECS {
763                        // Mark peer as snubbed via the command's choking algorithm
764                        self.mark_peer_snubbed(peer_id);
765                        newly_snubbed.push(peer_id);
766                        debug!(
767                            "[BT] Peer {} marked as snubbed (no data for {}s)",
768                            peer_id,
769                            last_time.elapsed().as_secs()
770                        );
771                    }
772                }
773                if !newly_snubbed.is_empty() {
774                    debug!(
775                        "[BT] Snub check: {} peers newly snubbed",
776                        newly_snubbed.len()
777                    );
778                }
779
780                // Also run the PeerStats-level snub check (timeout-based)
781                let stats_snubbed = self.check_snubbed_peers();
782                if !stats_snubbed.is_empty() {
783                    debug!(
784                        "[BT] PeerStats snub check: {} peers timed out",
785                        stats_snubbed.len()
786                    );
787                }
788            }
789
790            // PEX Integration: Periodic PEX message sending (BEP 11)
791            // Send PEX messages to peers that support ut_pex every 60 seconds
792            if last_pex_send.elapsed().as_secs() >= pex_send_interval_secs
793                && !pex_enabled_peers.is_empty()
794                && !self.pex_known_peers.is_empty()
795            {
796                *last_pex_send = Instant::now();
797                let pex_peers_count = self.pex_known_peers.len();
798
799                for peer_idx in pex_enabled_peers.iter() {
800                    if let Some(_conn) = active_connections.get(*peer_idx) {
801                        // Build PEX message for this peer
802                        let _remote_addr =
803                            aria2_protocol::bittorrent::peer::connection::PeerAddr::new(
804                                "0.0.0.0", // Placeholder - actual address would come from connection
805                                0,
806                            );
807
808                        // Note: In a full implementation, we would:
809                        // 1. Get the actual remote address from the connection
810                        // 2. Send the PEX extension message via the connection
811                        // 3. Handle incoming PEX messages in read_message loop
812
813                        debug!(
814                            "[PEX] Would send PEX to peer {} ({} known peers available)",
815                            peer_idx, pex_peers_count
816                        );
817                    }
818                }
819
820                info!(
821                    "[PEX] Periodic PEX exchange triggered: {} peers enabled, {} known peers",
822                    pex_enabled_peers.len(),
823                    pex_peers_count
824                );
825            }
826
827            let remaining = piece_picker.remaining_count();
828
829            let selection = piece_selector.select_next_piece(&mut piece_picker, remaining as usize);
830
831            let next_piece_idx = match selection.piece_index {
832                Some(idx) => idx,
833                None => {
834                    tracing::debug!("[BT] No piece available, waiting...");
835                    tokio::time::sleep(Duration::from_millis(PEER_CONNECTION_DELAY_MS)).await;
836                    continue;
837                }
838            };
839
840            tracing::info!("[BT] Downloading piece {}...", next_piece_idx);
841
842            let actual_piece_len =
843                piece_selector.calculate_piece_length(next_piece_idx, piece_length, total_size);
844
845            let num_blocks = BtPieceSelector::calculate_num_blocks(actual_piece_len, BLOCK_SIZE);
846            tracing::debug!(
847                "[BT] Piece {} has {} blocks (size: {} bytes)",
848                next_piece_idx,
849                num_blocks,
850                actual_piece_len
851            );
852            let mut piece_ok = false;
853
854            // Phase 14 - B1: Use endgame-aware download when in endgame mode
855            let download_result = if endgame_state.is_endgame_active() {
856                info!(
857                    "[BT] Endgame: downloading piece {} with duplicate requests ({} peers available)",
858                    next_piece_idx,
859                    active_connections.len()
860                );
861                BtMessageHandler::download_piece_blocks_endgame(
862                    active_connections,
863                    next_piece_idx as u32,
864                    actual_piece_len,
865                    num_blocks,
866                    &mut endgame_state,
867                )
868                .await
869            } else {
870                BtMessageHandler::download_piece_blocks(
871                    active_connections,
872                    next_piece_idx as u32,
873                    actual_piece_len,
874                    num_blocks,
875                )
876                .await
877            };
878
879            match download_result {
880                Ok(piece_data) => {
881                    self.completed_bytes += piece_data.len() as u64;
882
883                    // G1: Update last-data-time for all active peers on successful receive
884                    // (In a full implementation, this would be per-peer; here we update all
885                    // since the download loop processes pieces from any peer)
886                    for idx in 0..active_connections.len() {
887                        peer_last_data_time.insert(idx, Instant::now());
888                    }
889
890                    tracing::info!(
891                        "[BT] All blocks received for piece {}, verifying...",
892                        next_piece_idx
893                    );
894                    if piece_manager.verify_piece_hash(next_piece_idx as u32, &piece_data) {
895                        tracing::info!("[BT] Piece {} verified OK", next_piece_idx);
896                        piece_manager.mark_piece_complete(next_piece_idx as u32);
897                        piece_picker.mark_completed(next_piece_idx as u32);
898
899                        if let Some(ref layout) = self.multi_file_layout {
900                            // Phase 14 / I4: use coalesced writer to reduce syscalls
901                            write_piece_to_multi_files_coalesced(
902                                layout,
903                                next_piece_idx as u32,
904                                &piece_data,
905                                layout.piece_length(),
906                            )
907                            .await?;
908                        } else {
909                            writer.write(&piece_data).await?;
910                        }
911
912                        // Sync bitfield to RequestGroup for session persistence
913                        {
914                            let bitfield = piece_picker.export_bitfield();
915                            let g = self.group.read().await;
916                            g.set_bt_bitfield(Some(bitfield)).await;
917                        }
918
919                        BtPeerInteraction::broadcast_have(
920                            active_connections,
921                            next_piece_idx as u32,
922                        )
923                        .await;
924                        piece_ok = true;
925
926                        // PEX Integration: Trigger PEX send on piece completion
927                        // This ensures peers are exchanged when progress is made
928                        if !self.pex_known_peers.is_empty() && self.should_send_pex() {
929                            let dummy_remote =
930                                aria2_protocol::bittorrent::peer::connection::PeerAddr::new(
931                                    "0.0.0.0", 0,
932                                );
933                            if let Some(_pex_data) = self.maybe_send_pex(&dummy_remote) {
934                                debug!(
935                                    "[PEX] PEX message ready after piece {} completion",
936                                    next_piece_idx
937                                );
938                                // Note: Actual sending would happen via extension message channel
939                                // when full extension protocol integration is implemented
940                            }
941                        }
942
943                        // P1 集成: 定期保存下载进度到 .aria2 文件
944                        if let Some(ref mgr) = self.progress_manager
945                            && last_progress_save.elapsed() >= self.progress_save_interval
946                        {
947                            // 构造当前进度快照
948                            let progress = BtProgress {
949                                info_hash: meta.info_hash.bytes,
950                                bitfield: vec![], // 将在后续完善
951                                peers: vec![],
952                                stats: ProgressDownloadStats {
953                                    downloaded_bytes: self.completed_bytes,
954                                    uploaded_bytes: self.total_uploaded,
955                                    upload_speed: 0.0,
956                                    download_speed: 0.0,
957                                    elapsed_seconds: start_time.elapsed().as_secs(),
958                                },
959                                piece_length,
960                                total_size,
961                                num_pieces,
962                                save_time: std::time::SystemTime::now(),
963                                version: 1,
964                            };
965
966                            if let Err(e) = mgr.save_progress(&meta.info_hash.bytes, &progress) {
967                                warn!(
968                                    error = %e,
969                                    "Failed to save BT progress"
970                                );
971                            } else {
972                                debug!(
973                                    pieces_completed = next_piece_idx + 1,
974                                    total_pieces = num_pieces,
975                                    "BT progress saved successfully"
976                                );
977                            }
978                            last_progress_save = Instant::now();
979                        }
980                    } else {
981                        tracing::warn!(
982                            "[BT] SHA1 mismatch on piece {}, retrying...",
983                            next_piece_idx
984                        );
985
986                        // H3: Bad peer detection - increment bad data counter for peers
987                        // that contributed to this failed piece.
988                        // Note: In a full implementation, we would track which specific peer
989                        // sent each block and only penalize that peer. For now, we log
990                        // the event and the caller can use record_bad_piece_for_peer() if
991                        // they know which peer sent the invalid data.
992                        tracing::warn!(
993                            "[BT] Piece {} hash verification FAILED - potential bad peer detected",
994                            next_piece_idx
995                        );
996                    }
997                }
998                Err(_) => {
999                    tracing::warn!(
1000                        "[BT] Incomplete piece {}, needed {} blocks",
1001                        next_piece_idx,
1002                        num_blocks
1003                    );
1004                }
1005            }
1006
1007            if !piece_ok {
1008                // Try Web Seeds as fallback (BEP 19)
1009                if let Some(ws_mgr) = web_seed_manager {
1010                    info!(
1011                        "[BT] Piece {} failed from peers, trying web seeds...",
1012                        next_piece_idx
1013                    );
1014
1015                    match ws_mgr.request_piece(next_piece_idx as u32).await {
1016                        Ok(web_seed_data) => {
1017                            info!(
1018                                "[BT] Piece {} downloaded from web seed ({} bytes)",
1019                                next_piece_idx,
1020                                web_seed_data.len()
1021                            );
1022                            // Verify hash
1023                            if piece_manager
1024                                .verify_piece_hash(next_piece_idx as u32, &web_seed_data)
1025                            {
1026                                tracing::info!(
1027                                    "[BT] Piece {} from web seed verified OK",
1028                                    next_piece_idx
1029                                );
1030                                piece_manager.mark_piece_complete(next_piece_idx as u32);
1031                                piece_picker.mark_completed(next_piece_idx as u32);
1032
1033                                if let Some(ref layout) = self.multi_file_layout {
1034                                    write_piece_to_multi_files_coalesced(
1035                                        layout,
1036                                        next_piece_idx as u32,
1037                                        &web_seed_data,
1038                                        layout.piece_length(),
1039                                    )
1040                                    .await?;
1041                                } else {
1042                                    writer.write(&web_seed_data).await.ok();
1043                                }
1044
1045                                // Sync bitfield to RequestGroup for session persistence (Task 4)
1046                                {
1047                                    let bitfield = piece_picker.export_bitfield();
1048                                    let g = self.group.read().await;
1049                                    g.set_bt_bitfield(Some(bitfield)).await;
1050                                }
1051
1052                                self.completed_bytes += web_seed_data.len() as u64;
1053                                piece_ok = true;
1054                            } else {
1055                                tracing::warn!(
1056                                    "[BT] Piece {} from web seed failed hash verification",
1057                                    next_piece_idx
1058                                );
1059                            }
1060                        }
1061                        Err(e) => {
1062                            tracing::warn!(
1063                                "[BT] Web seed download failed for piece {}: {}",
1064                                next_piece_idx,
1065                                e
1066                            );
1067                        }
1068                    }
1069                }
1070
1071                if !piece_ok {
1072                    tracing::error!(
1073                        "[BT] Piece {} failed after {} retries (peers and web seeds)",
1074                        next_piece_idx,
1075                        MAX_RETRIES
1076                    );
1077                    return Err(Aria2Error::Fatal(FatalError::Config(format!(
1078                        "Piece {} download failed after {} retries",
1079                        next_piece_idx, MAX_RETRIES
1080                    ))));
1081                }
1082            }
1083
1084            {
1085                let g = self.group.write().await;
1086                g.update_progress(self.completed_bytes).await;
1087                g.set_completed_length(self.completed_bytes);
1088
1089                let elapsed = last_speed_update.elapsed();
1090                if elapsed.as_millis() >= 500 {
1091                    let delta = self.completed_bytes - last_completed;
1092                    let speed = (delta as f64 / elapsed.as_secs_f64()) as u64;
1093                    g.update_speed(speed, 0).await;
1094                    g.set_download_speed_cached(speed);
1095                    last_speed_update = Instant::now();
1096                    last_completed = self.completed_bytes;
1097                }
1098            }
1099        }
1100
1101        tracing::info!("[BT] Finalizing writer...");
1102        writer.finalize().await.ok();
1103        tracing::info!("[BT] Writer finalized OK");
1104        info!(
1105            "BT download done: {} ({} bytes)",
1106            self.output_path.display(),
1107            self.completed_bytes
1108        );
1109
1110        Ok(())
1111    }
1112
1113    async fn finalize_download(
1114        &mut self,
1115        start_time: Instant,
1116        meta: &aria2_protocol::bittorrent::torrent::parser::TorrentMeta,
1117    ) -> Result<()> {
1118        let final_speed = {
1119            let elapsed = start_time.elapsed().as_secs_f64();
1120            if elapsed > 0.0 {
1121                (self.completed_bytes as f64 / elapsed) as u64
1122            } else {
1123                0
1124            }
1125        };
1126        {
1127            let mut g = self.group.write().await;
1128            g.update_progress(self.completed_bytes).await;
1129            g.update_speed(final_speed, self.total_uploaded).await;
1130            g.set_completed_length(self.completed_bytes);
1131            g.set_download_speed_cached(final_speed);
1132            g.set_uploaded_length(self.total_uploaded);
1133            g.complete().await?;
1134        }
1135
1136        info!(
1137            "BT command done: downloaded={} uploaded={}",
1138            self.completed_bytes, self.total_uploaded
1139        );
1140
1141        if let Some(ref engine) = self.dht_engine {
1142            if let Err(e) = engine.announce_peer(&meta.info_hash.bytes, 0).await {
1143                warn!("[BT] DHT announce failed: {}", e);
1144            } else {
1145                info!(
1146                    "[BT] DHT announce_peer sent for {}",
1147                    meta.info_hash.as_hex()
1148                );
1149            }
1150            engine.shutdown();
1151        }
1152
1153        // P1 集成: 清理已完成的下载进度文件
1154        if let Some(ref mgr) = self.progress_manager {
1155            if let Err(e) = mgr.remove_progress(&meta.info_hash.bytes) {
1156                warn!(
1157                    error = %e,
1158                    "Failed to remove progress file after completion"
1159                );
1160            } else {
1161                info!("BT progress file removed after successful download");
1162            }
1163        }
1164
1165        // P2 集成: 触发下载完成后处理钩子
1166        if let Some(ref hm) = self.hook_manager {
1167            // 获取 gid(从 group 中提取)
1168            let gid = {
1169                let g = self.group.read().await;
1170                g.gid()
1171            };
1172
1173            let ctx = HookContext {
1174                gid,
1175                file_path: self.output_path.clone(),
1176                status: DownloadStatus::Complete,
1177                stats: crate::engine::bt_post_download_handler::DownloadStats {
1178                    uploaded_bytes: self.total_uploaded,
1179                    downloaded_bytes: self.completed_bytes,
1180                    upload_speed: 0.0,
1181                    download_speed: final_speed as f64,
1182                    elapsed_seconds: start_time.elapsed().as_secs(),
1183                },
1184                error: None,
1185            };
1186
1187            match hm.fire_complete(&ctx).await {
1188                Ok(results) => {
1189                    info!(
1190                        hook_count = results.len(),
1191                        "All post-download hooks executed successfully"
1192                    );
1193                    for result in &results {
1194                        debug!(result = %result, "Hook execution result");
1195                    }
1196                }
1197                Err(e) => {
1198                    warn!(
1199                        error = %e,
1200                        "Post-download hook execution failed (non-fatal)"
1201                    );
1202                }
1203            }
1204        }
1205
1206        Ok(())
1207    }
1208
1209    /// Check if both local and remote peer support ut_pex extension
1210    #[allow(dead_code)] // PEX support check; not yet called from production download loop
1211    pub fn check_pex_support(
1212        local_extension_ids: &[Option<u8>],
1213        remote_extension_ids: &[Option<u8>],
1214    ) -> bool {
1215        let local_supports = local_extension_ids.contains(&Some(PexHandler::EXTENSION_ID));
1216        let remote_supports = remote_extension_ids.contains(&Some(PexHandler::EXTENSION_ID));
1217        local_supports && remote_supports
1218    }
1219
1220    /// Build and optionally send a PEX message to connected peers
1221    /// Returns the encoded PEX message (or None if not ready to send)
1222    pub fn maybe_send_pex(
1223        &mut self,
1224        remote_peer_addr: &aria2_protocol::bittorrent::peer::connection::PeerAddr,
1225    ) -> Option<Vec<u8>> {
1226        // BEP 0027 (Private Torrent): PEX must never be sent for private
1227        // torrents. This is a defense-in-depth guard; the download loop also
1228        // leaves pex_known_peers empty for private torrents.
1229        if self.is_private {
1230            return None;
1231        }
1232
1233        if !self.should_send_pex() {
1234            return None;
1235        }
1236
1237        if self.pex_known_peers.is_empty() {
1238            debug!("[PEX] No known peers to exchange");
1239            return None;
1240        }
1241
1242        debug!(
1243            known_peers = self.pex_known_peers.len(),
1244            remote = %format!("{}:{}", remote_peer_addr.ip, remote_peer_addr.port),
1245            "[PEX] Building PEX message"
1246        );
1247
1248        let pex_msg = PexHandler::build_pex_added(
1249            &self.pex_known_peers,
1250            remote_peer_addr,
1251            PexHandler::DEFAULT_MAX_PEERS,
1252        );
1253
1254        let encoded = pex_msg.encode();
1255        self.update_pex_last_send();
1256
1257        debug!(
1258            size = encoded.len(),
1259            "[PEX] PEX message built and ready to send"
1260        );
1261        Some(encoded)
1262    }
1263
1264    /// Process an incoming PEX message and extract discovered/dropped peers
1265    pub fn handle_incoming_pex(
1266        &mut self,
1267        pex_data: &[u8],
1268        local_addr: &aria2_protocol::bittorrent::peer::connection::PeerAddr,
1269    ) -> Result<(
1270        Vec<aria2_protocol::bittorrent::peer::connection::PeerAddr>,
1271        Vec<aria2_protocol::bittorrent::peer::connection::PeerAddr>,
1272    )> {
1273        // BEP 0027 (Private Torrent): ignore any incoming PEX message for
1274        // private torrents. We must not incorporate peers learned through PEX
1275        // because the swarm is supposed to be tracker-controlled only.
1276        if self.is_private {
1277            debug!("[PEX] Ignoring incoming PEX message for private torrent (BEP 0027)");
1278            return Ok((Vec::new(), Vec::new()));
1279        }
1280
1281        match PexHandler::process_received_pex(pex_data, local_addr) {
1282            Ok((added, dropped)) => {
1283                if !added.is_empty() {
1284                    info!(count = added.len(), "[PEX] Discovered new peers from PEX");
1285                    for peer in &added {
1286                        self.add_pex_peer(peer.clone());
1287                    }
1288                }
1289                if !dropped.is_empty() {
1290                    debug!(count = dropped.len(), "[PEX] Peers to drop from PEX");
1291                }
1292                Ok((added, dropped))
1293            }
1294            Err(e) => {
1295                warn!(error = %e, "[PEX] Failed to process incoming PEX message");
1296                Err(Aria2Error::Recoverable(
1297                    RecoverableError::TemporaryNetworkFailure {
1298                        message: format!("PEX processing failed: {}", e),
1299                    },
1300                ))
1301            }
1302        }
1303    }
1304
1305    /// Connect to peers discovered via PEX
1306    ///
1307    /// This method attempts to establish connections with peers that were
1308    /// discovered through PEX (Peer Exchange, BEP 11). It's called when
1309    /// new peers are added to the PEX known peers list.
1310    ///
1311    /// # Arguments
1312    /// * `new_peers` - List of peer addresses discovered via PEX
1313    /// * `info_hash_raw` - Torrent info hash for handshake
1314    /// * `num_pieces` - Total number of pieces for bitfield size
1315    /// * `active_connections` - Current active connections (to avoid duplicates)
1316    ///
1317    /// # Returns
1318    /// * Number of successfully connected new peers
1319    pub async fn connect_to_pex_discovered_peers(
1320        &mut self,
1321        new_peers: &[aria2_protocol::bittorrent::peer::connection::PeerAddr],
1322        _info_hash_raw: &[u8; 20],
1323        _num_pieces: u32,
1324        active_connections: &[BtPeerConn],
1325    ) -> usize {
1326        // BEP 0027 (Private Torrent): never connect to peers discovered via PEX
1327        // for private torrents.
1328        if self.is_private || new_peers.is_empty() {
1329            return 0;
1330        }
1331
1332        // Filter out peers we're already connected to
1333        let already_connected: HashSet<(String, u16)> = active_connections
1334            .iter()
1335            .filter_map(|_conn| {
1336                // In a full implementation, we'd get the actual remote address
1337                // For now, we use a placeholder check
1338                None
1339            })
1340            .collect();
1341
1342        let peers_to_connect: Vec<aria2_protocol::bittorrent::peer::connection::PeerAddr> =
1343            new_peers
1344                .iter()
1345                .filter(|peer| !already_connected.contains(&(peer.ip.clone(), peer.port)))
1346                .take(10) // Limit to 10 new connections per PEX batch
1347                .cloned()
1348                .collect();
1349
1350        if peers_to_connect.is_empty() {
1351            debug!("[PEX] All discovered peers already connected");
1352            return 0;
1353        }
1354
1355        info!(
1356            "[PEX] Attempting to connect to {} new peers discovered via PEX",
1357            peers_to_connect.len()
1358        );
1359
1360        // Note: In a full implementation, this would:
1361        // 1. Use BtPeerInteraction::connect_to_peers to establish connections
1362        // 2. Add successful connections to active_connections
1363        // 3. Update pex_enabled_peers for new connections
1364
1365        // For now, we log the intent and return the count
1366        for peer in &peers_to_connect {
1367            debug!("[PEX] Would connect to peer {}:{}", peer.ip, peer.port);
1368        }
1369
1370        peers_to_connect.len()
1371    }
1372
1373    // ==================== BEP 6 Fast Extension (AllowedFast / Suggest) ====================
1374
1375    /// Maximum number of AllowedFast messages to send to a single peer
1376    const MAX_ALLOWED_FAST_PER_PEER: usize = 10;
1377
1378    /// Maximum number of Suggest messages to send per session per peer
1379    const MAX_SUGGEST_PER_PEER: usize = 5;
1380
1381    /// Check if a bitfield has a specific piece index set
1382    ///
1383    /// BitTorrent bitfields use MSB-first ordering within each byte.
1384    #[allow(dead_code)] // BEP 6 utility; used in tests only
1385    fn is_bitfield_set(bitfield: &[u8], piece_index: u32) -> bool {
1386        let byte_idx = (piece_index as usize) / 8;
1387        let bit_idx = 7 - ((piece_index as usize) % 8);
1388
1389        if byte_idx >= bitfield.len() {
1390            return false;
1391        }
1392
1393        (bitfield[byte_idx] & (1 << bit_idx)) != 0
1394    }
1395
1396    /// Calculate the set of pieces to send as AllowedFast to a peer
1397    ///
1398    /// Selects up to `MAX_ALLOWED_FAST_PER_PEER` pieces that:
1399    /// - We still need (not completed)
1400    /// - The peer has (based on their bitfield)
1401    /// - We haven't already sent AllowedFast for
1402    #[allow(dead_code)] // BEP 6 utility; used in tests only
1403    fn calculate_fast_set(
1404        needed_pieces: &[u32],
1405        peer_bitfield: &[u8],
1406        already_sent: &HashSet<u32>,
1407    ) -> Vec<u32> {
1408        let mut fast_set = Vec::new();
1409
1410        for &piece_idx in needed_pieces.iter() {
1411            if fast_set.len() >= Self::MAX_ALLOWED_FAST_PER_PEER {
1412                break;
1413            }
1414            if already_sent.contains(&piece_idx) {
1415                continue;
1416            }
1417
1418            // Check if peer has this piece (bitfield check)
1419            if Self::is_bitfield_set(peer_bitfield, piece_idx) {
1420                fast_set.push(piece_idx);
1421            }
1422        }
1423
1424        fast_set
1425    }
1426
1427    /// Send AllowedFast messages to a peer that supports BEP 6 Fast Extension
1428    ///
1429    /// This should be called after the extension handshake completes and we've received
1430    /// the peer's bitfield. It allows us to request specific pieces even when choked.
1431    #[allow(dead_code)] // BEP 6 method; not yet called from production download loop
1432    async fn send_allowed_fast_to_peer(
1433        peer_conn: &mut BtPeerConn,
1434        needed_pieces: &[u32],
1435        peer_bitfield: &[u8],
1436        already_sent: &mut HashSet<u32>,
1437    ) -> Result<usize> {
1438        let fast_set = Self::calculate_fast_set(needed_pieces, peer_bitfield, already_sent);
1439        let count = fast_set.len();
1440
1441        for piece_idx in fast_set {
1442            let _msg = serializer::serialize_allowed_fast(piece_idx);
1443
1444            // Note: In a full implementation, this would use a proper message queue/channel.
1445            // For now, we log and track what would be sent.
1446            debug!("[BEP6] Would send AllowedFast for piece {}", piece_idx);
1447
1448            already_sent.insert(piece_idx);
1449            peer_conn.add_allowed_fast(piece_idx);
1450        }
1451
1452        if count > 0 {
1453            info!("[BEP6] Sent {} AllowedFast messages to peer", count);
1454        }
1455
1456        Ok(count)
1457    }
1458
1459    /// Initialize BEP 6 tracking structures for all active connections
1460    #[allow(dead_code)]
1461    fn init_bep6_tracking(&mut self, num_connections: usize) {
1462        self.allowed_fast_sent_peers = HashMap::with_capacity(num_connections);
1463        self.suggest_sent_counts = HashMap::with_capacity(num_connections);
1464    }
1465
1466    /// Send AllowedFast messages to all peers after handshake/bitfield exchange
1467    ///
1468    /// This is called once during initialization to establish fast extension
1469    /// support with compatible peers.
1470    #[allow(dead_code)]
1471    async fn broadcast_allowed_fast(
1472        &mut self,
1473        active_connections: &mut [BtPeerConn],
1474        needed_pieces: &[u32],
1475        peer_bitfields: &[Vec<u8>],
1476    ) -> Result<u64> {
1477        self.init_bep6_tracking(active_connections.len());
1478
1479        let mut total_sent = 0u64;
1480
1481        for (idx, conn) in active_connections.iter_mut().enumerate() {
1482            let peer_bf = if idx < peer_bitfields.len() {
1483                &peer_bitfields[idx]
1484            } else {
1485                continue;
1486            };
1487
1488            let mut sent_for_peer = HashSet::new();
1489            match Self::send_allowed_fast_to_peer(conn, needed_pieces, peer_bf, &mut sent_for_peer)
1490                .await
1491            {
1492                Ok(count) => {
1493                    total_sent += count as u64;
1494                    if !sent_for_peer.is_empty() {
1495                        self.allowed_fast_sent_peers.insert(idx, sent_for_peer);
1496                    }
1497                }
1498                Err(e) => {
1499                    warn!("[BEP6] Failed to send AllowedFast to peer {}: {}", idx, e);
1500                }
1501            }
1502        }
1503
1504        if total_sent > 0 {
1505            info!(
1506                "[BEP6] Broadcast {} total AllowedFast messages to {} peers",
1507                total_sent,
1508                active_connections.len()
1509            );
1510        }
1511
1512        Ok(total_sent)
1513    }
1514
1515    /// Send Suggest messages to a peer to guide them toward pieces we need most
1516    ///
1517    /// Called after unchoking a peer, this sends up to `MAX_SUGGEST_PER_PEER` Suggest
1518    /// messages for high-priority, low-availability pieces we need urgently.
1519    ///
1520    /// # Arguments
1521    /// * `peer_idx` - Index of the peer in active_connections
1522    /// * `piece_picker` - The piece picker for selecting which pieces to suggest
1523    #[allow(dead_code)] // BEP 6 method; not yet called from production download loop
1524    async fn send_suggest_to_peer(
1525        &mut self,
1526        peer_idx: usize,
1527        piece_picker: &aria2_protocol::bittorrent::piece::picker::PiecePicker,
1528    ) -> Result<usize> {
1529        // Check if we've already sent too many suggests to this peer
1530        let sent_count = self
1531            .suggest_sent_counts
1532            .get(&peer_idx)
1533            .copied()
1534            .unwrap_or(0);
1535        if sent_count >= Self::MAX_SUGGEST_PER_PEER {
1536            debug!(
1537                "[BEP6] Already sent {} suggests to peer {}, skipping",
1538                sent_count, peer_idx
1539            );
1540            return Ok(0);
1541        }
1542
1543        let remaining = Self::MAX_SUGGEST_PER_PEER - sent_count;
1544
1545        // Select high-priority, low-availability pieces we need most urgently
1546        let mut suggestions: Vec<u32> = piece_picker
1547            .pieces_iter()
1548            .filter(|p| !p.completed && !p.in_progress && p.frequency > 0)
1549            .take(remaining)
1550            .map(|p| p.index)
1551            .collect();
1552
1553        // Sort by priority (highest first), then by rarity (lowest frequency)
1554        suggestions.sort_by(|&a, &b| {
1555            let pa = piece_picker.get_piece_info(a).unwrap();
1556            let pb = piece_picker.get_piece_info(b).unwrap();
1557            pb.priority
1558                .cmp(&pa.priority) // Higher priority first
1559                .then(pa.frequency.cmp(&pb.frequency)) // Then rarer
1560        });
1561
1562        let count = suggestions.len();
1563
1564        for piece_idx in suggestions {
1565            let _msg = serializer::serialize_suggest(piece_idx);
1566
1567            // Note: In a full implementation, this would use a proper message queue/channel.
1568            debug!(
1569                "[BEP6] Would send Suggest for piece {} to peer {}",
1570                piece_idx, peer_idx
1571            );
1572        }
1573
1574        if count > 0 {
1575            // Update suggest count for this peer
1576            let new_count = sent_count + count;
1577            self.suggest_sent_counts.insert(peer_idx, new_count);
1578
1579            info!(
1580                "[BEP6] Sent {} Suggest messages to peer {} (total: {})",
1581                count, peer_idx, new_count
1582            );
1583        }
1584
1585        Ok(count)
1586    }
1587}
1588
1589#[cfg(test)]
1590mod tests {
1591    use super::*;
1592
1593    #[test]
1594    fn test_endgame_state_new_is_inactive() {
1595        let es = EndgameState::new();
1596        assert!(!es.is_endgame_active());
1597        assert_eq!(es.tracked_count(), 0);
1598    }
1599
1600    #[test]
1601    fn test_endgame_state_default_is_inactive() {
1602        let es = EndgameState::default();
1603        assert!(!es.is_endgame_active());
1604    }
1605
1606    #[test]
1607    fn test_endgame_enter_and_exit() {
1608        let mut es = EndgameState::new();
1609        assert!(!es.is_endgame_active());
1610
1611        es.enter_endgame();
1612        assert!(es.is_endgame_active());
1613
1614        // Double enter should be idempotent
1615        es.enter_endgame();
1616        assert!(es.is_endgame_active());
1617
1618        es.exit_endgame();
1619        assert!(!es.is_endgame_active());
1620    }
1621
1622    #[test]
1623    fn test_endgame_track_request() {
1624        let mut es = EndgameState::new();
1625        es.enter_endgame();
1626
1627        // Track requests from 3 peers for the same block
1628        es.track_request(0, 0, 16384, 0);
1629        es.track_request(0, 0, 16384, 1);
1630        es.track_request(0, 0, 16384, 2);
1631
1632        assert_eq!(es.tracked_count(), 1); // One unique block tracked
1633
1634        let targets = es.get_cancel_targets(0, 0, 16384);
1635        assert_eq!(targets.len(), 3);
1636        assert!(targets.contains(&0));
1637        assert!(targets.contains(&1));
1638        assert!(targets.contains(&2));
1639    }
1640
1641    #[test]
1642    fn test_endgame_cancel_removes_on_arrival() {
1643        let mut es = EndgameState::new();
1644        es.enter_endgame();
1645
1646        es.track_request(5, 0, 16384, 0);
1647        es.track_request(5, 0, 16384, 1);
1648
1649        let targets = es.get_cancel_targets(5, 0, 16384);
1650        assert_eq!(targets.len(), 2);
1651
1652        // After removal, no more targets
1653        es.remove_request(5, 0, 16384);
1654        let targets_after = es.get_cancel_targets(5, 0, 16384);
1655        assert!(targets_after.is_empty());
1656        assert_eq!(es.tracked_count(), 0);
1657    }
1658
1659    #[test]
1660    fn test_endgame_multiple_blocks_tracked_independently() {
1661        let mut es = EndgameState::new();
1662        es.enter_endgame();
1663
1664        // Track different blocks
1665        es.track_request(0, 0, 16384, 0);
1666        es.track_request(0, 0, 16384, 1);
1667        es.track_request(0, 16384, 16384, 0);
1668        es.track_request(0, 16384, 16384, 2);
1669
1670        assert_eq!(es.tracked_count(), 2);
1671
1672        // Cancel one block doesn't affect the other
1673        es.remove_request(0, 0, 16384);
1674        assert_eq!(es.tracked_count(), 1);
1675
1676        let remaining = es.get_cancel_targets(0, 16384, 16384);
1677        assert_eq!(remaining.len(), 2);
1678        assert!(remaining.contains(&0));
1679        assert!(remaining.contains(&2));
1680    }
1681
1682    #[test]
1683    fn test_endgame_exit_clears_all_tracking() {
1684        let mut es = EndgameState::new();
1685        es.enter_endgame();
1686
1687        es.track_request(10, 0, 16384, 0);
1688        es.track_request(10, 0, 16384, 1);
1689        es.track_request(11, 0, 8192, 0);
1690        assert_eq!(es.tracked_count(), 2);
1691
1692        es.exit_endgame();
1693        assert!(!es.is_endgame_active());
1694        assert_eq!(es.tracked_count(), 0);
1695    }
1696
1697    #[test]
1698    fn test_endgame_get_cancel_targets_empty_when_inactive() {
1699        let es = EndgameState::new();
1700        // Even if we somehow track (shouldn't happen when inactive), targets should be empty
1701        // Actually tracking works regardless, but is_endgate_active gates usage
1702        let targets = es.get_cancel_targets(99, 0, 16384);
1703        assert!(targets.is_empty());
1704    }
1705
1706    #[test]
1707    fn test_endgame_track_different_piece_offsets_lengths() {
1708        let mut es = EndgameState::new();
1709        es.enter_endgame();
1710
1711        // Last block might be shorter
1712        es.track_request(0, 32768, 8000, 0);
1713        es.track_request(0, 32768, 8000, 1);
1714
1715        let targets = es.get_cancel_targets(0, 32768, 8000);
1716        assert_eq!(targets.len(), 2);
1717    }
1718
1719    #[test]
1720    fn test_endgame_remove_nonexistent_is_noop() {
1721        let mut es = EndgameState::new();
1722        es.enter_endgame();
1723
1724        // Remove something that was never tracked - should not panic
1725        es.remove_request(999, 999, 999);
1726        assert_eq!(es.tracked_count(), 0);
1727    }
1728
1729    // ==================== BEP 6 Fast Extension Tests ====================
1730
1731    #[test]
1732    fn test_is_bitfield_set_basic() {
1733        // Test bitfield: [0b11000000] = pieces 0 and 1 set (MSB first)
1734        let bf = vec![0xC0];
1735        assert!(BtDownloadCommand::is_bitfield_set(&bf, 0));
1736        assert!(BtDownloadCommand::is_bitfield_set(&bf, 1));
1737        assert!(!BtDownloadCommand::is_bitfield_set(&bf, 2));
1738        assert!(!BtDownloadCommand::is_bitfield_set(&bf, 7));
1739    }
1740
1741    #[test]
1742    fn test_is_bitfield_set_multi_byte() {
1743        // Bitfield for 16 pieces: all set
1744        let bf = vec![0xFF, 0xFF];
1745        for i in 0..16u32 {
1746            assert!(
1747                BtDownloadCommand::is_bitfield_set(&bf, i),
1748                "Piece {} should be set",
1749                i
1750            );
1751        }
1752    }
1753
1754    #[test]
1755    fn test_is_bitfield_set_out_of_range() {
1756        let bf = vec![0xFF];
1757        assert!(!BtDownloadCommand::is_bitfield_set(&bf, 8)); // Beyond bitfield length
1758        assert!(!BtDownloadCommand::is_bitfield_set(&bf, 100));
1759    }
1760
1761    #[test]
1762    fn test_calculate_fast_set_basic() {
1763        let needed = vec![0u32, 1, 2, 3, 4, 5];
1764        let peer_bf = vec![0b11111100]; // Peer has pieces 0-5
1765
1766        let already_sent = HashSet::new();
1767        let fast_set = BtDownloadCommand::calculate_fast_set(&needed, &peer_bf, &already_sent);
1768
1769        assert_eq!(fast_set.len(), 6); // All pieces should be selected (<10 limit)
1770        assert!(fast_set.contains(&0));
1771        assert!(fast_set.contains(&5));
1772    }
1773
1774    #[test]
1775    fn test_calculate_fast_set_respects_max_limit() {
1776        // Create 15 needed pieces
1777        let needed: Vec<u32> = (0..15).collect();
1778        let peer_bf = vec![0xFF, 0xFF]; // Peer has first 16 pieces
1779
1780        let already_sent = HashSet::new();
1781        let fast_set = BtDownloadCommand::calculate_fast_set(&needed, &peer_bf, &already_sent);
1782
1783        assert_eq!(fast_set.len(), 10); // Should cap at MAX_ALLOWED_FAST_PER_PEER
1784    }
1785
1786    #[test]
1787    fn test_calculate_fast_set_excludes_already_sent() {
1788        let needed = vec![0u32, 1, 2, 3, 4];
1789        let peer_bf = vec![0b11111000];
1790
1791        let mut already_sent = HashSet::new();
1792        already_sent.insert(0);
1793        already_sent.insert(1);
1794
1795        let fast_set = BtDownloadCommand::calculate_fast_set(&needed, &peer_bf, &already_sent);
1796
1797        assert_eq!(fast_set.len(), 3); // Only 2,3,4 should be selected
1798        assert!(!fast_set.contains(&0));
1799        assert!(!fast_set.contains(&1));
1800        assert!(fast_set.contains(&2));
1801    }
1802
1803    #[test]
1804    fn test_calculate_fast_set_filters_by_peer_bitfield() {
1805        let needed = vec![0u32, 1, 2, 3, 4];
1806        let peer_bf = vec![0b00011000]; // bitfield byte: bits 3 and 4 set (pieces 3,4)
1807
1808        let already_sent = HashSet::new();
1809        let fast_set = BtDownloadCommand::calculate_fast_set(&needed, &peer_bf, &already_sent);
1810
1811        assert_eq!(fast_set.len(), 2); // Only pieces that peer has
1812        assert!(fast_set.contains(&3));
1813        assert!(fast_set.contains(&4));
1814        assert!(!fast_set.contains(&0));
1815        assert!(!fast_set.contains(&1));
1816        assert!(!fast_set.contains(&2));
1817    }
1818}