Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion .github/workflows/checks.yml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ on:

env:
CARGO_TERM_COLOR: always
RUSTFLAGS: "-C target-cpu=native"
CARGO_INCREMENTAL: 0
RUST_BACKTRACE: 1

Expand Down
32 changes: 32 additions & 0 deletions crates/libtortillas/src/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,26 @@ pub enum TorrentError {
reason: String,
},

/// Tried to connect a peer that is already part of this torrent's swarm.
#[error("Peer is already connected: {peer_id}")]
PeerAlreadyConnected { peer_id: PeerId },

/// Tried to disconnect a peer that is not part of this torrent's swarm.
#[error("Peer not found: {peer_id}")]
PeerNotFound { peer_id: PeerId },

/// Tried to add a tracker that is already configured for this torrent.
#[error("Tracker already exists: {endpoint}")]
TrackerAlreadyExists { endpoint: String },

/// Tried to operate on a tracker that is not configured for this torrent.
#[error("Tracker not found: {endpoint}")]
TrackerNotFound { endpoint: String },

/// Tried to add a tracker protocol that the runtime cannot operate.
#[error("Unsupported tracker protocol: {protocol}")]
UnsupportedTrackerProtocol { protocol: &'static str },

/// Bitfield operation failed
#[error("Bitfield operation failed: {reason}")]
BitfieldError { reason: String },
Expand Down Expand Up @@ -385,6 +405,18 @@ pub(crate) fn map_torrent_send_error<M>(
},
}
}

