Skip to main content

dhaar_torrent/
lib.rs

1pub mod config;
2pub mod error;
3pub mod helpers;
4pub mod peer_connection;
5pub mod peer_explorer;
6pub mod peer_manager;
7pub mod piece_manager;
8pub mod status;
9pub mod store;
10pub mod torrent_parser;
11pub mod wire_protocol;
12
13use std::{path::Path, sync::Arc, time::Duration};
14
15use tokio::{sync::watch, task::JoinSet, time};
16
17use crate::{
18    error::Result,
19    helpers::generate_random_peer_id,
20    peer_explorer::{
21        PeerExplorer, PeerSource, channel::new_peer_explorer_channel, tracker::TrackerManager,
22    },
23    peer_manager::{
24        PeerManager,
25        peer_selection_strategy::{PeerSelectionStrategy, RetryAfterDelayPeerSelectionStrategy},
26    },
27    piece_manager::{PieceManager, channel::new_piece_manager_channel},
28    status::{DownloadState, DownloadStats, DownloadStatus, PieceProgress},
29    store::{DiskStore, Store},
30    torrent_parser::{TorrentParser, metadata::Torrent, parser::TorrentFileParser},
31};
32
33/// How often the sampled [`DownloadStatus`] is republished. Rates are measured
34/// across this window, so it is also how quickly they answer to a change.
35const STATUS_INTERVAL: Duration = Duration::from_secs(1);
36
37/// One torrent, and everything needed to fetch it.
38///
39/// The pieces of this crate are separate actors that only talk over channels,
40/// which makes them easy to test in isolation and tedious to assemble by hand.
41/// `Download` owns that assembly: the channels between the actors are an
42/// internal detail, and callers hold a value rather than a set of tasks.
43///
44/// The parts a caller might reasonably want to swap are type parameters —
45/// where bytes land (`W`), and which peer to try next (`S`) — while peer
46/// sources are trait objects because a download usually draws on several at
47/// once. [`Download::from_torrent_file`] fills all three in with the defaults.
48pub struct Download<W, S>
49where
50    W: Store + Send + Sync + 'static,
51    W::Error: std::error::Error + Send + Sync + 'static,
52    S: PeerSelectionStrategy + Send + Sync + 'static,
53{
54    torrent: Torrent,
55    peer_id: [u8; 20],
56    listening_port: u16,
57    store: W,
58    peer_selection_strategy: S,
59    peer_sources: Vec<Box<dyn PeerSource + Send>>,
60    stats: Arc<DownloadStats>,
61    progress_sender: watch::Sender<PieceProgress>,
62    status_sender: watch::Sender<DownloadStatus>,
63}
64
65impl Download<DiskStore, RetryAfterDelayPeerSelectionStrategy> {
66    /// Reads a `.torrent` file and takes the usual defaults: written to disk
67    /// beside the current directory, peers from the trackers the file names,
68    /// failed peers retried after a delay.
69    pub fn from_torrent_file(path: &Path, listening_port: u16) -> Result<Self> {
70        Ok(Self::from_torrent_with_port(
71            TorrentFileParser::parse_from_file_path(path)?,
72            listening_port,
73        ))
74    }
75
76    pub fn from_torrent_file_with_port(path: &Path, listening_port: u16) -> Result<Self> {
77        Ok(Self::from_torrent_with_port(
78            TorrentFileParser::parse_from_file_path(path)?,
79            listening_port,
80        ))
81    }
82
83    pub fn from_torrent_with_port(torrent: Torrent, listening_port: u16) -> Self {
84        let peer_id = generate_random_peer_id();
85        // Built first: the tracker announces these figures, so it needs the
86        // same counters the connections will be writing to.
87        let stats = Arc::new(DownloadStats::default());
88        let tracker_manager = TrackerManager::new(
89            torrent.announce_urls(),
90            &torrent.info_hash,
91            &peer_id,
92            stats.clone(),
93            listening_port,
94        );
95        let store = DiskStore::new(
96            torrent.info.total_length(),
97            torrent.info.piece_length,
98            &torrent.info.name,
99            &torrent.info.md5sum,
100            &torrent.info.files,
101            torrent.info_hash,
102        );
103
104        Self::new(
105            torrent,
106            peer_id,
107            store,
108            RetryAfterDelayPeerSelectionStrategy::new(),
109            vec![Box::new(tracker_manager)],
110            stats,
111            listening_port,
112        )
113    }
114}
115
116impl<W, S> Download<W, S>
117where
118    W: Store + Send + Sync + 'static,
119    W::Error: std::error::Error + Send + Sync + 'static,
120    S: PeerSelectionStrategy + Send + Sync + 'static,
121{
122    /// Every part chosen explicitly. `peer_id` identifies us to peers for the
123    /// life of the download and must stay fixed once announced to a tracker.
124    ///
125    /// `stats` is taken rather than created because peer sources are built by
126    /// the caller and may need to read it — a tracker announce is meaningless
127    /// without the byte counts.
128    pub fn new(
129        torrent: Torrent,
130        peer_id: [u8; 20],
131        store: W,
132        peer_selection_strategy: S,
133        peer_sources: Vec<Box<dyn PeerSource + Send>>,
134        stats: Arc<DownloadStats>,
135        listening_port: u16,
136    ) -> Self {
137        Self {
138            torrent,
139            peer_id,
140            listening_port,
141            store,
142            peer_selection_strategy,
143            peer_sources,
144            stats,
145            progress_sender: watch::Sender::new(PieceProgress::default()),
146            status_sender: watch::Sender::new(DownloadStatus::default()),
147        }
148    }
149
150    pub fn torrent(&self) -> &Torrent {
151        &self.torrent
152    }
153
154    pub fn peer_id(&self) -> &[u8; 20] {
155        &self.peer_id
156    }
157
158    /// The live counters. Readable at any time and from anywhere, including
159    /// before the download starts.
160    pub fn stats(&self) -> Arc<DownloadStats> {
161        self.stats.clone()
162    }
163
164    /// A feed of sampled status. Subscribing before [`Download::spawn`] is
165    /// fine — the receiver holds the empty starting value until the first
166    /// sample lands.
167    pub fn subscribe(&self) -> watch::Receiver<DownloadStatus> {
168        self.status_sender.subscribe()
169    }
170
171    /// Starts every actor and returns at once, handing back the means to
172    /// watch and to stop them.
173    pub fn spawn(self) -> DownloadHandle {
174        let (peer_explorer_channel_sender, peer_explorer_channel_receiver) =
175            new_peer_explorer_channel();
176        let (piece_manager_channel_sender, piece_manager_channel_receiver) =
177            new_piece_manager_channel();
178
179        let peer_explorer = PeerExplorer::new(self.peer_sources);
180        // One store, shared. The piece manager lays it out and owns the
181        // bitfield; every connection writes the pieces it finishes.
182        let store = Arc::new(self.store);
183        let piece_manager = PieceManager::new(
184            &self.torrent.info.pieces,
185            self.torrent.info.piece_length,
186            self.torrent.info.total_length(),
187            store.clone(),
188            self.stats.clone(),
189            self.progress_sender.clone(),
190            self.torrent.info_hash,
191        );
192        let peer_manager = PeerManager::new(
193            self.peer_selection_strategy,
194            &self.torrent.info_hash,
195            &self.peer_id,
196            self.stats.clone(),
197            self.listening_port,
198            store,
199        );
200
201        let status = self.status_sender.subscribe();
202        let progress = self.progress_sender.subscribe();
203
204        let mut tasks: JoinSet<()> = JoinSet::new();
205        tasks.spawn(peer_explorer.start(peer_explorer_channel_sender));
206        tasks.spawn(piece_manager.start(piece_manager_channel_receiver));
207        tasks.spawn(
208            peer_manager.start(peer_explorer_channel_receiver, piece_manager_channel_sender),
209        );
210        tasks.spawn(sample_status(
211            self.stats.clone(),
212            progress,
213            self.status_sender,
214        ));
215
216        DownloadHandle {
217            stats: self.stats,
218            status,
219            tasks,
220        }
221    }
222
223    pub fn set_listening_port(&mut self, port: u16) {
224        self.listening_port = port;
225    }
226
227    /// Runs until the first actor stops. Any one of them stopping means the
228    /// download cannot continue — they are only useful in each other's
229    /// company — so the rest are dropped rather than left running.
230    pub async fn start(self) {
231        self.spawn().wait().await;
232    }
233}
234
235/// A running download.
236///
237/// Dropping this aborts every actor, so it has to be held for as long as the
238/// download should live.
239pub struct DownloadHandle {
240    stats: Arc<DownloadStats>,
241    status: watch::Receiver<DownloadStatus>,
242    tasks: JoinSet<()>,
243}
244
245impl DownloadHandle {
246    /// The live counters, updated as the work happens rather than on the
247    /// sampling clock.
248    pub fn stats(&self) -> &Arc<DownloadStats> {
249        &self.stats
250    }
251
252    /// The most recent sample. Cheap, and never waits.
253    pub fn status(&self) -> DownloadStatus {
254        self.status.borrow().clone()
255    }
256
257    pub fn subscribe(&self) -> watch::Receiver<DownloadStatus> {
258        self.status.clone()
259    }
260
261    /// Waits for the next sample. `false` once the download has stopped and
262    /// no further sample will arrive.
263    pub async fn changed(&mut self) -> bool {
264        self.status.changed().await.is_ok()
265    }
266
267    /// Waits until the first actor stops, then drops the rest.
268    pub async fn wait(mut self) {
269        self.tasks.join_next().await;
270    }
271
272    /// Stops every actor now.
273    pub fn shutdown(mut self) {
274        self.tasks.abort_all();
275    }
276}
277
278/// Publishes a [`DownloadStatus`] every [`STATUS_INTERVAL`].
279///
280/// Rates are the reason this exists: they cannot be read from a counter at one
281/// instant, only measured between two of them.
282async fn sample_status(
283    stats: Arc<DownloadStats>,
284    mut progress: watch::Receiver<PieceProgress>,
285    status: watch::Sender<DownloadStatus>,
286) {
287    let mut ticker = time::interval(STATUS_INTERVAL);
288    let mut previous_downloaded = 0;
289    let mut previous_uploaded = 0;
290    let seconds = STATUS_INTERVAL.as_secs().max(1);
291
292    loop {
293        ticker.tick().await;
294
295        let downloaded_bytes = stats.downloaded_bytes();
296        let uploaded_bytes = stats.uploaded_bytes();
297        let pieces = progress.borrow_and_update().clone();
298
299        // Extraction is checked before completeness, not after: every piece is
300        // verified for the whole of `Finalizing` too, so testing `is_complete`
301        // first would report `Seeding` throughout and the new state would be
302        // unreachable.
303        let state = if pieces.extracted {
304            DownloadState::Seeding
305        } else if stats.is_complete() {
306            DownloadState::Finalizing
307        } else if pieces.completed_pieces == 0 {
308            DownloadState::Starting
309        } else {
310            DownloadState::Downloading
311        };
312
313        status.send_replace(DownloadStatus {
314            state,
315            pieces,
316            downloaded_bytes,
317            uploaded_bytes,
318            wasted_bytes: stats.wasted_bytes(),
319            hash_failures: stats.hash_failures(),
320            in_flight_pieces: stats.in_flight_pieces(),
321            active_peers: stats.active_peers(),
322            download_rate: downloaded_bytes.saturating_sub(previous_downloaded) / seconds,
323            upload_rate: uploaded_bytes.saturating_sub(previous_uploaded) / seconds,
324        });
325
326        previous_downloaded = downloaded_bytes;
327        previous_uploaded = uploaded_bytes;
328    }
329}