pub mod config;
pub mod error;
pub mod helpers;
pub mod peer_connection;
pub mod peer_explorer;
pub mod peer_manager;
pub mod piece_manager;
pub mod status;
pub mod store;
pub mod torrent_parser;
pub mod wire_protocol;
use std::{path::Path, sync::Arc, time::Duration};
use tokio::{sync::watch, task::JoinSet, time};
use crate::{
error::Result,
helpers::generate_random_peer_id,
peer_explorer::{
PeerExplorer, PeerSource, channel::new_peer_explorer_channel, tracker::TrackerManager,
},
peer_manager::{
PeerManager,
peer_selection_strategy::{PeerSelectionStrategy, RetryAfterDelayPeerSelectionStrategy},
},
piece_manager::{PieceManager, channel::new_piece_manager_channel},
status::{DownloadState, DownloadStats, DownloadStatus, PieceProgress},
store::{DiskStore, Store},
torrent_parser::{TorrentParser, metadata::Torrent, parser::TorrentFileParser},
};
const STATUS_INTERVAL: Duration = Duration::from_secs(1);
pub struct Download<W, S>
where
W: Store + Send + Sync + 'static,
W::Error: std::error::Error + Send + Sync + 'static,
S: PeerSelectionStrategy + Send + Sync + 'static,
{
torrent: Torrent,
peer_id: [u8; 20],
listening_port: u16,
store: W,
peer_selection_strategy: S,
peer_sources: Vec<Box<dyn PeerSource + Send>>,
stats: Arc<DownloadStats>,
progress_sender: watch::Sender<PieceProgress>,
status_sender: watch::Sender<DownloadStatus>,
}
impl Download<DiskStore, RetryAfterDelayPeerSelectionStrategy> {
pub fn from_torrent_file(path: &Path, listening_port: u16) -> Result<Self> {
Ok(Self::from_torrent_with_port(
TorrentFileParser::parse_from_file_path(path)?,
listening_port,
))
}
pub fn from_torrent_file_with_port(path: &Path, listening_port: u16) -> Result<Self> {
Ok(Self::from_torrent_with_port(
TorrentFileParser::parse_from_file_path(path)?,
listening_port,
))
}
pub fn from_torrent_with_port(torrent: Torrent, listening_port: u16) -> Self {
let peer_id = generate_random_peer_id();
let stats = Arc::new(DownloadStats::default());
let tracker_manager = TrackerManager::new(
torrent.announce_urls(),
&torrent.info_hash,
&peer_id,
stats.clone(),
listening_port,
);
let store = DiskStore::new(
torrent.info.total_length(),
torrent.info.piece_length,
&torrent.info.name,
&torrent.info.md5sum,
&torrent.info.files,
torrent.info_hash,
);
Self::new(
torrent,
peer_id,
store,
RetryAfterDelayPeerSelectionStrategy::new(),
vec![Box::new(tracker_manager)],
stats,
listening_port,
)
}
}
impl<W, S> Download<W, S>
where
W: Store + Send + Sync + 'static,
W::Error: std::error::Error + Send + Sync + 'static,
S: PeerSelectionStrategy + Send + Sync + 'static,
{
pub fn new(
torrent: Torrent,
peer_id: [u8; 20],
store: W,
peer_selection_strategy: S,
peer_sources: Vec<Box<dyn PeerSource + Send>>,
stats: Arc<DownloadStats>,
listening_port: u16,
) -> Self {
Self {
torrent,
peer_id,
listening_port,
store,
peer_selection_strategy,
peer_sources,
stats,
progress_sender: watch::Sender::new(PieceProgress::default()),
status_sender: watch::Sender::new(DownloadStatus::default()),
}
}
pub fn torrent(&self) -> &Torrent {
&self.torrent
}
pub fn peer_id(&self) -> &[u8; 20] {
&self.peer_id
}
pub fn stats(&self) -> Arc<DownloadStats> {
self.stats.clone()
}
pub fn subscribe(&self) -> watch::Receiver<DownloadStatus> {
self.status_sender.subscribe()
}
pub fn spawn(self) -> DownloadHandle {
let (peer_explorer_channel_sender, peer_explorer_channel_receiver) =
new_peer_explorer_channel();
let (piece_manager_channel_sender, piece_manager_channel_receiver) =
new_piece_manager_channel();
let peer_explorer = PeerExplorer::new(self.peer_sources);
let store = Arc::new(self.store);
let piece_manager = PieceManager::new(
&self.torrent.info.pieces,
self.torrent.info.piece_length,
self.torrent.info.total_length(),
store.clone(),
self.stats.clone(),
self.progress_sender.clone(),
self.torrent.info_hash,
);
let peer_manager = PeerManager::new(
self.peer_selection_strategy,
&self.torrent.info_hash,
&self.peer_id,
self.stats.clone(),
self.listening_port,
store,
);
let status = self.status_sender.subscribe();
let progress = self.progress_sender.subscribe();
let mut tasks: JoinSet<()> = JoinSet::new();
tasks.spawn(peer_explorer.start(peer_explorer_channel_sender));
tasks.spawn(piece_manager.start(piece_manager_channel_receiver));
tasks.spawn(
peer_manager.start(peer_explorer_channel_receiver, piece_manager_channel_sender),
);
tasks.spawn(sample_status(
self.stats.clone(),
progress,
self.status_sender,
));
DownloadHandle {
stats: self.stats,
status,
tasks,
}
}
pub fn set_listening_port(&mut self, port: u16) {
self.listening_port = port;
}
pub async fn start(self) {
self.spawn().wait().await;
}
}
pub struct DownloadHandle {
stats: Arc<DownloadStats>,
status: watch::Receiver<DownloadStatus>,
tasks: JoinSet<()>,
}
impl DownloadHandle {
pub fn stats(&self) -> &Arc<DownloadStats> {
&self.stats
}
pub fn status(&self) -> DownloadStatus {
self.status.borrow().clone()
}
pub fn subscribe(&self) -> watch::Receiver<DownloadStatus> {
self.status.clone()
}
pub async fn changed(&mut self) -> bool {
self.status.changed().await.is_ok()
}
pub async fn wait(mut self) {
self.tasks.join_next().await;
}
pub fn shutdown(mut self) {
self.tasks.abort_all();
}
}
async fn sample_status(
stats: Arc<DownloadStats>,
mut progress: watch::Receiver<PieceProgress>,
status: watch::Sender<DownloadStatus>,
) {
let mut ticker = time::interval(STATUS_INTERVAL);
let mut previous_downloaded = 0;
let mut previous_uploaded = 0;
let seconds = STATUS_INTERVAL.as_secs().max(1);
loop {
ticker.tick().await;
let downloaded_bytes = stats.downloaded_bytes();
let uploaded_bytes = stats.uploaded_bytes();
let pieces = progress.borrow_and_update().clone();
let state = if pieces.extracted {
DownloadState::Seeding
} else if stats.is_complete() {
DownloadState::Finalizing
} else if pieces.completed_pieces == 0 {
DownloadState::Starting
} else {
DownloadState::Downloading
};
status.send_replace(DownloadStatus {
state,
pieces,
downloaded_bytes,
uploaded_bytes,
wasted_bytes: stats.wasted_bytes(),
hash_failures: stats.hash_failures(),
in_flight_pieces: stats.in_flight_pieces(),
active_peers: stats.active_peers(),
download_rate: downloaded_bytes.saturating_sub(previous_downloaded) / seconds,
upload_rate: uploaded_bytes.saturating_sub(previous_uploaded) / seconds,
});
previous_downloaded = downloaded_bytes;
previous_uploaded = uploaded_bytes;
}
}