pub(crate) fn map_torrent_communication_error<M, E>(
operation: &'static str, error: SendError<M, E>,
) -> TorrentError
where
E: std::fmt::Debug + std::fmt::Display,
{
TorrentError::ActorCommunicationFailed {
operation,
reason: error.to_string(),
}
}
// Conversion implementations for backward compatibility during transition
impl From<num_enum::TryFromPrimitiveError<super::tracker::udp::Action>> for TrackerActorError {
fn from(err: num_enum::TryFromPrimitiveError<super::tracker::udp::Action>) -> Self {
Expand Down
1 change: 1 addition & 0 deletions crates/libtortillas/src/facade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
pub use crate::{
engine::{Engine, EngineSnapshot, EngineStatus, TorrentSource},
torrent::{RestoreVerification, Torrent, TorrentSnapshot},
tracker::Tracker,
};
#[cfg(feature = "live")]
pub use crate::{
Expand Down
4 changes: 3 additions & 1 deletion crates/libtortillas/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -571,7 +571,9 @@ pub(crate) mod testing {
stream.write_all(&message.to_bytes()?).await?;
}

Ok(())
// Keep the connection open so the peer does not observe EOF until its
// accept task and JoinSet are aborted by LocalPeer::drop.
std::future::pending().await
}

pub(crate) fn test_info_hash() -> InfoHash {
Expand Down
31 changes: 26 additions & 5 deletions crates/libtortillas/src/live/hub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -76,11 +76,7 @@ where
}

fn remove_value(&self, value: &Arc<V>) -> bool {
let key = self
.values
.iter()
.find(|entry| Arc::ptr_eq(entry.value(), value))
.map(|entry| entry.key().clone());
let key = self.key_for_value(value);
key.is_some_and(|key| {
self
.values
Expand All @@ -89,6 +85,14 @@ where
})
}

fn key_for_value(&self, value: &Arc<V>) -> Option<K> {
self
.values
.iter()
.find(|entry| Arc::ptr_eq(entry.value(), value))
.map(|entry| entry.key().clone())
}

fn values(&self) -> Vec<Arc<V>> {
self
.values
Expand Down Expand Up @@ -607,6 +611,23 @@ impl Hub {
})
}

pub(crate) fn tracker_handle(
&self, torrent: InfoHash, source: &Tracker,
) -> Option<TrackerHandle> {
let inner = self.inner()?;
let scope = inner.torrents.get(&torrent)?;
scope
.trackers
.get(source)
.map(|inner| TrackerHandle { inner })
}

pub(crate) fn tracker_source(&self, tracker: &TrackerHandle) -> Option<Tracker> {
let inner = self.inner()?;
let scope = inner.torrents.get(&tracker.torrent())?;
scope.trackers.key_for_value(&tracker.inner)
}

pub(crate) fn register_tracker_scope(
&self, torrent: InfoHash, source: &Tracker, view: TrackerView,
) -> Option<TrackerHandle> {
Expand Down
105 changes: 35 additions & 70 deletions crates/libtortillas/src/torrent/actor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,17 +35,12 @@ use crate::{
BLOCK_SIZE, PieceBlockSnapshot, PieceStorageStrategy, TORRENT_SNAPSHOT_VERSION,
TorrentSnapshot, TorrentState,
},
tracker::{
Announce, Event, Tracker, TrackerActor, TrackerActorArgs, TrackerUpdate, udp::UdpServer,
},
tracker::{Announce, Event, Tracker, TrackerActor, TrackerUpdate, udp::UdpServer},
};
#[cfg(feature = "live")]
use crate::{
live::{Hub, LiveHealthLevel, TorrentEventKind, TorrentView, TrackerStatus, TrackerView},
metrics::{
ByteCount, ContentProgress, HasTransferMetrics, TorrentMetrics, TrackerMetrics,
TransferMetrics,
},
live::{Hub, LiveHealthLevel, TorrentEventKind, TorrentView},
metrics::{ByteCount, ContentProgress, HasTransferMetrics, TorrentMetrics, TransferMetrics},
};

/// A hook that is called when the torrent is ready to start downloading.
Expand Down Expand Up @@ -117,6 +112,7 @@ pub(crate) struct TorrentActor {
pub(super) metainfo: MetaInfo,
#[allow(dead_code)]
pub(super) tracker_server: UdpServer,
pub(super) primary_addr: SocketAddr,
pub(super) scheduler: ActorRef<Scheduler>,
/// Should only be used to create new connections
pub(super) utp_server: Arc<UtpSocketUdp>,
Expand Down Expand Up @@ -477,10 +473,13 @@ impl TorrentActor {
);
}

pub(super) fn remove_peer(&mut self, id: PeerId) {
pub(super) fn remove_peer(&mut self, id: PeerId) -> bool {
self.piece_scheduler.peer_disconnected(id);
if let Some(actor) = self.peers.remove(&id) {
actor.kill();
true
} else {
false
}
}

Expand Down Expand Up @@ -757,61 +756,7 @@ impl Actor for TorrentActor {
TorrentState::ResolvingMetadata
};
let bitfield = BitVec::repeat(false, piece_count);
let initial_left = info.map(Info::total_length);

// Create tracker actors
let tracker_list = metainfo.announce_list();
let mut trackers = HashMap::new();
for tracker in tracker_list {
if matches!(tracker, Tracker::Websocket(_)) {
warn!(
tracker_uri = %tracker.redacted_endpoint(),
"Skipping unsupported websocket tracker"
);
continue;
}
#[cfg(feature = "live")]
let endpoint = tracker.redacted_endpoint();
#[cfg(feature = "live")]
let Some(tracker_handle) = hub.register_tracker_scope(
torrent_id,
&tracker,
TrackerView {
endpoint,
status: TrackerStatus::Pending,
metrics: TrackerMetrics::default(),
},
) else {
return Err(TorrentError::ActorCommunicationFailed {
operation: "register tracker live scope",
reason: "live hub is no longer available".to_string(),
});
};
let actor = TrackerActor::supervise(
&us,
TrackerActorArgs {
tracker: tracker.clone(),
peer_id,
server: tracker_server.clone(),
socket_addr: primary_addr,
initial_left,
supervisor: us.clone(),
scheduler: scheduler.clone(),
settings: settings.tracker.clone(),
#[cfg(feature = "live")]
live_handle: tracker_handle,
},
)
.restart_policy(RestartPolicy::Transient)
.restart_limit(
settings.torrent.tracker_restart.limit,
settings.torrent.tracker_restart.period,
)
.spawn()
.await;

trackers.insert(tracker, actor);
}
let default_manager = FilePieceManager(base_path, info.cloned());
let piece_store = PieceStoreActor::supervise(&us, ())
.restart_policy(RestartPolicy::Permanent)
Expand All @@ -822,15 +767,16 @@ impl Actor for TorrentActor {
.spawn()
.await;

let actor = Self {
let mut actor = Self {
#[cfg(feature = "live")]
hub,
peers: HashMap::new(),
bitfield,
tracker_server,
primary_addr,
scheduler,
utp_server,
trackers,
trackers: HashMap::new(),
id: peer_id,
info_hash: torrent_id,
metainfo,
Expand All @@ -853,6 +799,15 @@ impl Actor for TorrentActor {
piece_manager: PieceManagerProxy::Default(default_manager),
settings,
};
for tracker in tracker_list {
if let Err(error) = actor.register_tracker_actor(tracker.clone()).await {
warn!(
%error,
tracker_uri = %tracker.redacted_endpoint(),
"Skipping unavailable tracker"
);
}
}
crate::live_only!(actor.hub.initialize_torrent_projection(actor.live_view()));

Ok(actor)
Expand Down Expand Up @@ -938,9 +893,11 @@ mod tests {
use super::*;
use crate::{
hashes::HashVec,
live::{PeerIdentity, PeerView},
live::{PeerIdentity, PeerView, TrackerStatus, TrackerView},
metainfo::{InfoKeys, MetaInfo, TorrentFile},
metrics::{BytesPerSecond, PeerMetrics, TrafficTotals, TransferRates, TransferSample},
metrics::{
BytesPerSecond, PeerMetrics, TrackerMetrics, TrafficTotals, TransferRates, TransferSample,
},
protocol::{
messages::{Handshake, PeerMessages},
stream::{PeerRecv, PeerSend, PeerStream},
Expand All @@ -949,8 +906,8 @@ mod tests {
testing,
torrent::{
BLOCK_SIZE, Torrent, TorrentSnapshot,
commands::{GetState, HasInfoDict, SetState, SnapshotState},
events::{AddPeer, IncomingPiece},
commands::{AddPeer, GetState, HasInfoDict, SetState, SnapshotState},
events::IncomingPiece,
},
tracker::Tracker,
};
Expand Down Expand Up @@ -1346,7 +1303,12 @@ mod tests {
hub: Hub::default(),
});
let torrent = Torrent::new(info_hash, actor.clone());
actor.tell(AddPeer { peer: seed.peer() }).await.unwrap();
let connection = actor.ask(AddPeer { peer: seed.peer() }).await.unwrap();
timeout(Duration::from_secs(5), connection)
.await
.expect("local seed connection should complete")
.unwrap()
.unwrap();

timeout(Duration::from_secs(5), torrent.poll_ready())
.await
Expand Down Expand Up @@ -1701,6 +1663,7 @@ mod tests {
resolved_magnet_info: None,
metainfo: metainfo.clone(),
tracker_server: udp_server.clone(),
primary_addr: utp_server.bind_addr(),
scheduler: Scheduler::spawn(Scheduler::new()),
utp_server,
actor_ref: actor_ref.clone(),
Expand Down Expand Up @@ -1861,6 +1824,7 @@ mod tests {
resolved_magnet_info: None,
metainfo: metainfo.clone(),
tracker_server: udp_server,
primary_addr: utp_server.bind_addr(),
scheduler: Scheduler::spawn(Scheduler::new()),
utp_server,
actor_ref: actor_ref.clone(),
Expand Down Expand Up @@ -2083,6 +2047,7 @@ mod tests {
resolved_magnet_info: None,
metainfo,
tracker_server: udp_server,
primary_addr: utp_server.bind_addr(),
scheduler: Scheduler::spawn(Scheduler::new()),
utp_server,
actor_ref: actor_ref.clone(),
Expand Down
Loading
Loading