From 0f7207381c5c0a16387990329f962ff81392b2f3 Mon Sep 17 00:00:00 2001 From: Frigyes Erdosi Szucs Date: Fri, 31 Jul 2026 10:00:54 -0700 Subject: [PATCH 1/2] I2P support --- Cargo.lock | 3 + .../src/v1/extractors/announce_request.rs | 4 +- .../src/v1/handlers/announce.rs | 19 +- .../receiving_an_announce_request.rs | 90 +++++-- packages/http-core/benches/helpers/util.rs | 4 +- packages/http-core/src/lib.rs | 5 +- packages/http-core/src/services/announce.rs | 70 ++++- packages/http-core/src/services/scrape.rs | 6 +- packages/http-protocol/Cargo.toml | 1 + packages/http-protocol/src/v1/query.rs | 53 ++-- .../http-protocol/src/v1/requests/announce.rs | 94 ++++++- .../src/v1/responses/announce/data.rs | 14 +- .../src/v1/responses/announce/encoding.rs | 175 +++++++------ .../src/v1/responses/announce/mod.rs | 2 +- packages/primitives/Cargo.toml | 2 + packages/primitives/src/lib.rs | 1 + packages/primitives/src/peer.rs | 245 ++++++++++++++++-- .../src/v1/conversion.rs | 2 +- .../examples/bench_peers.rs | 4 +- .../swarm-coordination-registry/src/lib.rs | 10 +- .../src/statistics/event/handler.rs | 10 +- .../src/swarm/coordinator.rs | 219 +++++++++------- .../src/swarm/registry.rs | 35 ++- .../benches/helpers/utils.rs | 2 +- .../src/entry/peer_list.rs | 25 +- .../src/entry/single.rs | 8 +- .../tests/entry/mod.rs | 19 +- .../tests/repository/mod.rs | 2 +- packages/tracker-core/src/announce_handler.rs | 18 +- packages/tracker-core/src/peer_tests.rs | 2 +- packages/tracker-core/src/test_helpers.rs | 6 +- packages/tracker-core/src/torrent/services.rs | 4 +- .../tracker-core/tests/common/fixtures.rs | 2 +- .../tracker-core/tests/common/test_env.rs | 2 +- packages/udp-core/src/peer_builder.rs | 2 +- packages/udp-server/src/handlers/announce.rs | 21 +- packages/udp-server/src/lib.rs | 2 +- 37 files changed, 847 insertions(+), 336 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index b87d50f73..907e0a889 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5124,6 +5124,7 @@ dependencies = [ "torrust-info-hash", "torrust-located-error", "torrust-peer-id", + "torrust-tracker-primitives", ] [[package]] @@ -5148,10 +5149,12 @@ dependencies = [ name = "torrust-tracker-primitives" version = "3.0.0" dependencies = [ + "base64", "binascii", "derive_more 2.1.1", "serde", "serde_json", + "sha2 0.11.0", "tdyne-peer-id", "tdyne-peer-id-registry", "thiserror 2.0.19", diff --git a/packages/axum-http-server/src/v1/extractors/announce_request.rs b/packages/axum-http-server/src/v1/extractors/announce_request.rs index b6072d29c..e21e23c76 100644 --- a/packages/axum-http-server/src/v1/extractors/announce_request.rs +++ b/packages/axum-http-server/src/v1/extractors/announce_request.rs @@ -88,7 +88,7 @@ mod tests { use std::str::FromStr; use torrust_info_hash::InfoHash; - use torrust_tracker_http_protocol::v1::requests::announce::{Announce, Compact, Event, NumberOfBytes}; + use torrust_tracker_http_protocol::v1::requests::announce::{Announce, AnnounceAddress, Compact, Event, NumberOfBytes}; use torrust_tracker_http_protocol::v1::responses::error::Error; use torrust_tracker_primitives::PeerId; @@ -113,7 +113,7 @@ mod tests { info_hash: InfoHash::from_str("3b245504cf5f11bbdbe1201cea6a6bf45aee1bc0").unwrap(), // DevSkim: ignore DS173237 peer_id: PeerId(*b"-qB00000000000000001"), port: 17548, - ip: Some(IpAddr::V4(Ipv4Addr::new(2, 137, 87, 41))), + ip: Some(AnnounceAddress::Ip(IpAddr::V4(Ipv4Addr::new(2, 137, 87, 41)))), downloaded: Some(NumberOfBytes::new(0)), uploaded: Some(NumberOfBytes::new(0)), left: Some(NumberOfBytes::new(0)), diff --git a/packages/axum-http-server/src/v1/handlers/announce.rs b/packages/axum-http-server/src/v1/handlers/announce.rs index 73e4868bd..c3625d6d6 100644 --- a/packages/axum-http-server/src/v1/handlers/announce.rs +++ b/packages/axum-http-server/src/v1/handlers/announce.rs @@ -106,9 +106,22 @@ fn to_protocol_announce_data(domain_data: DomainAnnounceData) -> responses::anno peers: domain_data .peers .into_iter() - .map(|peer| responses::announce::Peer { - peer_id: peer.peer_id, - peer_addr: peer.peer_addr, + .map(|peer| { + let peer_addr = match &peer.peer_addr { + torrust_tracker_primitives::PeerAddress::Clearnet(address) => { + responses::announce::PeerAddress::Clearnet(*address) + } + torrust_tracker_primitives::PeerAddress::I2p(address) => responses::announce::PeerAddress::I2p { + destination: address.destination.to_string(), + destination_hash: *address.destination.hash(), + port: 1, + }, + }; + + responses::announce::Peer { + peer_id: peer.peer_id, + peer_addr, + } }) .collect(), stats: responses::announce::SwarmMetadata { diff --git a/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs b/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs index bbd6c68c6..fc5d45152 100644 --- a/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs +++ b/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs @@ -24,10 +24,10 @@ use torrust_tracker_client::http::client::Client; use torrust_tracker_http_protocol::percent_encoding::percent_encode_byte_array; use torrust_tracker_http_protocol::v1::requests::announce::{AnnounceBuilder, Compact}; use torrust_tracker_http_protocol::v1::responses::announce::deserialization::{ - CompactPeer, CompactPeerList, DeserializedNormal, DictionaryPeer, + CompactPeer, CompactPeerList, DeserializedCompact, DeserializedNormal, DictionaryPeer, }; -use torrust_tracker_primitives::PeerId as DomainPeerId; use torrust_tracker_primitives::peer::fixture::PeerBuilder; +use torrust_tracker_primitives::{I2pDestination, PeerId as DomainPeerId}; use torrust_tracker_test_helpers::{configuration, logging}; use crate::common::fixtures::invalid_info_hashes; @@ -594,7 +594,7 @@ async fn should_return_the_list_of_previously_announced_peers() { min_interval: announce_policy.interval_min, peers: vec![DictionaryPeer { peer_id: previously_announced_peer.peer_id.as_bytes().to_vec(), - ip: previously_announced_peer.peer_addr.ip().to_string(), + ip: previously_announced_peer.peer_addr.ip().unwrap().to_string(), port: previously_announced_peer.peer_addr.port(), }], }, @@ -658,12 +658,12 @@ async fn should_return_the_list_of_previously_announced_peers_including_peers_us peers: vec![ DictionaryPeer { peer_id: peer_using_ipv4.peer_id.as_bytes().to_vec(), - ip: peer_using_ipv4.peer_addr.ip().to_string(), + ip: peer_using_ipv4.peer_addr.ip().unwrap().to_string(), port: peer_using_ipv4.peer_addr.port(), }, DictionaryPeer { peer_id: peer_using_ipv6.peer_id.as_bytes().to_vec(), - ip: peer_using_ipv6.peer_addr.ip().to_string(), + ip: peer_using_ipv6.peer_addr.ip().unwrap().to_string(), port: peer_using_ipv6.peer_addr.port(), }, ], @@ -689,14 +689,14 @@ async fn should_consider_two_peers_to_be_the_same_when_they_have_the_same_socket let announce_query_1 = AnnounceBuilder::default() .with_info_hash(&info_hash) .with_peer_id(&PeerId(peer.peer_id.0)) - .with_ip(peer.peer_addr.ip()) + .with_ip(peer.peer_addr.ip().unwrap()) .with_port(peer.peer_addr.port()) .query(); let announce_query_2 = AnnounceBuilder::default() .with_info_hash(&info_hash) .with_peer_id(&PeerId(*b"-qB00000000000000002")) // Different peer ID - .with_ip(peer.peer_addr.ip()) + .with_ip(peer.peer_addr.ip().unwrap()) .with_port(peer.peer_addr.port()) .query(); @@ -776,7 +776,7 @@ async fn should_return_the_compact_response() { incomplete: 0, interval: 120, min_interval: 120, - peers: CompactPeerList::new([CompactPeer::new(&previously_announced_peer.peer_addr)].to_vec()), + peers: CompactPeerList::new([CompactPeer::new(&previously_announced_peer.peer_addr.socket_addr().unwrap())].to_vec()), }; assert_compact_announce_response(response, &expected_response).await; @@ -784,6 +784,58 @@ async fn should_return_the_compact_response() { env.stop().await; } +#[tokio::test] +async fn it_should_return_i2p_destination_hashes_in_a_compact_response() { + logging::setup(); + + let cfg = configuration::ephemeral_public(); + let core_config = Arc::new(cfg.core.clone()); + let http_tracker_config = Arc::new(cfg.http_trackers.unwrap()[0].clone()); + let env = Started::new(&core_config, &http_tracker_config).await; + let client = Client::new(env.base_url(), Duration::from_secs(5)).unwrap(); + let info_hash = InfoHash::from_str("9c38422213e30bff212b30c360d26f9a02136422").unwrap(); // DevSkim: ignore DS173237 + let first_destination = format!("{}BQAEAAAAAA==.i2p", "A".repeat(512)) + .parse::() + .unwrap(); + + client + .announce( + &AnnounceBuilder::default() + .with_info_hash(&info_hash) + .with_peer_id(&PeerId(*b"-qB00000000000000001")) + .with_port(1) + .with_i2p_destination(first_destination.clone()) + .with_compact(Compact::Accepted) + .query(), + ) + .await + .unwrap(); + + let response = client + .announce( + &AnnounceBuilder::default() + .with_info_hash(&info_hash) + .with_peer_id(&PeerId(*b"-qB00000000000000002")) + .with_port(1) + .with_i2p_destination( + format!("B{}BQAEAAAAAA==.i2p", "A".repeat(511)) + .parse::() + .unwrap(), + ) + .with_compact(Compact::Accepted) + .query(), + ) + .await + .unwrap(); + let bytes = response.bytes().await.unwrap(); + let announce = DeserializedCompact::from_bytes(&bytes).unwrap(); + + assert_eq!(announce.peers, *first_destination.hash()); + assert!(announce.peers6.is_empty()); + + env.stop().await; +} + #[tokio::test] async fn should_return_the_compact_response_by_default() { logging::setup(); @@ -942,10 +994,10 @@ async fn should_assign_to_the_peer_ip_the_remote_client_ip_instead_of_the_peer_a .in_memory_torrent_repository .get_torrent_peers(&info_hash, usize::MAX) .await; - let peer_addr = peers[0].peer_addr; + let peer_addr = &peers[0].peer_addr; - assert_eq!(peer_addr.ip(), client_ip); - assert_ne!(peer_addr.ip(), IpAddr::from_str("2.2.2.2").unwrap()); + assert_eq!(peer_addr.ip(), Some(client_ip)); + assert_ne!(peer_addr.ip(), Some(IpAddr::from_str("2.2.2.2").unwrap())); env.stop().await; } @@ -986,7 +1038,7 @@ async fn when_the_client_ip_is_a_loopback_ipv4_it_should_assign_to_the_peer_ip_t .in_memory_torrent_repository .get_torrent_peers(&info_hash, usize::MAX) .await; - let peer_addr = peers[0].peer_addr; + let peer_addr = &peers[0].peer_addr; let ext_ip: IpAddr = env .container @@ -996,8 +1048,8 @@ async fn when_the_client_ip_is_a_loopback_ipv4_it_should_assign_to_the_peer_ip_t .external_ip .unwrap() .into(); - assert_eq!(peer_addr.ip(), ext_ip); - assert_ne!(peer_addr.ip(), IpAddr::from_str("2.2.2.2").unwrap()); + assert_eq!(peer_addr.ip(), Some(ext_ip)); + assert_ne!(peer_addr.ip(), Some(IpAddr::from_str("2.2.2.2").unwrap())); env.stop().await; } @@ -1039,7 +1091,7 @@ async fn when_the_client_ip_is_a_loopback_ipv6_it_should_assign_to_the_peer_ip_t .in_memory_torrent_repository .get_torrent_peers(&info_hash, usize::MAX) .await; - let peer_addr = peers[0].peer_addr; + let peer_addr = &peers[0].peer_addr; let ext_ip: IpAddr = env .container @@ -1049,8 +1101,8 @@ async fn when_the_client_ip_is_a_loopback_ipv6_it_should_assign_to_the_peer_ip_t .external_ip .unwrap() .into(); - assert_eq!(peer_addr.ip(), ext_ip); - assert_ne!(peer_addr.ip(), IpAddr::from_str("2.2.2.2").unwrap()); + assert_eq!(peer_addr.ip(), Some(ext_ip)); + assert_ne!(peer_addr.ip(), Some(IpAddr::from_str("2.2.2.2").unwrap())); env.stop().await; } @@ -1095,9 +1147,9 @@ async fn when_the_tracker_is_behind_a_reverse_proxy_it_should_assign_to_the_peer .in_memory_torrent_repository .get_torrent_peers(&info_hash, usize::MAX) .await; - let peer_addr = peers[0].peer_addr; + let peer_addr = &peers[0].peer_addr; - assert_eq!(peer_addr.ip(), IpAddr::from_str("150.172.238.178").unwrap()); + assert_eq!(peer_addr.ip(), Some(IpAddr::from_str("150.172.238.178").unwrap())); env.stop().await; } diff --git a/packages/http-core/benches/helpers/util.rs b/packages/http-core/benches/helpers/util.rs index 7faf6f86e..62114e01c 100644 --- a/packages/http-core/benches/helpers/util.rs +++ b/packages/http-core/benches/helpers/util.rs @@ -93,7 +93,7 @@ pub async fn initialize_core_tracker_services_with_config( pub fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -123,7 +123,7 @@ pub fn sample_announce_request_for_peer(peer: Peer) -> (Announce, ClientIpSource let client_ip_sources = ClientIpSources { right_most_x_forwarded_for: None, - connection_info_socket_address: Some(SocketAddr::new(peer.peer_addr.ip(), 8080)), + connection_info_socket_address: Some(SocketAddr::new(peer.peer_addr.ip().unwrap(), 8080)), }; (announce_request, client_ip_sources) diff --git a/packages/http-core/src/lib.rs b/packages/http-core/src/lib.rs index fc6f5b068..13becb9aa 100644 --- a/packages/http-core/src/lib.rs +++ b/packages/http-core/src/lib.rs @@ -45,14 +45,15 @@ pub(crate) mod tests { peer.peer_addr = SocketAddr::new( IpAddr::V6(Ipv6Addr::new(0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969)), 8080, - ); + ) + .into(); peer } pub fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/http-core/src/services/announce.rs b/packages/http-core/src/services/announce.rs index 608943534..a5157f958 100644 --- a/packages/http-core/src/services/announce.rs +++ b/packages/http-core/src/services/announce.rs @@ -19,14 +19,14 @@ use torrust_tracker_core::authentication::{self, Key}; use torrust_tracker_core::error::{AnnounceError, TrackerCoreError, WhitelistError}; use torrust_tracker_core::whitelist; use torrust_tracker_http_protocol::v1::requests::announce::{ - Announce, Event as ProtocolAnnounceEvent, NumberOfBytes as ProtocolNumberOfBytes, + Announce, AnnounceAddress, Event as ProtocolAnnounceEvent, NumberOfBytes as ProtocolNumberOfBytes, }; use torrust_tracker_http_protocol::v1::responses::error::Error as HttpProtocolErrorResponse; use torrust_tracker_http_protocol::v1::services::peer_ip_resolver::{ ClientIpSources, PeerIpResolutionError, RemoteClientAddr, resolve_remote_client_addr, }; use torrust_tracker_primitives::peer::PeerAnnouncement; -use torrust_tracker_primitives::{AnnounceData, AnnounceEvent, NumberOfBytes}; +use torrust_tracker_primitives::{AnnounceData, AnnounceEvent, I2pPeerAddress, NumberOfBytes, PeerAddress}; use crate::event; use crate::event::Event; @@ -121,7 +121,12 @@ impl AnnounceService { PeerAnnouncement { peer_id: announce_request.peer_id, - peer_addr: std::net::SocketAddr::new(*peer_ip, announce_request.port), + peer_addr: match &announce_request.ip { + Some(AnnounceAddress::I2p(destination)) => PeerAddress::I2p(I2pPeerAddress { + destination: destination.clone(), + }), + _ => std::net::SocketAddr::new(*peer_ip, announce_request.port).into(), + }, updated: ::now(), uploaded: NumberOfBytes::new(uploaded.0), downloaded: NumberOfBytes::new(downloaded.0), @@ -358,7 +363,7 @@ mod tests { let client_ip_sources = ClientIpSources { right_most_x_forwarded_for: None, - connection_info_socket_address: Some(SocketAddr::new(peer.peer_addr.ip(), 8080)), + connection_info_socket_address: Some(SocketAddr::new(peer.peer_addr.ip().unwrap(), 8080)), }; (announce_request, client_ip_sources) @@ -392,9 +397,10 @@ mod tests { use mockall::predicate::{self}; use torrust_net_primitives::service_binding::{Protocol, ServiceBinding}; use torrust_tracker_configuration::Configuration; + use torrust_tracker_http_protocol::v1::requests::announce::AnnounceAddress; use torrust_tracker_http_protocol::v1::services::peer_ip_resolver::{RemoteClientAddr, ResolvedIp}; use torrust_tracker_primitives::swarm_metadata::SwarmMetadata; - use torrust_tracker_primitives::{AnnounceData, peer}; + use torrust_tracker_primitives::{AnnounceData, I2pDestination, PeerAddress, PeerId, peer}; use torrust_tracker_test_helpers::configuration; use crate::event::test::announce_events_match; @@ -443,11 +449,50 @@ mod tests { assert_eq!(announce_data, expected_announce_data); } + #[tokio::test] + async fn it_should_coordinate_i2p_peers_by_their_destinations() { + let (core_tracker_services, core_http_tracker_services) = initialize_core_tracker_services().await; + let server_socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 7070); + let server_service_binding = ServiceBinding::new(Protocol::HTTP, server_socket_addr).unwrap(); + let announce_service = AnnounceService::new( + core_tracker_services.core_config, + core_tracker_services.announce_handler, + core_tracker_services.authentication_service, + core_tracker_services.whitelist_authorization, + core_http_tracker_services.http_stats_event_sender, + ); + + let (mut first_request, client_ip_sources) = sample_announce_request_for_peer(sample_peer()); + first_request.ip = Some(AnnounceAddress::I2p( + format!("{}.i2p", "A".repeat(516)).parse::().unwrap(), + )); + announce_service + .handle_announce(&first_request, &client_ip_sources, &server_service_binding, None) + .await + .unwrap(); + + let mut second_peer = sample_peer(); + second_peer.peer_id = PeerId(*b"-qB00000000000000002"); + let (mut second_request, _) = sample_announce_request_for_peer(second_peer); + second_request.ip = Some(AnnounceAddress::I2p( + format!("{}.i2p", "B".repeat(516)).parse::().unwrap(), + )); + + let announce_data = announce_service + .handle_announce(&second_request, &client_ip_sources, &server_service_binding, None) + .await + .unwrap(); + + assert_eq!(announce_data.peers.len(), 1); + assert!(matches!(announce_data.peers[0].peer_addr, PeerAddress::I2p(_))); + } + #[tokio::test] async fn it_should_send_the_tcp_4_announce_event_when_the_peer_uses_ipv4() { let server_socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 7070); let server_service_binding = ServiceBinding::new(Protocol::HTTP, server_socket_addr).unwrap(); let peer = sample_peer_using_ipv4(); + let expected_peer = peer.clone(); let remote_client_ip = IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)); let server_service_binding_clone = server_service_binding.clone(); @@ -456,8 +501,8 @@ mod tests { http_stats_event_sender_mock .expect_send() .with(predicate::function(move |event| { - let mut announcement = peer; - announcement.peer_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080); + let mut announcement = expected_peer.clone(); + announcement.peer_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(); let expected_event = Event::TcpAnnounce { connection: ConnectionContext::new( @@ -507,7 +552,7 @@ mod tests { fn peer_with_the_ipv4_loopback_ip() -> peer::Peer { let loopback_ip = IpAddr::V4(Ipv4Addr::LOCALHOST); let mut peer = sample_peer(); - peer.peer_addr = SocketAddr::new(loopback_ip, 8080); + peer.peer_addr = SocketAddr::new(loopback_ip, 8080).into(); peer } @@ -519,6 +564,7 @@ mod tests { let server_socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 7070); let server_service_binding = ServiceBinding::new(Protocol::HTTP, server_socket_addr).unwrap(); let peer = peer_with_the_ipv4_loopback_ip(); + let expected_peer = peer.clone(); let remote_client_ip = IpAddr::V4(Ipv4Addr::LOCALHOST); let server_service_binding_clone = server_service_binding.clone(); @@ -527,11 +573,12 @@ mod tests { http_stats_event_sender_mock .expect_send() .with(predicate::function(move |event| { - let mut peer_announcement = peer; + let mut peer_announcement = expected_peer.clone(); peer_announcement.peer_addr = SocketAddr::new( IpAddr::V6(Ipv6Addr::new(0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969)), 8080, - ); + ) + .into(); let expected_event = Event::TcpAnnounce { connection: ConnectionContext::new( @@ -576,6 +623,7 @@ mod tests { let server_socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 7070); let server_service_binding = ServiceBinding::new(Protocol::HTTP, server_socket_addr).unwrap(); let peer = sample_peer_using_ipv6(); + let expected_peer = peer.clone(); let remote_client_ip = IpAddr::V6(Ipv6Addr::new(0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969)); let mut http_stats_event_sender_mock = MockHttpStatsEventSender::new(); @@ -588,7 +636,7 @@ mod tests { server_service_binding.clone(), ), info_hash: sample_info_hash(), - announcement: peer, + announcement: expected_peer.clone(), }; announce_events_match(event, &expected_event) })) diff --git a/packages/http-core/src/services/scrape.rs b/packages/http-core/src/services/scrape.rs index fa5b7dfe9..8a4d42f36 100644 --- a/packages/http-core/src/services/scrape.rs +++ b/packages/http-core/src/services/scrape.rs @@ -238,7 +238,7 @@ mod tests { fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -299,7 +299,7 @@ mod tests { // Announce a new peer to force scrape data to contain non zeroed data let mut peer = sample_peer(); - let original_peer_ip = peer.ip(); + let original_peer_ip = peer.ip().unwrap(); container .announce_handler .handle_announcement(&info_hash, &mut peer, &original_peer_ip, &PeersWanted::AsManyAsPossible) @@ -489,7 +489,7 @@ mod tests { // Announce a new peer to force scrape data to contain non zeroed data let mut peer = sample_peer(); - let original_peer_ip = peer.ip(); + let original_peer_ip = peer.ip().unwrap(); container .announce_handler .handle_announcement(&info_hash, &mut peer, &original_peer_ip, &PeersWanted::AsManyAsPossible) diff --git a/packages/http-protocol/Cargo.toml b/packages/http-protocol/Cargo.toml index ebaabcfa7..0d7377982 100644 --- a/packages/http-protocol/Cargo.toml +++ b/packages/http-protocol/Cargo.toml @@ -28,6 +28,7 @@ thiserror = "2" torrust-clock = "3.0.0" torrust-bencode = "3.0.0" torrust-located-error = "3.0.0" +torrust-tracker-primitives = { version = "3.0.0", path = "../primitives" } [package.metadata.cargo-machete] ignored = [ "serde_bytes" ] diff --git a/packages/http-protocol/src/v1/query.rs b/packages/http-protocol/src/v1/query.rs index 878033423..9aa29635b 100644 --- a/packages/http-protocol/src/v1/query.rs +++ b/packages/http-protocol/src/v1/query.rs @@ -97,8 +97,7 @@ impl Query { /// from a string. #[derive(Error, Debug)] pub enum ParseQueryError { - /// Invalid URL query param. For example: `"name=value=value"`. It contains - /// an unescaped `=` character. + /// Invalid URL query parameter without a name/value separator. #[error("invalid param {raw_param} in {location}")] InvalidParam { location: &'static Location<'static>, @@ -168,18 +167,14 @@ impl FromStr for NameValuePair { type Err = ParseQueryError; fn from_str(raw_param: &str) -> Result { - let pair = raw_param.split('=').collect::>(); - - if pair.len() != 2 { - return Err(ParseQueryError::InvalidParam { - location: Location::caller(), - raw_param: raw_param.to_owned(), - }); - } + let (name, value) = raw_param.split_once('=').ok_or_else(|| ParseQueryError::InvalidParam { + location: Location::caller(), + raw_param: raw_param.to_owned(), + })?; Ok(Self { - name: pair[0].to_owned(), - value: pair[1].to_owned(), + name: name.to_owned(), + value: value.to_owned(), }) } } @@ -257,12 +252,21 @@ mod tests { } #[test] - fn should_fail_parsing_an_invalid_query_string() { - let invalid_raw_query = "name=value=value"; + fn it_should_preserve_equals_characters_in_a_query_parameter_value() { + let raw_query = "name=value=="; + + let query = raw_query.parse::().unwrap(); + + assert_eq!(query.get_param("name"), Some("value==".to_string())); + } + + #[test] + fn it_should_reject_a_query_parameter_without_a_separator() { + let invalid_raw_query = "name"; - let query = invalid_raw_query.parse::(); + let result = invalid_raw_query.parse::(); - assert!(query.is_err()); + assert!(result.is_err()); } #[test] @@ -345,12 +349,21 @@ mod tests { } #[test] - fn should_fail_parsing_an_invalid_query_param() { - let invalid_raw_param = "name=value=value"; + fn it_should_preserve_equals_characters_in_the_value() { + let raw_param = "name=value=="; + + let param = raw_param.parse::().unwrap(); + + assert_eq!(param.value, "value=="); + } + + #[test] + fn it_should_reject_a_param_without_a_separator() { + let invalid_raw_param = "name"; - let query = invalid_raw_param.parse::(); + let result = invalid_raw_param.parse::(); - assert!(query.is_err()); + assert!(result.is_err()); } #[test] diff --git a/packages/http-protocol/src/v1/requests/announce.rs b/packages/http-protocol/src/v1/requests/announce.rs index 4f91e1ca1..055632d5c 100644 --- a/packages/http-protocol/src/v1/requests/announce.rs +++ b/packages/http-protocol/src/v1/requests/announce.rs @@ -12,6 +12,7 @@ use thiserror::Error; use torrust_info_hash::InfoHash; use torrust_located_error::{Located, LocatedError}; use torrust_peer_id::PeerId; +use torrust_tracker_primitives::I2pDestination; use crate::percent_encoding::{ PeerIdConversionError, percent_decode_info_hash, percent_decode_peer_id, percent_encode_byte_array, @@ -96,8 +97,8 @@ pub struct Announce { pub port: u16, // Optional params - /// The peer IP address (BEP 3 `ip` parameter). - pub ip: Option, + /// The peer IP address or I2P Destination (BEP 3 `ip` parameter). + pub ip: Option, /// The number of bytes downloaded by the peer. pub downloaded: Option, @@ -120,6 +121,22 @@ pub struct Announce { pub numwant: Option, } +/// Address supplied in the BEP 3 `ip` parameter. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum AnnounceAddress { + Ip(IpAddr), + I2p(I2pDestination), +} + +impl fmt::Display for AnnounceAddress { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Ip(address) => address.fmt(f), + Self::I2p(destination) => destination.fmt(f), + } + } +} + /// Errors that can occur when parsing the `Announce` request. /// /// The `info_hash` and `peer_id` query params are special because they contain @@ -291,7 +308,7 @@ impl TryFrom for Announce { event: extract_event(&query)?, compact: extract_compact(&query)?, numwant: extract_numwant(&query)?, - ip: extract_ip(&query), + ip: extract_ip(&query)?, }) } } @@ -376,7 +393,7 @@ impl AnnounceBuilder { info_hash: InfoHash::from_str("9c38422213e30bff212b30c360d26f9a02136422").unwrap(), // DevSkim: ignore DS173237 peer_id: PeerId(*b"-qB00000000000000001"), port: 17548, - ip: Some(IpAddr::V4(std::net::Ipv4Addr::new(192, 168, 1, 88))), + ip: Some(AnnounceAddress::Ip(IpAddr::V4(std::net::Ipv4Addr::new(192, 168, 1, 88)))), downloaded: None, uploaded: None, left: None, @@ -409,7 +426,13 @@ impl AnnounceBuilder { #[must_use] pub fn with_ip(mut self, ip: IpAddr) -> Self { - self.announce.ip = Some(ip); + self.announce.ip = Some(AnnounceAddress::Ip(ip)); + self + } + + #[must_use] + pub fn with_i2p_destination(mut self, destination: I2pDestination) -> Self { + self.announce.ip = Some(AnnounceAddress::I2p(destination)); self } @@ -563,10 +586,27 @@ fn extract_number_of_bytes_from_param(param_name: &str, query: &Query) -> Result } } -fn extract_ip(query: &Query) -> Option { +fn extract_ip(query: &Query) -> Result, ParseAnnounceQueryError> { match query.get_param(IP) { - Some(raw_param) => IpAddr::from_str(&raw_param).ok(), - None => None, + Some(raw_param) => { + if let Ok(ip) = IpAddr::from_str(&raw_param) { + return Ok(Some(AnnounceAddress::Ip(ip))); + } + + if raw_param.ends_with(".i2p") { + return I2pDestination::from_str(&raw_param) + .map(AnnounceAddress::I2p) + .map(Some) + .map_err(|_| ParseAnnounceQueryError::InvalidParam { + param_name: IP.to_owned(), + param_value: raw_param, + location: Location::caller(), + }); + } + + Ok(None) + } + None => Ok(None), } } @@ -608,8 +648,8 @@ mod tests { use crate::v1::query::Query; use crate::v1::requests::announce::{ - Announce, COMPACT, Compact, DOWNLOADED, EVENT, Event, INFO_HASH, LEFT, NUMWANT, NumberOfBytes, PEER_ID, PORT, - UPLOADED, + Announce, AnnounceAddress, COMPACT, Compact, DOWNLOADED, EVENT, Event, INFO_HASH, IP, LEFT, NUMWANT, NumberOfBytes, + PEER_ID, PORT, UPLOADED, }; #[test] @@ -678,6 +718,40 @@ mod tests { ); } + #[test] + fn it_should_parse_a_padded_i2p_destination_from_the_ip_param() { + // 391 decoded bytes: 384 key bytes, a key certificate with its + // four-byte key-type payload, and `==` Base64 padding. + let destination = format!("{}BQAEAAAAAA==.i2p", "A".repeat(512)); + let raw_query = Query::from(vec![ + (INFO_HASH, "%3B%24U%04%CF%5F%11%BB%DB%E1%20%1C%EAjk%F4Z%EE%1B%C0"), + (PEER_ID, "-RC3000-000000000001"), + (PORT, "1"), + (IP, &destination), + ]) + .to_string(); + + let announce_request = Announce::try_from(raw_query.parse::().unwrap()).unwrap(); + + assert!(matches!(announce_request.ip, Some(AnnounceAddress::I2p(_)))); + } + + #[test] + fn it_should_reject_an_invalid_i2p_destination() { + let destination = "invalid.i2p"; + let raw_query = Query::from(vec![ + (INFO_HASH, "%3B%24U%04%CF%5F%11%BB%DB%E1%20%1C%EAjk%F4Z%EE%1B%C0"), + (PEER_ID, "-RC3000-000000000001"), + (PORT, "1"), + (IP, destination), + ]) + .to_string(); + + let result = Announce::try_from(raw_query.parse::().unwrap()); + + assert!(result.is_err()); + } + mod when_it_is_instantiated_from_the_url_query_params { use crate::v1::query::Query; diff --git a/packages/http-protocol/src/v1/responses/announce/data.rs b/packages/http-protocol/src/v1/responses/announce/data.rs index 0034ca854..6f20b633d 100644 --- a/packages/http-protocol/src/v1/responses/announce/data.rs +++ b/packages/http-protocol/src/v1/responses/announce/data.rs @@ -51,8 +51,18 @@ impl SwarmMetadata { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash)] pub struct Peer { pub peer_id: PeerId, - pub peer_addr: SocketAddr, + pub peer_addr: PeerAddress, +} + +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub enum PeerAddress { + Clearnet(SocketAddr), + I2p { + destination: String, + destination_hash: [u8; 32], + port: u16, + }, } diff --git a/packages/http-protocol/src/v1/responses/announce/encoding.rs b/packages/http-protocol/src/v1/responses/announce/encoding.rs index d1c69a3ba..e11ac7766 100644 --- a/packages/http-protocol/src/v1/responses/announce/encoding.rs +++ b/packages/http-protocol/src/v1/responses/announce/encoding.rs @@ -2,13 +2,12 @@ //! //! Types for encoding announce responses into bencoded bytes. //! Supports two encoding forms: [`Normal`] (dictionary-based) and [`Compact`] (packed binary). -use std::io::Write; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr}; -use derive_more::{AsRef, Constructor, From}; +use derive_more::{AsRef, Constructor}; use torrust_bencode::{BMutAccess, BencodeMut, ben_bytes, ben_int, ben_list, ben_map}; -use crate::v1::responses::announce::data::{AnnounceData, Peer}; +use crate::v1::responses::announce::data::{AnnounceData, Peer, PeerAddress}; /// An [`Announce`] response, that can be anything that is convertible from [`AnnounceData`]. /// @@ -95,21 +94,30 @@ pub struct Compact { impl From for Compact { fn from(data: AnnounceData) -> Self { - let compact_peers: Vec = data.peers.into_iter().map(CompactPeer::from).collect(); + let mut peers = vec![]; + let mut peers6 = vec![]; - let (peers, peers6): (Vec>, Vec>) = - compact_peers.into_iter().collect(); - - let peers_encoded: CompactPeersEncoded = peers.into_iter().collect(); - let peers_encoded_6: CompactPeersEncoded = peers6.into_iter().collect(); + for peer in data.peers.into_iter().map(CompactPeer::from) { + match peer { + CompactPeer::V4(peer) => { + peers.extend(u32::from(peer.ip).to_be_bytes()); + peers.extend(peer.port.to_be_bytes()); + } + CompactPeer::V6(peer) => { + peers6.extend(u128::from(peer.ip).to_be_bytes()); + peers6.extend(peer.port.to_be_bytes()); + } + CompactPeer::I2p(hash) => peers.extend(hash), + } + } Self { complete: data.stats.complete.into(), incomplete: data.stats.incomplete.into(), interval: data.policy.interval.into(), min_interval: data.policy.interval_min.into(), - peers: peers_encoded.0, - peers6: peers_encoded_6.0, + peers, + peers6, } } } @@ -132,13 +140,12 @@ impl Into> for Compact { /// A [`NormalPeer`], for the [`Normal`] form. /// /// ```rust -/// use std::net::{IpAddr, Ipv4Addr}; /// use torrust_tracker_http_protocol::v1::responses::announce::{Normal, NormalPeer}; /// /// let peer = NormalPeer { /// peer_id: *b"-RC3000-000000000001", -/// ip: IpAddr::V4(Ipv4Addr::new(0x69, 0x69, 0x69, 0x69)), // 105.105.105.105 -/// port: 0x7070, // 28784 +/// ip: "105.105.105.105".to_owned(), +/// port: 0x7070, // 28784 /// }; /// /// ``` @@ -147,17 +154,24 @@ pub struct NormalPeer { /// The peer's ID. pub peer_id: [u8; 20], /// The peer's IP address. - pub ip: IpAddr, + pub ip: String, /// The peer's port number. pub port: u16, } impl From for NormalPeer { fn from(peer: Peer) -> Self { - NormalPeer { - peer_id: peer.peer_id.0, - ip: peer.peer_addr.ip(), - port: peer.peer_addr.port(), + match peer.peer_addr { + PeerAddress::Clearnet(address) => NormalPeer { + peer_id: peer.peer_id.0, + ip: address.ip().to_string(), + port: address.port(), + }, + PeerAddress::I2p { destination, port, .. } => NormalPeer { + peer_id: peer.peer_id.0, + ip: destination, + port, + }, } } } @@ -166,7 +180,7 @@ impl From<&NormalPeer> for BencodeMut<'_> { fn from(value: &NormalPeer) -> Self { ben_map! { "peer id" => ben_bytes!(value.peer_id.clone().to_vec()), - "ip" => ben_bytes!(value.ip.to_string()), + "ip" => ben_bytes!(value.ip.clone()), "port" => ben_int!(i64::from(value.port)) } } @@ -200,6 +214,8 @@ pub enum CompactPeer { V4(CompactPeerData), /// The peer's port number. V6(CompactPeerData), + /// The SHA-256 hash of an I2P Destination. + I2p([u8; 32]), } impl CompactPeer { @@ -246,9 +262,16 @@ impl CompactPeer { impl From for CompactPeer { fn from(peer: Peer) -> Self { - match (peer.peer_addr.ip(), peer.peer_addr.port()) { - (IpAddr::V4(ip), port) => Self::V4(CompactPeerData { ip, port }), - (IpAddr::V6(ip), port) => Self::V6(CompactPeerData { ip, port }), + match peer.peer_addr { + PeerAddress::Clearnet(SocketAddr::V4(address)) => Self::V4(CompactPeerData { + ip: *address.ip(), + port: address.port(), + }), + PeerAddress::Clearnet(SocketAddr::V6(address)) => Self::V6(CompactPeerData { + ip: *address.ip(), + port: address.port(), + }), + PeerAddress::I2p { destination_hash, .. } => Self::I2p(destination_hash), } } } @@ -263,54 +286,6 @@ pub struct CompactPeerData { pub port: u16, } -impl FromIterator for (Vec>, Vec>) { - fn from_iter>(iter: T) -> Self { - let mut peers_v4: Vec> = vec![]; - let mut peers_v6: Vec> = vec![]; - - for peer in iter { - match peer { - CompactPeer::V4(peer) => peers_v4.push(peer), - CompactPeer::V6(peer6) => peers_v6.push(peer6), - } - } - - (peers_v4, peers_v6) - } -} - -#[derive(From, PartialEq)] -struct CompactPeersEncoded(Vec); - -impl FromIterator> for CompactPeersEncoded { - fn from_iter>>(iter: T) -> Self { - let mut bytes: Vec = vec![]; - - for peer in iter { - bytes - .write_all(&u32::from(peer.ip).to_be_bytes()) - .expect("it should write peer ip"); - bytes.write_all(&peer.port.to_be_bytes()).expect("it should write peer port"); - } - - bytes.into() - } -} - -impl FromIterator> for CompactPeersEncoded { - fn from_iter>>(iter: T) -> Self { - let mut bytes: Vec = Vec::new(); - - for peer in iter { - bytes - .write_all(&u128::from(peer.ip).to_be_bytes()) - .expect("it should write peer ip"); - bytes.write_all(&peer.port.to_be_bytes()).expect("it should write peer port"); - } - bytes.into() - } -} - #[cfg(test)] mod tests { @@ -318,7 +293,9 @@ mod tests { use torrust_peer_id::PeerId; - use crate::v1::responses::announce::{Announce, AnnounceData, AnnouncePolicy, Compact, Normal, Peer, SwarmMetadata}; + use crate::v1::responses::announce::{ + Announce, AnnounceData, AnnouncePolicy, Compact, Normal, Peer, PeerAddress, SwarmMetadata, + }; // Some ascii values used in tests: // @@ -337,15 +314,15 @@ mod tests { let peer_ipv4 = Peer { peer_id: PeerId(*b"-RC3000-000000000001"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(0x69, 0x69, 0x69, 0x69)), 0x7070), + peer_addr: PeerAddress::Clearnet(SocketAddr::new(IpAddr::V4(Ipv4Addr::new(0x69, 0x69, 0x69, 0x69)), 0x7070)), }; let peer_ipv6 = Peer { peer_id: PeerId(*b"-RC3000-000000000002"), - peer_addr: SocketAddr::new( + peer_addr: PeerAddress::Clearnet(SocketAddr::new( IpAddr::V6(Ipv6Addr::new(0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969, 0x6969)), 0x7070, - ), + )), }; let peers = vec![peer_ipv4, peer_ipv6]; @@ -382,4 +359,52 @@ mod tests { String::from_utf8(expected_bytes.to_vec()).unwrap() ); } + + #[test] + fn it_should_encode_an_i2p_peer_as_a_destination_in_a_non_compact_response() { + let destination = format!("{}.i2p", "A".repeat(516)); + let data = AnnounceData::new( + vec![Peer { + peer_id: PeerId(*b"-RC3000-000000000001"), + peer_addr: PeerAddress::I2p { + destination: destination.clone(), + destination_hash: [7; 32], + port: 1, + }, + }], + SwarmMetadata::default(), + AnnouncePolicy::default(), + ); + + let response: Announce = data.into(); + let bytes: Vec = response.data.into(); + let decoded = serde_bencode::from_bytes::(&bytes).unwrap(); + + assert_eq!(decoded.peers[0].ip, destination); + assert_eq!(decoded.peers[0].port, 1); + } + + #[test] + fn it_should_encode_an_i2p_peer_hash_in_a_compact_response() { + let destination_hash = [7; 32]; + let data = AnnounceData::new( + vec![Peer { + peer_id: PeerId(*b"-RC3000-000000000001"), + peer_addr: PeerAddress::I2p { + destination: format!("{}.i2p", "A".repeat(516)), + destination_hash, + port: 1, + }, + }], + SwarmMetadata::default(), + AnnouncePolicy::default(), + ); + + let response: Announce = data.into(); + let bytes: Vec = response.data.into(); + let decoded = crate::v1::responses::announce::DeserializedCompact::from_bytes(&bytes).unwrap(); + + assert_eq!(decoded.peers, destination_hash); + assert!(decoded.peers6.is_empty()); + } } diff --git a/packages/http-protocol/src/v1/responses/announce/mod.rs b/packages/http-protocol/src/v1/responses/announce/mod.rs index 57d746382..eb5b27d8a 100644 --- a/packages/http-protocol/src/v1/responses/announce/mod.rs +++ b/packages/http-protocol/src/v1/responses/announce/mod.rs @@ -3,6 +3,6 @@ pub mod data; pub mod deserialization; pub mod encoding; -pub use data::{AnnounceData, AnnouncePolicy, Peer, SwarmMetadata}; +pub use data::{AnnounceData, AnnouncePolicy, Peer, PeerAddress, SwarmMetadata}; pub use deserialization::{CompactPeerList, DeserializedCompact, DeserializedCompactParsed, DeserializedNormal, DictionaryPeer}; pub use encoding::{Announce, Compact, CompactPeer, CompactPeerData, Normal, NormalPeer}; diff --git a/packages/primitives/Cargo.toml b/packages/primitives/Cargo.toml index e2f35549b..f03029c27 100644 --- a/packages/primitives/Cargo.toml +++ b/packages/primitives/Cargo.toml @@ -16,10 +16,12 @@ version = "3.0.0" [dependencies] torrust-peer-id = "0.1.0" +base64 = "0.22.1" binascii = "0" torrust-info-hash = "=0.2.0" derive_more = { version = "2", features = [ "constructor", "display" ] } serde = { version = "1", features = [ "derive" ] } +sha2 = "0.11.0" tdyne-peer-id = "1" tdyne-peer-id-registry = "0" thiserror = "2" diff --git a/packages/primitives/src/lib.rs b/packages/primitives/src/lib.rs index f24f7462c..e2b23c179 100644 --- a/packages/primitives/src/lib.rs +++ b/packages/primitives/src/lib.rs @@ -29,6 +29,7 @@ pub use configuration_instance_id::ConfigurationInstanceId; pub use driver::Driver; pub use mode::PrivateMode; pub use number_of_bytes::NumberOfBytes; +pub use peer::{I2pDestination, I2pPeerAddress, PeerAddress}; pub use policy::TrackerPolicy; pub use scrape::ScrapeData; pub use service_role::ServiceRole; diff --git a/packages/primitives/src/peer.rs b/packages/primitives/src/peer.rs index 1e3678e78..16050a5e0 100644 --- a/packages/primitives/src/peer.rs +++ b/packages/primitives/src/peer.rs @@ -13,7 +13,7 @@ //! //! peer::Peer { //! peer_id: PeerId(*b"-qB00000000000000000"), -//! peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), +//! peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), //! updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), //! uploaded: NumberOfBytes::new(0), //! downloaded: NumberOfBytes::new(0), @@ -28,13 +28,156 @@ use std::ops::{Deref, DerefMut}; use std::str::FromStr; use std::sync::Arc; +use base64::Engine; +use base64::alphabet::Alphabet; +use base64::engine::{GeneralPurpose, GeneralPurposeConfig}; use serde::Serialize; +use sha2::{Digest, Sha256}; +use thiserror::Error; use torrust_clock::DurationSinceUnixEpoch; use crate::{AnnounceEvent, NumberOfBytes, PeerId}; pub type PeerAnnouncement = Peer; +const I2P_BASE64_ALPHABET: &str = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-~"; +const I2P_SUFFIX: &str = ".i2p"; +const MIN_I2P_DESTINATION_BYTES: usize = 387; +const I2P_CERTIFICATE_LENGTH_OFFSET: usize = 385; + +/// A validated I2P Base64 Destination. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct I2pDestination { + value: Box, + hash: [u8; 32], +} + +impl I2pDestination { + /// Returns the SHA-256 hash of the decoded binary Destination. + #[must_use] + pub const fn hash(&self) -> &[u8; 32] { + &self.hash + } +} + +impl fmt::Display for I2pDestination { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(&self.value) + } +} + +impl FromStr for I2pDestination { + type Err = ParseI2pDestinationError; + + fn from_str(value: &str) -> Result { + let encoded = value.strip_suffix(I2P_SUFFIX).unwrap_or(value); + let alphabet = + Alphabet::new(I2P_BASE64_ALPHABET).expect("the I2P Base64 alphabet must contain 64 unique ASCII characters"); + let engine = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()); + let decoded = engine.decode(encoded).map_err(|_| ParseI2pDestinationError::InvalidBase64)?; + + if decoded.len() < MIN_I2P_DESTINATION_BYTES { + return Err(ParseI2pDestinationError::TooShort { actual: decoded.len() }); + } + + let certificate_payload_length = usize::from(u16::from_be_bytes([ + decoded[I2P_CERTIFICATE_LENGTH_OFFSET], + decoded[I2P_CERTIFICATE_LENGTH_OFFSET + 1], + ])); + let expected_length = MIN_I2P_DESTINATION_BYTES + certificate_payload_length; + + if decoded.len() != expected_length { + return Err(ParseI2pDestinationError::InvalidCertificateLength { + declared: certificate_payload_length, + actual: decoded.len() - MIN_I2P_DESTINATION_BYTES, + }); + } + + Ok(Self { + value: format!("{encoded}{I2P_SUFFIX}").into_boxed_str(), + hash: Sha256::digest(decoded).into(), + }) + } +} + +/// Error returned when parsing an I2P Destination. +#[derive(Debug, Error, PartialEq, Eq)] +pub enum ParseI2pDestinationError { + #[error("the I2P Destination is not valid I2P Base64")] + InvalidBase64, + #[error("the decoded I2P Destination must contain at least {MIN_I2P_DESTINATION_BYTES} bytes, got {actual}")] + TooShort { actual: usize }, + #[error("the I2P certificate declares a {declared}-byte payload, but the Destination contains {actual} payload bytes")] + InvalidCertificateLength { declared: usize, actual: usize }, +} + +/// An I2P peer address. I2P routes by Destination and has no peer port. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct I2pPeerAddress { + pub destination: I2pDestination, +} + +/// A peer endpoint on either the public Internet or I2P. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub enum PeerAddress { + Clearnet(SocketAddr), + I2p(I2pPeerAddress), +} + +impl PeerAddress { + #[must_use] + pub const fn port(&self) -> u16 { + match self { + Self::Clearnet(address) => address.port(), + // I2P clients ignore the port, but legacy tracker response parsers + // expect the key to exist in non-compact responses. + Self::I2p(_) => 1, + } + } + + #[must_use] + pub const fn ip(&self) -> Option { + match self { + Self::Clearnet(address) => Some(address.ip()), + Self::I2p(_) => None, + } + } + + #[must_use] + pub const fn is_i2p(&self) -> bool { + matches!(self, Self::I2p(_)) + } + + #[must_use] + pub const fn socket_addr(&self) -> Option { + match self { + Self::Clearnet(address) => Some(*address), + Self::I2p(_) => None, + } + } +} + +impl From for PeerAddress { + fn from(value: SocketAddr) -> Self { + Self::Clearnet(value) + } +} + +impl fmt::Display for PeerAddress { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Clearnet(address) => address.fmt(f), + Self::I2p(address) => address.destination.fmt(f), + } + } +} + +impl Serialize for PeerAddress { + fn serialize(&self, serializer: S) -> Result { + serializer.serialize_str(&self.to_string()) + } +} + #[derive(Debug, Serialize, Copy, Clone, PartialEq, Eq, Hash)] #[serde(rename_all_fields = "lowercase")] pub enum PeerRole { @@ -101,7 +244,7 @@ pub enum ParsePeerRoleError { /// /// peer::Peer { /// peer_id: PeerId(*b"-qB00000000000000000"), -/// peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), +/// peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), /// updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), /// uploaded: NumberOfBytes::new(0), /// downloaded: NumberOfBytes::new(0), @@ -109,13 +252,13 @@ pub enum ParsePeerRoleError { /// event: AnnounceEvent::Started, /// }; /// ``` -#[derive(Debug, Clone, Serialize, Copy, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, Serialize, PartialEq, Eq, Hash)] pub struct Peer { /// ID used by the downloader peer #[serde(serialize_with = "ser_peer_id")] pub peer_id: PeerId, - /// The IP and port this peer is listening on - pub peer_addr: SocketAddr, + /// The clearnet socket address or I2P Destination for this peer. + pub peer_addr: PeerAddress, /// The last time the the tracker receive an announce request from this peer (timestamp) #[serde(serialize_with = "ser_unix_time_value")] pub updated: DurationSinceUnixEpoch, @@ -203,7 +346,7 @@ pub trait ReadInfo { fn get_event(&self) -> AnnounceEvent; fn get_id(&self) -> PeerId; fn get_updated(&self) -> DurationSinceUnixEpoch; - fn get_address(&self) -> SocketAddr; + fn get_address(&self) -> &PeerAddress; } impl ReadInfo for Peer { @@ -227,8 +370,8 @@ impl ReadInfo for Peer { self.updated } - fn get_address(&self) -> SocketAddr { - self.peer_addr + fn get_address(&self) -> &PeerAddress { + &self.peer_addr } } @@ -253,8 +396,8 @@ impl ReadInfo for Arc { self.updated } - fn get_address(&self) -> SocketAddr { - self.peer_addr + fn get_address(&self) -> &PeerAddress { + &self.peer_addr } } @@ -283,12 +426,14 @@ impl Peer { } } - pub fn ip(&mut self) -> IpAddr { + pub fn ip(&self) -> Option { self.peer_addr.ip() } pub fn change_ip(&mut self, new_ip: &IpAddr) { - self.peer_addr = SocketAddr::new(*new_ip, self.peer_addr.port()); + if let PeerAddress::Clearnet(address) = &mut self.peer_addr { + address.set_ip(*new_ip); + } } pub fn mark_as_completed(&mut self) { @@ -314,8 +459,6 @@ impl Peer { use std::panic::Location; -use thiserror::Error; - /// Error returned when trying to convert an invalid peer id from another type. /// /// Usually because the source format does not contain 20 bytes. @@ -518,7 +661,7 @@ pub mod fixture { pub fn seeder() -> Self { let peer = Peer { peer_id: PeerId(*b"-qB00000000000000001"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -534,7 +677,7 @@ pub mod fixture { pub fn leecher() -> Self { let peer = Peer { peer_id: PeerId(*b"-qB00000000000000002"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -555,13 +698,13 @@ pub mod fixture { #[allow(dead_code)] #[must_use] pub fn with_peer_addr(mut self, peer_addr: &SocketAddr) -> Self { - self.peer.peer_addr = *peer_addr; + self.peer.peer_addr = (*peer_addr).into(); self } #[must_use] pub fn with_peer_address(mut self, peer_addr: SocketAddr) -> Self { - self.peer.peer_addr = peer_addr; + self.peer.peer_addr = peer_addr.into(); self } @@ -622,7 +765,7 @@ pub mod fixture { fn default() -> Self { Self { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -642,6 +785,70 @@ pub mod fixture { #[cfg(test)] pub mod test { + mod i2p_destination { + use std::str::FromStr; + + use base64::Engine; + use base64::alphabet::Alphabet; + use base64::engine::{GeneralPurpose, GeneralPurposeConfig}; + + use crate::peer::{I2pDestination, ParseI2pDestinationError}; + + #[test] + fn it_should_parse_a_valid_i2p_base64_destination() { + let destination = "A".repeat(516); + + let parsed = I2pDestination::from_str(&destination).unwrap(); + + assert_eq!(parsed.to_string(), format!("{destination}.i2p")); + assert_eq!(parsed.hash().len(), 32); + } + + #[test] + fn it_should_reject_an_i2p_destination_with_invalid_base64() { + let destination = format!("{}.i2p", "!".repeat(516)); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!(error, ParseI2pDestinationError::InvalidBase64); + } + + #[test] + fn it_should_reject_an_i2p_destination_shorter_than_the_minimum_length() { + let destination = "A".repeat(512); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!(error, ParseI2pDestinationError::TooShort { actual: 384 }); + } + + #[test] + fn it_should_parse_a_long_padded_destination_when_the_certificate_length_matches() { + let certificate_payload_length = 91_u16; + let mut decoded = vec![0; 387 + usize::from(certificate_payload_length)]; + decoded[384] = 5; + decoded[385..387].copy_from_slice(&certificate_payload_length.to_be_bytes()); + let alphabet = Alphabet::new("ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-~").unwrap(); + let encoded = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()).encode(decoded); + + let parsed = I2pDestination::from_str(&encoded).unwrap(); + + assert!(encoded.ends_with("==")); + assert_eq!(parsed.to_string(), format!("{encoded}.i2p")); + } + + #[test] + fn it_should_reject_a_destination_when_the_certificate_length_does_not_match() { + let destination = "A".repeat(520); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!( + error, + ParseI2pDestinationError::InvalidCertificateLength { declared: 0, actual: 3 } + ); + } + } mod peer { use crate::peer::fixture::PeerBuilder; diff --git a/packages/rest-api-runtime-adapter/src/v1/conversion.rs b/packages/rest-api-runtime-adapter/src/v1/conversion.rs index 8eecb0ea9..78abfa8ea 100644 --- a/packages/rest-api-runtime-adapter/src/v1/conversion.rs +++ b/packages/rest-api-runtime-adapter/src/v1/conversion.rs @@ -80,7 +80,7 @@ mod tests { fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/swarm-coordination-registry/examples/bench_peers.rs b/packages/swarm-coordination-registry/examples/bench_peers.rs index 7235b1ac1..69125393f 100644 --- a/packages/swarm-coordination-registry/examples/bench_peers.rs +++ b/packages/swarm-coordination-registry/examples/bench_peers.rs @@ -16,7 +16,7 @@ fn make_peer(ip_last_octet: u8, port: u16, seed: u8) -> Peer { id[0] = ip_last_octet; Peer { peer_id: PeerId(id), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, ip_last_octet)), port), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, ip_last_octet)), port).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -46,7 +46,7 @@ fn bench_peers_excluding(num_peers: usize, limit: usize, iterations: u64) -> f64 rt.block_on(coordinator.handle_announcement(&peer)); } - let requesting_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 254)), 6999); + let requesting_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 254)), 6999).into(); // Warm up for _ in 0..1000 { diff --git a/packages/swarm-coordination-registry/src/lib.rs b/packages/swarm-coordination-registry/src/lib.rs index 992db6010..34cfcc5cd 100644 --- a/packages/swarm-coordination-registry/src/lib.rs +++ b/packages/swarm-coordination-registry/src/lib.rs @@ -68,7 +68,7 @@ pub(crate) mod tests { pub fn sample_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -81,7 +81,7 @@ pub(crate) mod tests { pub fn sample_peer_one() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000001"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -94,7 +94,7 @@ pub(crate) mod tests { pub fn sample_peer_two() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000002"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8082), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8082).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -120,7 +120,7 @@ pub(crate) mod tests { pub fn complete_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -134,7 +134,7 @@ pub(crate) mod tests { pub fn incomplete_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/swarm-coordination-registry/src/statistics/event/handler.rs b/packages/swarm-coordination-registry/src/statistics/event/handler.rs index 03952e137..bea313718 100644 --- a/packages/swarm-coordination-registry/src/statistics/event/handler.rs +++ b/packages/swarm-coordination-registry/src/statistics/event/handler.rs @@ -192,7 +192,7 @@ mod tests { // It returns a peer with the opposite role of the given peer. fn make_opposite_role_peer(peer: &Peer) -> Peer { - let mut opposite_role_peer = *peer; + let mut opposite_role_peer = peer.clone(); match peer.role() { PeerRole::Seeder => { @@ -491,7 +491,7 @@ mod tests { handle_event( Event::PeerAdded { info_hash: sample_info_hash(), - peer: old_peer, + peer: old_peer.clone(), }, &stats_repository, CurrentClock::now(), @@ -533,7 +533,7 @@ mod tests { handle_event( Event::PeerAdded { info_hash: sample_info_hash(), - peer, + peer: peer.clone(), }, &stats_repository, CurrentClock::now(), @@ -560,7 +560,7 @@ mod tests { handle_event( Event::PeerRemoved { info_hash: sample_info_hash(), - peer, + peer: peer.clone(), }, &stats_repository, CurrentClock::now(), @@ -588,7 +588,7 @@ mod tests { Event::PeerUpdated { info_hash: sample_info_hash(), old_peer: sample_peer(), - new_peer, + new_peer: new_peer.clone(), }, &stats_repository, CurrentClock::now(), diff --git a/packages/swarm-coordination-registry/src/swarm/coordinator.rs b/packages/swarm-coordination-registry/src/swarm/coordinator.rs index 562408af5..4c53d8e24 100644 --- a/packages/swarm-coordination-registry/src/swarm/coordinator.rs +++ b/packages/swarm-coordination-registry/src/swarm/coordinator.rs @@ -1,14 +1,13 @@ //! A swarm is a collection of peers that are all trying to download the same //! torrent. use std::collections::BTreeMap; -use std::net::SocketAddr; use std::sync::Arc; use torrust_clock::DurationSinceUnixEpoch; use torrust_info_hash::InfoHash; use torrust_tracker_primitives::peer::{self, Peer, PeerAnnouncement}; use torrust_tracker_primitives::swarm_metadata::SwarmMetadata; -use torrust_tracker_primitives::{AnnounceEvent, TrackerPolicy}; +use torrust_tracker_primitives::{AnnounceEvent, PeerAddress, TrackerPolicy}; use crate::event::Event; use crate::event::sender::Sender; @@ -16,7 +15,7 @@ use crate::event::sender::Sender; #[derive(Clone)] pub struct Coordinator { info_hash: InfoHash, - peers: BTreeMap>, + peers: BTreeMap>, metadata: SwarmMetadata, event_sender: Sender, } @@ -35,7 +34,7 @@ impl Coordinator { pub async fn handle_announcement(&mut self, incoming_announce: &PeerAnnouncement) { let _previous_peer = match peer::ReadInfo::get_event(incoming_announce) { AnnounceEvent::Started | AnnounceEvent::None | AnnounceEvent::Completed => { - self.upsert_peer(Arc::new(*incoming_announce)).await + self.upsert_peer(Arc::new(incoming_announce.clone())).await } AnnounceEvent::Stopped => self.remove_peer(&incoming_announce.peer_addr).await, }; @@ -52,7 +51,7 @@ impl Coordinator { } #[must_use] - pub fn get(&self, peer_addr: &SocketAddr) -> Option<&Arc> { + pub fn get(&self, peer_addr: &PeerAddress) -> Option<&Arc> { self.peers.get(peer_addr) } @@ -65,24 +64,16 @@ impl Coordinator { } #[must_use] - pub fn peers_excluding(&self, peer_addr: &SocketAddr, limit: Option) -> Vec> { + pub fn peers_excluding(&self, peer_addr: &PeerAddress, limit: Option) -> Vec> { + let peers = self + .peers + .values() + .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != peer_addr) + .filter(|peer| peer.peer_addr.is_i2p() == peer_addr.is_i2p()); + match limit { - Some(limit) => self - .peers - .values() - // Take peers which are not the client peer - .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != *peer_addr) - // Limit the number of peers on the result - .take(limit) - .cloned() - .collect(), - None => self - .peers - .values() - // Take peers which are not the client peer - .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != *peer_addr) - .cloned() - .collect(), + Some(limit) => peers.take(limit).cloned().collect(), + None => peers.cloned().collect(), } } @@ -157,7 +148,7 @@ impl Coordinator { async fn upsert_peer(&mut self, incoming_announce: Arc) -> Option> { let announcement = incoming_announce.clone(); - if let Some(previous_announce) = self.peers.insert(incoming_announce.peer_addr, incoming_announce) { + if let Some(previous_announce) = self.peers.insert(incoming_announce.peer_addr.clone(), incoming_announce) { let downloads_increased = self.update_metadata_on_update(&previous_announce, &announcement); self.trigger_peer_updated_event(&previous_announce, &announcement).await; @@ -176,7 +167,7 @@ impl Coordinator { } } - async fn remove_peer(&mut self, peer_addr: &SocketAddr) -> Option> { + async fn remove_peer(&mut self, peer_addr: &PeerAddress) -> Option> { if let Some(old_peer) = self.peers.remove(peer_addr) { self.update_metadata_on_removal(&old_peer); @@ -189,11 +180,11 @@ impl Coordinator { } #[must_use] - fn inactive_peers(&self, current_cutoff: DurationSinceUnixEpoch) -> Vec { + fn inactive_peers(&self, current_cutoff: DurationSinceUnixEpoch) -> Vec { self.peers .iter() .filter(|(_, peer)| peer::ReadInfo::get_updated(&**peer) <= current_cutoff) - .map(|(addr, _)| *addr) + .map(|(addr, _)| addr.clone()) .collect() } @@ -249,7 +240,7 @@ impl Coordinator { event_sender .send(Event::PeerAdded { info_hash: self.info_hash, - peer: *announcement.clone(), + peer: announcement.as_ref().clone(), }) .await; } @@ -260,7 +251,7 @@ impl Coordinator { event_sender .send(Event::PeerRemoved { info_hash: self.info_hash, - peer: *old_peer.clone(), + peer: old_peer.as_ref().clone(), }) .await; } @@ -271,8 +262,8 @@ impl Coordinator { event_sender .send(Event::PeerUpdated { info_hash: self.info_hash, - old_peer: *old_announce.clone(), - new_peer: *new_announce.clone(), + old_peer: old_announce.as_ref().clone(), + new_peer: new_announce.as_ref().clone(), }) .await; } @@ -283,7 +274,7 @@ impl Coordinator { event_sender .send(Event::PeerDownloadCompleted { info_hash: self.info_hash, - peer: *new_announce.clone(), + peer: new_announce.as_ref().clone(), }) .await; } @@ -321,9 +312,9 @@ mod tests { use std::sync::Arc; use torrust_clock::DurationSinceUnixEpoch; - use torrust_tracker_primitives::PeerId; use torrust_tracker_primitives::peer::fixture::PeerBuilder; use torrust_tracker_primitives::swarm_metadata::SwarmMetadata; + use torrust_tracker_primitives::{I2pDestination, I2pPeerAddress, PeerAddress, PeerId}; use crate::swarm::coordinator::Coordinator; use crate::tests::sample_info_hash; @@ -348,7 +339,7 @@ mod tests { let peer = PeerBuilder::default().build(); - assert_eq!(swarm.upsert_peer(peer.into()).await, None); + assert_eq!(swarm.upsert_peer(peer.clone().into()).await, None); } #[tokio::test] @@ -357,9 +348,9 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; - assert_eq!(swarm.upsert_peer(peer.into()).await, Some(Arc::new(peer))); + assert_eq!(swarm.upsert_peer(peer.clone().into()).await, Some(Arc::new(peer))); } #[tokio::test] @@ -368,7 +359,7 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.peers(None), [Arc::new(peer)]); } @@ -379,7 +370,7 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.get(&peer.peer_addr), Some(Arc::new(peer)).as_ref()); } @@ -390,7 +381,7 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.len(), 1); } @@ -401,7 +392,7 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; swarm.remove_peer(&peer.peer_addr).await; @@ -414,11 +405,11 @@ mod tests { let peer = PeerBuilder::default().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let old = swarm.remove_peer(&peer.peer_addr).await; - assert_eq!(old, Some(Arc::new(peer))); + assert_eq!(old, Some(Arc::new(peer.clone()))); assert_eq!(swarm.get(&peer.peer_addr), None); } @@ -439,17 +430,34 @@ mod tests { .with_peer_id(&PeerId(*b"-qB00000000000000001")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6969)) .build(); - swarm.upsert_peer(peer1.into()).await; + swarm.upsert_peer(peer1.clone().into()).await; let peer2 = PeerBuilder::default() .with_peer_id(&PeerId(*b"-qB00000000000000002")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 6969)) .build(); - swarm.upsert_peer(peer2.into()).await; + swarm.upsert_peer(peer2.clone().into()).await; assert_eq!(swarm.peers_excluding(&peer2.peer_addr, None), [Arc::new(peer1)]); } + #[tokio::test] + async fn it_should_not_return_clearnet_peers_to_an_i2p_peer() { + let mut swarm = Coordinator::new(&sample_info_hash(), 0, None); + let clearnet_peer = PeerBuilder::default().build(); + swarm.upsert_peer(clearnet_peer.clone().into()).await; + + let mut i2p_peer = PeerBuilder::default().with_peer_id(&PeerId(*b"-qB00000000000000002")).build(); + i2p_peer.peer_addr = PeerAddress::I2p(I2pPeerAddress { + destination: format!("{}.i2p", "A".repeat(516)).parse::().unwrap(), + }); + swarm.upsert_peer(i2p_peer.clone().into()).await; + + let peers = swarm.peers_excluding(&i2p_peer.peer_addr, None); + + assert!(peers.is_empty()); + } + #[tokio::test] async fn it_should_count_inactive_peers() { let mut swarm = Coordinator::new(&sample_info_hash(), 0, None); @@ -459,7 +467,7 @@ mod tests { // Insert the peer let last_update_time = DurationSinceUnixEpoch::new(1_669_397_478_934, 0); let peer = PeerBuilder::default().last_updated_on(last_update_time).build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let inactive_peers_total = swarm.count_inactive_peers(last_update_time + one_second); @@ -475,7 +483,7 @@ mod tests { // Insert the peer let last_update_time = DurationSinceUnixEpoch::new(1_669_397_478_934, 0); let peer = PeerBuilder::default().last_updated_on(last_update_time).build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; // Remove peers not updated since one second after inserting the peer swarm.remove_inactive(last_update_time + one_second).await; @@ -492,7 +500,7 @@ mod tests { // Insert the peer let last_update_time = DurationSinceUnixEpoch::new(1_669_397_478_934, 0); let peer = PeerBuilder::default().last_updated_on(last_update_time).build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; // Remove peers not updated since one second before inserting the peer. swarm.remove_inactive(last_update_time.checked_sub(one_second).unwrap()).await; @@ -523,11 +531,11 @@ mod tests { let mut peer = PeerBuilder::leecher().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; peer.event = torrust_tracker_primitives::AnnounceEvent::Completed; - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert!(swarm.metadata().downloads() > 0); @@ -608,12 +616,12 @@ mod tests { let peer1 = PeerBuilder::default() .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6969)) .build(); - swarm.upsert_peer(peer1.into()).await; + swarm.upsert_peer(peer1.clone().into()).await; let peer2 = PeerBuilder::default() .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 6969)) .build(); - swarm.upsert_peer(peer2.into()).await; + swarm.upsert_peer(peer2.clone().into()).await; assert_eq!(swarm.len(), 2); } @@ -629,13 +637,13 @@ mod tests { .with_peer_id(&PeerId(*b"-qB00000000000000001")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6969)) .build(); - swarm.upsert_peer(peer1.into()).await; + swarm.upsert_peer(peer1.clone().into()).await; let peer2 = PeerBuilder::default() .with_peer_id(&PeerId(*b"-qB00000000000000002")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6969)) .build(); - swarm.upsert_peer(peer2.into()).await; + swarm.upsert_peer(peer2.clone().into()).await; assert_eq!(swarm.len(), 1); } @@ -647,8 +655,8 @@ mod tests { let seeder = PeerBuilder::seeder().build(); let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(seeder.into()).await; - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; + swarm.upsert_peer(leecher.clone().into()).await; assert_eq!( swarm.metadata(), @@ -667,8 +675,8 @@ mod tests { let seeder = PeerBuilder::seeder().build(); let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(seeder.into()).await; - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; + swarm.upsert_peer(leecher.clone().into()).await; let (seeders, _leechers) = swarm.seeders_and_leechers(); @@ -682,8 +690,8 @@ mod tests { let seeder = PeerBuilder::seeder().build(); let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(seeder.into()).await; - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; + swarm.upsert_peer(leecher.clone().into()).await; let (_seeders, leechers) = swarm.seeders_and_leechers(); @@ -712,7 +720,7 @@ mod tests { let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(leecher.clone().into()).await; assert_eq!(swarm.metadata().leechers(), leechers + 1); } @@ -725,7 +733,7 @@ mod tests { let seeder = PeerBuilder::seeder().build(); - swarm.upsert_peer(seeder.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; assert_eq!(swarm.metadata().seeders(), seeders + 1); } @@ -739,7 +747,7 @@ mod tests { let seeder = PeerBuilder::seeder().build(); - swarm.upsert_peer(seeder.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; assert_eq!(swarm.metadata().downloads(), downloads); } @@ -757,7 +765,7 @@ mod tests { let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(leecher.clone().into()).await; let leechers = swarm.metadata().leechers(); @@ -772,7 +780,7 @@ mod tests { let seeder = PeerBuilder::seeder().build(); - swarm.upsert_peer(seeder.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; let seeders = swarm.metadata().seeders(); @@ -796,7 +804,7 @@ mod tests { let leecher = PeerBuilder::leecher().build(); - swarm.upsert_peer(leecher.into()).await; + swarm.upsert_peer(leecher.clone().into()).await; let leechers = swarm.metadata().leechers(); @@ -811,7 +819,7 @@ mod tests { let seeder = PeerBuilder::seeder().build(); - swarm.upsert_peer(seeder.into()).await; + swarm.upsert_peer(seeder.clone().into()).await; let seeders = swarm.metadata().seeders(); @@ -834,14 +842,14 @@ mod tests { let mut peer = PeerBuilder::leecher().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let leechers = swarm.metadata().leechers(); let seeders = swarm.metadata().seeders(); peer.left = NumberOfBytes::new(0); // Convert to seeder - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.metadata().seeders(), seeders + 1); assert_eq!(swarm.metadata().leechers(), leechers - 1); @@ -853,14 +861,14 @@ mod tests { let mut peer = PeerBuilder::seeder().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let leechers = swarm.metadata().leechers(); let seeders = swarm.metadata().seeders(); peer.left = NumberOfBytes::new(10); // Convert to leecher - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.metadata().leechers(), leechers + 1); assert_eq!(swarm.metadata().seeders(), seeders - 1); @@ -872,13 +880,13 @@ mod tests { let mut peer = PeerBuilder::leecher().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let downloads = swarm.metadata().downloads(); peer.event = torrust_tracker_primitives::AnnounceEvent::Completed; - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.metadata().downloads(), downloads + 1); } @@ -889,15 +897,15 @@ mod tests { let mut peer = PeerBuilder::leecher().build(); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; let downloads = swarm.metadata().downloads(); peer.event = torrust_tracker_primitives::AnnounceEvent::Completed; - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; assert_eq!(swarm.metadata().downloads(), downloads + 1); } @@ -924,11 +932,17 @@ mod tests { let mut event_sender_mock = MockEventSender::new(); - expect_event_sequence(&mut event_sender_mock, vec![Event::PeerAdded { info_hash, peer }]); + expect_event_sequence( + &mut event_sender_mock, + vec![Event::PeerAdded { + info_hash, + peer: peer.clone(), + }], + ); let mut swarm = Coordinator::new(&sample_info_hash(), 0, Some(Arc::new(event_sender_mock))); - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; } #[tokio::test] @@ -940,13 +954,22 @@ mod tests { expect_event_sequence( &mut event_sender_mock, - vec![Event::PeerAdded { info_hash, peer }, Event::PeerRemoved { info_hash, peer }], + vec![ + Event::PeerAdded { + info_hash, + peer: peer.clone(), + }, + Event::PeerRemoved { + info_hash, + peer: peer.clone(), + }, + ], ); let mut swarm = Coordinator::new(&info_hash, 0, Some(Arc::new(event_sender_mock))); // Insert the peer - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; swarm.remove_peer(&peer.peer_addr).await; } @@ -960,13 +983,22 @@ mod tests { expect_event_sequence( &mut event_sender_mock, - vec![Event::PeerAdded { info_hash, peer }, Event::PeerRemoved { info_hash, peer }], + vec![ + Event::PeerAdded { + info_hash, + peer: peer.clone(), + }, + Event::PeerRemoved { + info_hash, + peer: peer.clone(), + }, + ], ); let mut swarm = Coordinator::new(&info_hash, 0, Some(Arc::new(event_sender_mock))); // Insert the peer - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; // Peers not updated after this time will be removed let current_cutoff = peer.updated + DurationSinceUnixEpoch::from_secs(1); @@ -984,11 +1016,14 @@ mod tests { expect_event_sequence( &mut event_sender_mock, vec![ - Event::PeerAdded { info_hash, peer }, + Event::PeerAdded { + info_hash, + peer: peer.clone(), + }, Event::PeerUpdated { info_hash, - old_peer: peer, - new_peer: peer, + old_peer: peer.clone(), + new_peer: peer.clone(), }, ], ); @@ -996,17 +1031,17 @@ mod tests { let mut swarm = Coordinator::new(&info_hash, 0, Some(Arc::new(event_sender_mock))); // Insert the peer - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; // Update the peer - swarm.upsert_peer(peer.into()).await; + swarm.upsert_peer(peer.clone().into()).await; } #[tokio::test] async fn it_should_trigger_an_event_when_a_peer_completes_a_download() { let info_hash = sample_info_hash(); let started_peer = PeerBuilder::leecher().with_event(Started).build(); - let completed_peer = started_peer.into_completed(); + let completed_peer = started_peer.clone().into_completed(); let mut event_sender_mock = MockEventSender::new(); @@ -1015,16 +1050,16 @@ mod tests { vec![ Event::PeerAdded { info_hash, - peer: started_peer, + peer: started_peer.clone(), }, Event::PeerUpdated { info_hash, - old_peer: started_peer, - new_peer: completed_peer, + old_peer: started_peer.clone(), + new_peer: completed_peer.clone(), }, Event::PeerDownloadCompleted { info_hash, - peer: completed_peer, + peer: completed_peer.clone(), }, ], ); @@ -1032,10 +1067,10 @@ mod tests { let mut swarm = Coordinator::new(&info_hash, 0, Some(Arc::new(event_sender_mock))); // Insert the peer - swarm.upsert_peer(started_peer.into()).await; + swarm.upsert_peer(started_peer.clone().into()).await; // Announce as completed - swarm.upsert_peer(completed_peer.into()).await; + swarm.upsert_peer(completed_peer.clone().into()).await; } } } diff --git a/packages/swarm-coordination-registry/src/swarm/registry.rs b/packages/swarm-coordination-registry/src/swarm/registry.rs index cbac4b826..686c03dae 100644 --- a/packages/swarm-coordination-registry/src/swarm/registry.rs +++ b/packages/swarm-coordination-registry/src/swarm/registry.rs @@ -68,7 +68,7 @@ impl Registry { event_sender .send(Event::TorrentAdded { info_hash: *info_hash, - announcement: *peer, + announcement: peer.clone(), }) .await; } @@ -652,7 +652,7 @@ mod tests { for idx in 1..=75 { let peer = Peer { peer_id: numeric_peer_id(idx), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, idx.try_into().unwrap())), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, idx.try_into().unwrap())), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -723,7 +723,8 @@ mod tests { for idx in 2..=75 { let peer = Peer { peer_id: numeric_peer_id(idx), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, idx.try_into().unwrap())), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, idx.try_into().unwrap())), 8080) + .into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -874,7 +875,7 @@ mod tests { fn into(self) -> TorrentEntryInfo { TorrentEntryInfo { swarm_metadata: self.metadata(), - peers: self.peers(None).iter().map(|peer| *peer.clone()).collect(), + peers: self.peers(None).iter().map(|peer| peer.as_ref().clone()).collect(), number_of_peers: self.len(), } } @@ -1366,9 +1367,12 @@ mod tests { vec![ Event::TorrentAdded { info_hash, - announcement: peer, + announcement: peer.clone(), + }, + Event::PeerAdded { + info_hash, + peer: peer.clone(), }, - Event::PeerAdded { info_hash, peer }, ], ); @@ -1389,9 +1393,12 @@ mod tests { vec![ Event::TorrentAdded { info_hash, - announcement: peer, + announcement: peer.clone(), + }, + Event::PeerAdded { + info_hash, + peer: peer.clone(), }, - Event::PeerAdded { info_hash, peer }, Event::TorrentRemoved { info_hash }, ], ); @@ -1415,10 +1422,16 @@ mod tests { vec![ Event::TorrentAdded { info_hash, - announcement: peer, + announcement: peer.clone(), + }, + Event::PeerAdded { + info_hash, + peer: peer.clone(), + }, + Event::PeerRemoved { + info_hash, + peer: peer.clone(), }, - Event::PeerAdded { info_hash, peer }, - Event::PeerRemoved { info_hash, peer }, Event::TorrentRemoved { info_hash }, ], ); diff --git a/packages/torrent-repository-benchmarking/benches/helpers/utils.rs b/packages/torrent-repository-benchmarking/benches/helpers/utils.rs index 99dd439cd..d477c4ae4 100644 --- a/packages/torrent-repository-benchmarking/benches/helpers/utils.rs +++ b/packages/torrent-repository-benchmarking/benches/helpers/utils.rs @@ -8,7 +8,7 @@ use torrust_tracker_primitives::{AnnounceEvent, NumberOfBytes, PeerId}; pub const DEFAULT_PEER: Peer = Peer { peer_id: PeerId([0; 20]), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080), + peer_addr: torrust_tracker_primitives::PeerAddress::Clearnet(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8080)), updated: DurationSinceUnixEpoch::from_secs(0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/torrent-repository-benchmarking/src/entry/peer_list.rs b/packages/torrent-repository-benchmarking/src/entry/peer_list.rs index aac071c6a..a1df7748e 100644 --- a/packages/torrent-repository-benchmarking/src/entry/peer_list.rs +++ b/packages/torrent-repository-benchmarking/src/entry/peer_list.rs @@ -1,9 +1,8 @@ //! A peer list. -use std::net::SocketAddr; use std::sync::Arc; use torrust_clock::DurationSinceUnixEpoch; -use torrust_tracker_primitives::{PeerId, peer}; +use torrust_tracker_primitives::{PeerAddress, PeerId, peer}; // code-review: the current implementation uses the peer Id as the ``BTreeMap`` // key. That would allow adding two identical peers except for the Id. @@ -61,13 +60,13 @@ impl PeerList { } #[must_use] - pub fn get_peers_excluding_addr(&self, peer_addr: &SocketAddr, limit: Option) -> Vec> { + pub fn get_peers_excluding_addr(&self, peer_addr: &PeerAddress, limit: Option) -> Vec> { limit.map_or_else( || { self.peers .values() // Take peers which are not the client peer - .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != *peer_addr) + .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != peer_addr) .cloned() .collect() }, @@ -75,7 +74,7 @@ impl PeerList { self.peers .values() // Take peers which are not the client peer - .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != *peer_addr) + .filter(|peer| peer::ReadInfo::get_address(peer.as_ref()) != peer_addr) // Limit the number of peers on the result .take(limit) .cloned() @@ -127,9 +126,9 @@ mod tests { let peer = PeerBuilder::default().build(); - peer_list.upsert(peer.into()); + peer_list.upsert(peer.clone().into()); - assert_eq!(peer_list.upsert(peer.into()), Some(Arc::new(peer))); + assert_eq!(peer_list.upsert(peer.clone().into()), Some(Arc::new(peer))); } #[test] @@ -138,7 +137,7 @@ mod tests { let peer = PeerBuilder::default().build(); - peer_list.upsert(peer.into()); + peer_list.upsert(peer.clone().into()); assert_eq!(peer_list.get_all(None), [Arc::new(peer)]); } @@ -149,7 +148,7 @@ mod tests { let peer = PeerBuilder::default().build(); - peer_list.upsert(peer.into()); + peer_list.upsert(peer.clone().into()); assert_eq!(peer_list.get(&peer.peer_id), Some(Arc::new(peer)).as_ref()); } @@ -171,7 +170,7 @@ mod tests { let peer = PeerBuilder::default().build(); - peer_list.upsert(peer.into()); + peer_list.upsert(peer.clone().into()); peer_list.remove(&peer.peer_id); @@ -184,7 +183,7 @@ mod tests { let peer = PeerBuilder::default().build(); - peer_list.upsert(peer.into()); + peer_list.upsert(peer.clone().into()); peer_list.remove(&peer.peer_id); @@ -199,13 +198,13 @@ mod tests { .with_peer_id(&PeerId(*b"-qB00000000000000001")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 6969)) .build(); - peer_list.upsert(peer1.into()); + peer_list.upsert(peer1.clone().into()); let peer2 = PeerBuilder::default() .with_peer_id(&PeerId(*b"-qB00000000000000002")) .with_peer_addr(&SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 2)), 6969)) .build(); - peer_list.upsert(peer2.into()); + peer_list.upsert(peer2.clone().into()); assert_eq!(peer_list.get_peers_excluding_addr(&peer2.peer_addr, None), [Arc::new(peer1)]); } diff --git a/packages/torrent-repository-benchmarking/src/entry/single.rs b/packages/torrent-repository-benchmarking/src/entry/single.rs index 8d949698c..dba6cf606 100644 --- a/packages/torrent-repository-benchmarking/src/entry/single.rs +++ b/packages/torrent-repository-benchmarking/src/entry/single.rs @@ -4,7 +4,7 @@ use std::sync::Arc; use torrust_clock::DurationSinceUnixEpoch; use torrust_tracker_primitives::peer::{self}; use torrust_tracker_primitives::swarm_metadata::SwarmMetadata; -use torrust_tracker_primitives::{AnnounceEvent, TrackerPolicy}; +use torrust_tracker_primitives::{AnnounceEvent, PeerAddress, TrackerPolicy}; use super::Entry; use crate::EntrySingle; @@ -46,7 +46,7 @@ impl Entry for EntrySingle { } fn get_peers_for_client(&self, client: &SocketAddr, limit: Option) -> Vec> { - self.swarm.get_peers_excluding_addr(client, limit) + self.swarm.get_peers_excluding_addr(&PeerAddress::Clearnet(*client), limit) } fn upsert_peer(&mut self, peer: &peer::Peer) -> bool { @@ -57,7 +57,7 @@ impl Entry for EntrySingle { drop(self.swarm.remove(&peer::ReadInfo::get_id(peer))); } AnnounceEvent::Completed => { - let previous = self.swarm.upsert(Arc::new(*peer)); + let previous = self.swarm.upsert(Arc::new(peer.clone())); // Don't count if peer was not previously known and not already completed. if previous.is_some_and(|p| p.event != AnnounceEvent::Completed) { self.downloaded += 1; @@ -67,7 +67,7 @@ impl Entry for EntrySingle { _ => { // `Started` event (first announced event) or // `None` event (announcements done at regular intervals). - drop(self.swarm.upsert(Arc::new(*peer))); + drop(self.swarm.upsert(Arc::new(peer.clone()))); } } diff --git a/packages/torrent-repository-benchmarking/tests/entry/mod.rs b/packages/torrent-repository-benchmarking/tests/entry/mod.rs index e06ad358b..a7871b235 100644 --- a/packages/torrent-repository-benchmarking/tests/entry/mod.rs +++ b/packages/torrent-repository-benchmarking/tests/entry/mod.rs @@ -271,7 +271,7 @@ async fn it_should_handle_a_peer_completed_announcement_and_update_the_downloade let downloaded = torrent.get_stats().await.downloaded; let peers = torrent.get_peers(None).await; - let mut peer = **peers.first().expect("there should be a peer"); + let mut peer = peers.first().expect("there should be a peer").as_ref().clone(); let is_already_completed = peer.event == AnnounceEvent::Completed; @@ -302,7 +302,7 @@ async fn it_should_update_a_peer_as_a_seeder( let completed = u32::try_from(peers.iter().filter(|p| p.is_seeder()).count()).expect("it_should_not_be_so_many"); let peers = torrent.get_peers(None).await; - let mut peer = **peers.first().expect("there should be a peer"); + let mut peer = peers.first().expect("there should be a peer").as_ref().clone(); let is_already_non_left = peer.left == NumberOfBytes::new(0); @@ -334,7 +334,7 @@ async fn it_should_update_a_peer_as_incomplete( let incomplete = u32::try_from(peers.iter().filter(|p| !p.is_seeder()).count()).expect("it should not be so many"); let peers = torrent.get_peers(None).await; - let mut peer = **peers.first().expect("there should be a peer"); + let mut peer = peers.first().expect("there should be a peer").as_ref().clone(); let completed_already = peer.left == NumberOfBytes::new(0); @@ -365,18 +365,23 @@ async fn it_should_get_peers_excluding_the_client_socket( make(&mut torrent, makes).await; let peers = torrent.get_peers(None).await; - let mut peer = **peers.first().expect("there should be a peer"); + let mut peer = peers.first().expect("there should be a peer").as_ref().clone(); let socket = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 8081); // for this test, we should not already use this socket. - assert_ne!(peer.peer_addr, socket); + assert_ne!(peer.peer_addr, socket.into()); // it should get the peer as it dose not share the socket. - assert!(torrent.get_peers_for_client(&socket, None).await.contains(&peer.into())); + assert!( + torrent + .get_peers_for_client(&socket, None) + .await + .contains(&peer.clone().into()) + ); // set the address to the socket. - peer.peer_addr = socket; + peer.peer_addr = socket.into(); torrent.upsert_peer(&peer).await; // Add peer // It should not include the peer that has the same socket. diff --git a/packages/torrent-repository-benchmarking/tests/repository/mod.rs b/packages/torrent-repository-benchmarking/tests/repository/mod.rs index a8469413a..1d0df2a31 100644 --- a/packages/torrent-repository-benchmarking/tests/repository/mod.rs +++ b/packages/torrent-repository-benchmarking/tests/repository/mod.rs @@ -583,7 +583,7 @@ async fn it_should_remove_inactive_peers( // Verify that this new peer was inserted into the repository. { let entry = repo.get(&info_hash).await.expect("it_should_get_some"); - assert!(entry.get_peers(None).contains(&peer.into())); + assert!(entry.get_peers(None).contains(&peer.clone().into())); } // Remove peers that have not been updated since the timeout (120 seconds ago). diff --git a/packages/tracker-core/src/announce_handler.rs b/packages/tracker-core/src/announce_handler.rs index b4339f692..304522ee1 100644 --- a/packages/tracker-core/src/announce_handler.rs +++ b/packages/tracker-core/src/announce_handler.rs @@ -27,7 +27,7 @@ //! //! let peer = peer::Peer { //! peer_id: PeerId(*b"-qB00000000000000001"), -//! peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081), +//! peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081).into(), //! updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), //! uploaded: NumberOfBytes::new(0), //! downloaded: NumberOfBytes::new(0), @@ -163,10 +163,12 @@ impl AnnounceHandler { ) -> Result { self.whitelist_authorization.authorize(info_hash).await?; - peer.change_ip(&assign_ip_address_to_peer( - remote_client_ip, - self.config.net.external_ip.map(Into::into), - )); + if !peer.peer_addr.is_i2p() { + peer.change_ip(&assign_ip_address_to_peer( + remote_client_ip, + self.config.net.external_ip.map(Into::into), + )); + } self.in_memory_torrent_repository .handle_announcement(info_hash, peer, self.load_downloads_metric_if_needed(info_hash).await?) @@ -314,7 +316,7 @@ mod tests { fn sample_peer_1() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000001"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8081).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -327,7 +329,7 @@ mod tests { fn sample_peer_2() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000002"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8082), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8082).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -340,7 +342,7 @@ mod tests { fn sample_peer_3() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000003"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 3)), 8082), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 3)), 8082).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/tracker-core/src/peer_tests.rs b/packages/tracker-core/src/peer_tests.rs index 6dcf08f14..82721c3bd 100644 --- a/packages/tracker-core/src/peer_tests.rs +++ b/packages/tracker-core/src/peer_tests.rs @@ -14,7 +14,7 @@ fn it_should_be_serializable() { let torrent_peer = peer::Peer { peer_id: PeerId(*b"-qB0000-000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: CurrentClock::now(), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/tracker-core/src/test_helpers.rs b/packages/tracker-core/src/test_helpers.rs index 4607eb205..459ebb811 100644 --- a/packages/tracker-core/src/test_helpers.rs +++ b/packages/tracker-core/src/test_helpers.rs @@ -69,7 +69,7 @@ pub(crate) mod tests { pub fn sample_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -105,7 +105,7 @@ pub(crate) mod tests { pub fn complete_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000001"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), @@ -119,7 +119,7 @@ pub(crate) mod tests { pub fn incomplete_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000002"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 2)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/tracker-core/src/torrent/services.rs b/packages/tracker-core/src/torrent/services.rs index de539c6d8..bb028189c 100644 --- a/packages/tracker-core/src/torrent/services.rs +++ b/packages/tracker-core/src/torrent/services.rs @@ -105,7 +105,7 @@ pub async fn get_torrent_info( let peers = torrent_entry.lock().await.peers(None); - let peers = Some(peers.iter().map(|peer| **peer).collect()); + let peers = Some(peers.iter().map(|peer| peer.as_ref().clone()).collect()); Some(Info { info_hash: *info_hash, @@ -212,7 +212,7 @@ mod tests { fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/tracker-core/tests/common/fixtures.rs b/packages/tracker-core/tests/common/fixtures.rs index 6e3d2680b..dfdcc03ba 100644 --- a/packages/tracker-core/tests/common/fixtures.rs +++ b/packages/tracker-core/tests/common/fixtures.rs @@ -36,7 +36,7 @@ pub fn sample_info_hash() -> InfoHash { pub fn sample_peer() -> Peer { Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(remote_client_ip(), 8080), + peer_addr: SocketAddr::new(remote_client_ip(), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), diff --git a/packages/tracker-core/tests/common/test_env.rs b/packages/tracker-core/tests/common/test_env.rs index 855fa0abb..bd6779f18 100644 --- a/packages/tracker-core/tests/common/test_env.rs +++ b/packages/tracker-core/tests/common/test_env.rs @@ -142,7 +142,7 @@ impl TestEnv { } pub async fn increase_number_of_downloads(&mut self, peer: Peer, remote_client_ip: &IpAddr, info_hash: &InfoHash) { - let _announce_data = self.announce_peer_started(peer, remote_client_ip, info_hash).await; + let _announce_data = self.announce_peer_started(peer.clone(), remote_client_ip, info_hash).await; let announce_data = self.announce_peer_completed(peer, remote_client_ip, info_hash).await; assert_eq!(announce_data.stats.downloads(), 1); diff --git a/packages/udp-core/src/peer_builder.rs b/packages/udp-core/src/peer_builder.rs index 5bef7d48e..8cff393f7 100644 --- a/packages/udp-core/src/peer_builder.rs +++ b/packages/udp-core/src/peer_builder.rs @@ -18,7 +18,7 @@ pub fn from_request(announce_request: &torrust_tracker_udp_protocol::AnnounceReq peer::Peer { peer_id: torrust_tracker_primitives::PeerId(announce_request.peer_id.0), - peer_addr: SocketAddr::new(*peer_ip, announce_request.port.0.into()), + peer_addr: SocketAddr::new(*peer_ip, announce_request.port.0.into()).into(), updated: CurrentClock::now(), uploaded: torrust_tracker_primitives::NumberOfBytes::new(announce_request.bytes_uploaded.0.get()), downloaded: torrust_tracker_primitives::NumberOfBytes::new(announce_request.bytes_downloaded.0.get()), diff --git a/packages/udp-server/src/handlers/announce.rs b/packages/udp-server/src/handlers/announce.rs index 5cc6bf459..0ab061350 100644 --- a/packages/udp-server/src/handlers/announce.rs +++ b/packages/udp-server/src/handlers/announce.rs @@ -132,7 +132,7 @@ fn build_response( .peers .iter() .filter_map(|peer| { - if let IpAddr::V4(ip) = peer.peer_addr.ip() { + if let Some(IpAddr::V4(ip)) = peer.peer_addr.ip() { Some(ResponsePeer:: { ip_address: ip.into(), port: Port(peer.peer_addr.port().into()), @@ -157,7 +157,7 @@ fn build_response( .peers .iter() .filter_map(|peer| { - if let IpAddr::V6(ip) = peer.peer_addr.ip() { + if let Some(IpAddr::V6(ip)) = peer.peer_addr.ip() { Some(ResponsePeer:: { ip_address: ip.into(), port: Port(peer.peer_addr.port().into()), @@ -413,7 +413,10 @@ pub(crate) mod tests { .get_torrent_peers(&info_hash.0.into(), usize::MAX) .await; - assert_eq!(peers[0].peer_addr, SocketAddr::new(IpAddr::V4(remote_client_ip), client_port)); + assert_eq!( + peers[0].peer_addr, + SocketAddr::new(IpAddr::V4(remote_client_ip), client_port).into() + ); } async fn add_a_torrent_peer_using_ipv6(in_memory_torrent_repository: &Arc) { @@ -770,7 +773,10 @@ pub(crate) mod tests { .await; // When using IPv6 the tracker converts the remote client ip into a IPv4 address - assert_eq!(peers[0].peer_addr, SocketAddr::new(IpAddr::V6(remote_client_ip), client_port)); + assert_eq!( + peers[0].peer_addr, + SocketAddr::new(IpAddr::V6(remote_client_ip), client_port).into() + ); } async fn add_a_torrent_peer_using_ipv4(in_memory_torrent_repository: &Arc) { @@ -943,7 +949,8 @@ pub(crate) mod tests { let peer_id = PeerId([255u8; 20]); let mut announcement = sample_peer(); announcement.peer_id = torrust_tracker_primitives::PeerId(peer_id.0); - announcement.peer_addr = SocketAddr::new(IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0, 0x7e00, 1)), client_port); + announcement.peer_addr = + SocketAddr::new(IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0, 0x7e00, 1)), client_port).into(); let client_socket_addr = SocketAddr::new(IpAddr::V6(client_ip_v6), client_port); let mut server_socket_addr = config.udp_trackers.clone().unwrap()[0].bind_address; @@ -979,7 +986,7 @@ pub(crate) mod tests { server_service_binding.clone(), ), info_hash: torrust_info_hash::InfoHash::from(info_hash.0), - announcement, + announcement: announcement.clone(), }; announce_events_match(event, &expected_event) @@ -1043,7 +1050,7 @@ pub(crate) mod tests { // 1111:2222:3333:4444:5555:6666:1.2.3.4 // // ::127.0.0.1 is the IPV6 representation for the IPV4 address 127.0.0.1. - assert_eq!(Ok(peers[0].peer_addr.ip()), "::126.0.0.1".parse()); + assert_eq!(peers[0].peer_addr.ip(), "::126.0.0.1".parse::().ok()); } } } diff --git a/packages/udp-server/src/lib.rs b/packages/udp-server/src/lib.rs index 75a54e25a..2d8d8a27e 100644 --- a/packages/udp-server/src/lib.rs +++ b/packages/udp-server/src/lib.rs @@ -683,7 +683,7 @@ pub(crate) mod tests { pub fn sample_peer() -> peer::Peer { peer::Peer { peer_id: PeerId(*b"-qB00000000000000000"), - peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080), + peer_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(126, 0, 0, 1)), 8080).into(), updated: DurationSinceUnixEpoch::new(1_669_397_478_934, 0), uploaded: NumberOfBytes::new(0), downloaded: NumberOfBytes::new(0), From ca3006f5b8e7e17e67bd10447d6edd946d624e26 Mon Sep 17 00:00:00 2001 From: Frigyes Erdosi Szucs Date: Fri, 31 Jul 2026 12:45:40 -0700 Subject: [PATCH 2/2] feat(http-tracker): support I2P peers --- README.md | 2 + cspell.json | 3 +- docs/packages.md | 7 + .../src/v1/handlers/announce.rs | 9 +- .../receiving_an_announce_request.rs | 6 +- .../server/v1/contract/context/torrent.rs | 2 +- packages/http-core/benches/helpers/sync.rs | 2 +- packages/http-core/benches/helpers/util.rs | 2 +- packages/http-core/src/services/announce.rs | 16 +- .../http-protocol/src/v1/requests/announce.rs | 32 +++- .../src/v1/responses/announce/data.rs | 1 - .../src/v1/responses/announce/encoding.rs | 11 +- packages/primitives/src/i2p.rs | 158 +++++++++++++++++ packages/primitives/src/lib.rs | 4 +- packages/primitives/src/peer.rs | 160 ++---------------- .../src/v1/adapters/torrent.rs | 2 +- .../src/v1/conversion.rs | 13 +- .../src/swarm/coordinator.rs | 8 +- packages/tracker-core/src/announce_handler.rs | 10 +- project-words.txt | 1 + 20 files changed, 254 insertions(+), 195 deletions(-) create mode 100644 packages/primitives/src/i2p.rs diff --git a/README.md b/README.md index 56932331c..83f8fc755 100644 --- a/README.md +++ b/README.md @@ -14,6 +14,7 @@ - [x] Good Performance in Busy Conditions. - [x] Support for `UDP`, `HTTP`, and `TLS` Sockets. - [x] Native `IPv4` and `IPv6` support. +- [x] [I2P] peer announces and matchmaking over HTTP. - [x] Private & Whitelisted mode. - [x] Tracker Management API. - [x] Support [newTrackon][newtrackon] checks. @@ -296,3 +297,4 @@ This project was a joint effort by [Nautilus Cyberneering GmbH][nautilus] and [D [Power2All]: https://github.com/power2all [torrust-demo]: https://github.com/torrust/torrust-demo [prometheus]: https://prometheus.io/ +[I2P]: https://i2p.net/en/docs/applications/bittorrent/ diff --git a/cspell.json b/cspell.json index be5f3d101..42d6e3199 100644 --- a/cspell.json +++ b/cspell.json @@ -18,6 +18,7 @@ ], "ignorePaths": [ ".tmp/**", + "storage/**", "target", "docs/media/*.svg", "contrib/bencode/benches/*.bencode", @@ -33,4 +34,4 @@ "contrib/dev-tools/git/github-merge.py", "docs/issues/**/evidence/*.html" ] -} \ No newline at end of file +} diff --git a/docs/packages.md b/docs/packages.md index 69eb24ef9..0b54b574e 100644 --- a/docs/packages.md +++ b/docs/packages.md @@ -277,6 +277,13 @@ Packages that have been extracted to their own standalone repositories. - Response bencoding - Error code mapping - Compact peer formatting + - I2P Destination parsing and compact Destination-hash formatting + +HTTP swarms keep I2P peers isolated from clearnet peers. An I2P announce may provide its full +Base64 Destination, with or without the `.i2p` suffix, in the `ip` query parameter. Non-compact +responses return that Destination, while compact responses return its 32-byte SHA-256 hash. +See the [I2P BitTorrent specification](https://i2p.net/en/docs/applications/bittorrent/) for the +wire-format details. ### UDP Tracker (BEP 15) diff --git a/packages/axum-http-server/src/v1/handlers/announce.rs b/packages/axum-http-server/src/v1/handlers/announce.rs index c3625d6d6..949af67c4 100644 --- a/packages/axum-http-server/src/v1/handlers/announce.rs +++ b/packages/axum-http-server/src/v1/handlers/announce.rs @@ -13,7 +13,7 @@ use torrust_tracker_http_core::services::announce::{AnnounceService, HttpAnnounc use torrust_tracker_http_protocol::v1::requests::announce::{Announce, Compact}; use torrust_tracker_http_protocol::v1::responses::{self}; use torrust_tracker_http_protocol::v1::services::peer_ip_resolver::ClientIpSources; -use torrust_tracker_primitives::AnnounceData as DomainAnnounceData; +use torrust_tracker_primitives::{AnnounceData as DomainAnnounceData, PeerAddress as DomainPeerAddress}; use crate::v1::extractors::announce_request::ExtractRequest; use crate::v1::extractors::authentication_key::Extract as ExtractKey; @@ -108,13 +108,10 @@ fn to_protocol_announce_data(domain_data: DomainAnnounceData) -> responses::anno .into_iter() .map(|peer| { let peer_addr = match &peer.peer_addr { - torrust_tracker_primitives::PeerAddress::Clearnet(address) => { - responses::announce::PeerAddress::Clearnet(*address) - } - torrust_tracker_primitives::PeerAddress::I2p(address) => responses::announce::PeerAddress::I2p { + DomainPeerAddress::Clearnet(address) => responses::announce::PeerAddress::Clearnet(*address), + DomainPeerAddress::I2p(address) => responses::announce::PeerAddress::I2p { destination: address.destination.to_string(), destination_hash: *address.destination.hash(), - port: 1, }, }; diff --git a/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs b/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs index fc5d45152..effb6462b 100644 --- a/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs +++ b/packages/axum-http-server/tests/server/v1/contract/for_all_config_modes/receiving_an_announce_request.rs @@ -105,7 +105,7 @@ async fn should_fail_when_url_query_parameters_are_invalid() { let http_tracker_config = Arc::new(cfg.http_trackers.unwrap()[0].clone()); let env = Started::new(&core_config, &http_tracker_config).await; - let invalid_query_param = "a=b=c"; + let invalid_query_param = "missing-value-separator"; let response = Client::new(env.base_url(), Duration::from_secs(5)) .unwrap() @@ -113,7 +113,7 @@ async fn should_fail_when_url_query_parameters_are_invalid() { .await .unwrap(); - assert_cannot_parse_query_param_error_response(response, "invalid param a=b=c").await; + assert_cannot_parse_query_param_error_response(response, "invalid param missing-value-separator").await; env.stop().await; } @@ -794,6 +794,7 @@ async fn it_should_return_i2p_destination_hashes_in_a_compact_response() { let env = Started::new(&core_config, &http_tracker_config).await; let client = Client::new(env.base_url(), Duration::from_secs(5)).unwrap(); let info_hash = InfoHash::from_str("9c38422213e30bff212b30c360d26f9a02136422").unwrap(); // DevSkim: ignore DS173237 + // cspell:disable-next-line let first_destination = format!("{}BQAEAAAAAA==.i2p", "A".repeat(512)) .parse::() .unwrap(); @@ -818,6 +819,7 @@ async fn it_should_return_i2p_destination_hashes_in_a_compact_response() { .with_peer_id(&PeerId(*b"-qB00000000000000002")) .with_port(1) .with_i2p_destination( + // cspell:disable-next-line format!("B{}BQAEAAAAAA==.i2p", "A".repeat(511)) .parse::() .unwrap(), diff --git a/packages/axum-rest-api-server/tests/server/v1/contract/context/torrent.rs b/packages/axum-rest-api-server/tests/server/v1/contract/context/torrent.rs index e0265ddc0..817103c43 100644 --- a/packages/axum-rest-api-server/tests/server/v1/contract/context/torrent.rs +++ b/packages/axum-rest-api-server/tests/server/v1/contract/context/torrent.rs @@ -333,7 +333,7 @@ async fn should_allow_getting_a_torrent_info() { seeders: 1, completed: 0, leechers: 0, - peers: Some(vec![conversion::from_domain_peer(peer)]), + peers: Some(vec![conversion::from_domain_peer(&peer)]), }, ) .await; diff --git a/packages/http-core/benches/helpers/sync.rs b/packages/http-core/benches/helpers/sync.rs index d487bc54a..1fb83d292 100644 --- a/packages/http-core/benches/helpers/sync.rs +++ b/packages/http-core/benches/helpers/sync.rs @@ -12,7 +12,7 @@ pub async fn return_announce_data_once(samples: u64) -> Duration { let peer = sample_peer(); - let (announce_request, client_ip_sources) = sample_announce_request_for_peer(peer); + let (announce_request, client_ip_sources) = sample_announce_request_for_peer(&peer); let announce_service = AnnounceService::new( core_tracker_services.core_config.clone(), diff --git a/packages/http-core/benches/helpers/util.rs b/packages/http-core/benches/helpers/util.rs index 62114e01c..364e56a2a 100644 --- a/packages/http-core/benches/helpers/util.rs +++ b/packages/http-core/benches/helpers/util.rs @@ -102,7 +102,7 @@ pub fn sample_peer() -> peer::Peer { } } -pub fn sample_announce_request_for_peer(peer: Peer) -> (Announce, ClientIpSources) { +pub fn sample_announce_request_for_peer(peer: &Peer) -> (Announce, ClientIpSources) { let announce_request = Announce { info_hash: sample_info_hash(), peer_id: peer.peer_id, diff --git a/packages/http-core/src/services/announce.rs b/packages/http-core/src/services/announce.rs index a5157f958..5a3826673 100644 --- a/packages/http-core/src/services/announce.rs +++ b/packages/http-core/src/services/announce.rs @@ -328,7 +328,7 @@ mod tests { ) } - fn sample_announce_request_for_peer(peer: Peer) -> (Announce, ClientIpSources) { + fn sample_announce_request_for_peer(peer: &Peer) -> (Announce, ClientIpSources) { let announce_request = Announce { info_hash: sample_info_hash(), peer_id: peer.peer_id, @@ -418,7 +418,7 @@ mod tests { let peer = sample_peer(); - let (announce_request, client_ip_sources) = sample_announce_request_for_peer(peer); + let (announce_request, client_ip_sources) = sample_announce_request_for_peer(&peer); let server_socket_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 7070); let server_service_binding = ServiceBinding::new(Protocol::HTTP, server_socket_addr).unwrap(); @@ -462,7 +462,7 @@ mod tests { core_http_tracker_services.http_stats_event_sender, ); - let (mut first_request, client_ip_sources) = sample_announce_request_for_peer(sample_peer()); + let (mut first_request, client_ip_sources) = sample_announce_request_for_peer(&sample_peer()); first_request.ip = Some(AnnounceAddress::I2p( format!("{}.i2p", "A".repeat(516)).parse::().unwrap(), )); @@ -473,9 +473,9 @@ mod tests { let mut second_peer = sample_peer(); second_peer.peer_id = PeerId(*b"-qB00000000000000002"); - let (mut second_request, _) = sample_announce_request_for_peer(second_peer); + let (mut second_request, _) = sample_announce_request_for_peer(&second_peer); second_request.ip = Some(AnnounceAddress::I2p( - format!("{}.i2p", "B".repeat(516)).parse::().unwrap(), + format!("B{}.i2p", "A".repeat(515)).parse::().unwrap(), )); let announce_data = announce_service @@ -523,7 +523,7 @@ mod tests { core_http_tracker_services.http_stats_event_sender = http_stats_event_sender; - let (announce_request, client_ip_sources) = sample_announce_request_for_peer(peer); + let (announce_request, client_ip_sources) = sample_announce_request_for_peer(&peer); let announce_service = AnnounceService::new( core_tracker_services.core_config.clone(), @@ -601,7 +601,7 @@ mod tests { core_http_tracker_services.http_stats_event_sender = http_stats_event_sender; - let (announce_request, client_ip_sources) = sample_announce_request_for_peer(peer); + let (announce_request, client_ip_sources) = sample_announce_request_for_peer(&peer); let announce_service = AnnounceService::new( core_tracker_services.core_config.clone(), @@ -647,7 +647,7 @@ mod tests { let (core_tracker_services, mut core_http_tracker_services) = initialize_core_tracker_services().await; core_http_tracker_services.http_stats_event_sender = http_stats_event_sender; - let (announce_request, client_ip_sources) = sample_announce_request_for_peer(peer); + let (announce_request, client_ip_sources) = sample_announce_request_for_peer(&peer); let announce_service = AnnounceService::new( core_tracker_services.core_config.clone(), diff --git a/packages/http-protocol/src/v1/requests/announce.rs b/packages/http-protocol/src/v1/requests/announce.rs index 055632d5c..4e4b86eb8 100644 --- a/packages/http-protocol/src/v1/requests/announce.rs +++ b/packages/http-protocol/src/v1/requests/announce.rs @@ -593,15 +593,20 @@ fn extract_ip(query: &Query) -> Result, ParseAnnounceQue return Ok(Some(AnnounceAddress::Ip(ip))); } - if raw_param.ends_with(".i2p") { - return I2pDestination::from_str(&raw_param) - .map(AnnounceAddress::I2p) - .map(Some) - .map_err(|_| ParseAnnounceQueryError::InvalidParam { + let has_i2p_suffix = raw_param + .rsplit_once('.') + .is_some_and(|(_, suffix)| suffix.eq_ignore_ascii_case("i2p")); + + match I2pDestination::from_str(&raw_param) { + Ok(destination) => return Ok(Some(AnnounceAddress::I2p(destination))), + Err(_) if has_i2p_suffix => { + return Err(ParseAnnounceQueryError::InvalidParam { param_name: IP.to_owned(), param_value: raw_param, location: Location::caller(), }); + } + Err(_) => {} } Ok(None) @@ -722,6 +727,7 @@ mod tests { fn it_should_parse_a_padded_i2p_destination_from_the_ip_param() { // 391 decoded bytes: 384 key bytes, a key certificate with its // four-byte key-type payload, and `==` Base64 padding. + // cspell:disable-next-line let destination = format!("{}BQAEAAAAAA==.i2p", "A".repeat(512)); let raw_query = Query::from(vec![ (INFO_HASH, "%3B%24U%04%CF%5F%11%BB%DB%E1%20%1C%EAjk%F4Z%EE%1B%C0"), @@ -736,6 +742,22 @@ mod tests { assert!(matches!(announce_request.ip, Some(AnnounceAddress::I2p(_)))); } + #[test] + fn it_should_parse_an_i2p_destination_without_the_i2p_suffix() { + let destination = "A".repeat(516); + let raw_query = Query::from(vec![ + (INFO_HASH, "%3B%24U%04%CF%5F%11%BB%DB%E1%20%1C%EAjk%F4Z%EE%1B%C0"), + (PEER_ID, "-RC3000-000000000001"), + (PORT, "1"), + (IP, &destination), + ]) + .to_string(); + + let announce_request = Announce::try_from(raw_query.parse::().unwrap()).unwrap(); + + assert!(matches!(announce_request.ip, Some(AnnounceAddress::I2p(_)))); + } + #[test] fn it_should_reject_an_invalid_i2p_destination() { let destination = "invalid.i2p"; diff --git a/packages/http-protocol/src/v1/responses/announce/data.rs b/packages/http-protocol/src/v1/responses/announce/data.rs index 6f20b633d..75d3272ff 100644 --- a/packages/http-protocol/src/v1/responses/announce/data.rs +++ b/packages/http-protocol/src/v1/responses/announce/data.rs @@ -63,6 +63,5 @@ pub enum PeerAddress { I2p { destination: String, destination_hash: [u8; 32], - port: u16, }, } diff --git a/packages/http-protocol/src/v1/responses/announce/encoding.rs b/packages/http-protocol/src/v1/responses/announce/encoding.rs index e11ac7766..11e1789c3 100644 --- a/packages/http-protocol/src/v1/responses/announce/encoding.rs +++ b/packages/http-protocol/src/v1/responses/announce/encoding.rs @@ -9,6 +9,8 @@ use torrust_bencode::{BMutAccess, BencodeMut, ben_bytes, ben_int, ben_list, ben_ use crate::v1::responses::announce::data::{AnnounceData, Peer, PeerAddress}; +const I2P_PLACEHOLDER_PORT: u16 = 1; + /// An [`Announce`] response, that can be anything that is convertible from [`AnnounceData`]. /// /// The [`Announce`] can built from any data that implements: [`From`] and [`Into>`]. @@ -25,6 +27,7 @@ use crate::v1::responses::announce::data::{AnnounceData, Peer, PeerAddress}; /// - [BEP 03: The `BitTorrent` Protocol Specification](https://www.bittorrent.org/beps/bep_0003.html) /// - [BEP 23: Tracker Returns Compact Peer Lists](https://www.bittorrent.org/beps/bep_0023.html) /// - [BEP 07: IPv6 Tracker Extension](https://www.bittorrent.org/beps/bep_0007.html) +/// - [I2P BitTorrent client protocol](https://i2p.net/en/docs/applications/bittorrent/) #[derive(Debug, AsRef, PartialEq, Constructor)] pub struct Announce @@ -153,7 +156,7 @@ impl Into> for Compact { pub struct NormalPeer { /// The peer's ID. pub peer_id: [u8; 20], - /// The peer's IP address. + /// The peer's IP address or I2P Destination. pub ip: String, /// The peer's port number. pub port: u16, @@ -167,10 +170,10 @@ impl From for NormalPeer { ip: address.ip().to_string(), port: address.port(), }, - PeerAddress::I2p { destination, port, .. } => NormalPeer { + PeerAddress::I2p { destination, .. } => NormalPeer { peer_id: peer.peer_id.0, ip: destination, - port, + port: I2P_PLACEHOLDER_PORT, }, } } @@ -369,7 +372,6 @@ mod tests { peer_addr: PeerAddress::I2p { destination: destination.clone(), destination_hash: [7; 32], - port: 1, }, }], SwarmMetadata::default(), @@ -393,7 +395,6 @@ mod tests { peer_addr: PeerAddress::I2p { destination: format!("{}.i2p", "A".repeat(516)), destination_hash, - port: 1, }, }], SwarmMetadata::default(), diff --git a/packages/primitives/src/i2p.rs b/packages/primitives/src/i2p.rs new file mode 100644 index 000000000..1fd6e68e1 --- /dev/null +++ b/packages/primitives/src/i2p.rs @@ -0,0 +1,158 @@ +//! I2P addressing primitives. + +use std::fmt; +use std::str::FromStr; + +use base64::Engine; +use base64::alphabet::Alphabet; +use base64::engine::{GeneralPurpose, GeneralPurposeConfig}; +use sha2::{Digest, Sha256}; +use thiserror::Error; + +const I2P_BASE64_ALPHABET: &str = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-~"; +const I2P_SUFFIX: &str = "i2p"; +const MIN_I2P_DESTINATION_BYTES: usize = 387; +const I2P_CERTIFICATE_LENGTH_OFFSET: usize = 385; + +/// A validated I2P Base64 Destination. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct I2pDestination { + value: Box, + hash: [u8; 32], +} + +impl I2pDestination { + /// Returns the SHA-256 hash of the decoded binary Destination. + #[must_use] + pub const fn hash(&self) -> &[u8; 32] { + &self.hash + } +} + +impl fmt::Display for I2pDestination { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(&self.value) + } +} + +impl FromStr for I2pDestination { + type Err = ParseI2pDestinationError; + + fn from_str(value: &str) -> Result { + let encoded = value + .rsplit_once('.') + .filter(|(_, suffix)| suffix.eq_ignore_ascii_case(I2P_SUFFIX)) + .map_or(value, |(encoded, _)| encoded); + let alphabet = + Alphabet::new(I2P_BASE64_ALPHABET).expect("the I2P Base64 alphabet must contain 64 unique ASCII characters"); + let engine = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()); + let decoded = engine.decode(encoded).map_err(|_| ParseI2pDestinationError::InvalidBase64)?; + + if decoded.len() < MIN_I2P_DESTINATION_BYTES { + return Err(ParseI2pDestinationError::TooShort { actual: decoded.len() }); + } + + let certificate_payload_length = usize::from(u16::from_be_bytes([ + decoded[I2P_CERTIFICATE_LENGTH_OFFSET], + decoded[I2P_CERTIFICATE_LENGTH_OFFSET + 1], + ])); + let expected_length = MIN_I2P_DESTINATION_BYTES + certificate_payload_length; + + if decoded.len() != expected_length { + return Err(ParseI2pDestinationError::InvalidCertificateLength { + declared: certificate_payload_length, + actual: decoded.len() - MIN_I2P_DESTINATION_BYTES, + }); + } + + Ok(Self { + value: format!("{encoded}.{I2P_SUFFIX}").into_boxed_str(), + hash: Sha256::digest(decoded).into(), + }) + } +} + +/// Error returned when parsing an I2P Destination. +#[derive(Debug, Error, PartialEq, Eq)] +pub enum ParseI2pDestinationError { + #[error("the I2P Destination is not valid I2P Base64")] + InvalidBase64, + #[error("the decoded I2P Destination must contain at least {MIN_I2P_DESTINATION_BYTES} bytes, got {actual}")] + TooShort { actual: usize }, + #[error("the I2P certificate declares a {declared}-byte payload, but the Destination contains {actual} payload bytes")] + InvalidCertificateLength { declared: usize, actual: usize }, +} + +/// An I2P peer address. I2P routes by Destination and has no peer port. +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct I2pPeerAddress { + /// The peer's full I2P Destination. + pub destination: I2pDestination, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn it_should_parse_and_normalize_a_valid_i2p_base64_destination() { + let destination = "A".repeat(516); + let destination_with_uppercase_suffix = format!("{destination}.I2P"); + + let parsed = I2pDestination::from_str(&destination_with_uppercase_suffix).unwrap(); + + assert_eq!(parsed.to_string(), format!("{destination}.i2p")); + assert_eq!( + parsed.hash(), + &[ + 0x31, 0x19, 0xfc, 0xeb, 0x0e, 0xad, 0x1d, 0x08, 0x04, 0xdb, 0x90, 0xfb, 0x0c, 0x87, 0xa3, 0x38, 0x10, 0x89, 0xf9, + 0xd2, 0x26, 0x4a, 0x37, 0x6c, 0x41, 0xa3, 0x9a, 0x06, 0xe5, 0x32, 0xa6, 0x41, + ] + ); + } + + #[test] + fn it_should_reject_an_i2p_destination_with_invalid_base64() { + let destination = format!("{}.i2p", "!".repeat(516)); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!(error, ParseI2pDestinationError::InvalidBase64); + } + + #[test] + fn it_should_reject_an_i2p_destination_shorter_than_the_minimum_length() { + let destination = "A".repeat(512); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!(error, ParseI2pDestinationError::TooShort { actual: 384 }); + } + + #[test] + fn it_should_parse_a_long_padded_destination_when_the_certificate_length_matches() { + let certificate_payload_length = 91_u16; + let mut decoded = vec![0; 387 + usize::from(certificate_payload_length)]; + decoded[384] = 5; + decoded[385..387].copy_from_slice(&certificate_payload_length.to_be_bytes()); + let alphabet = Alphabet::new(I2P_BASE64_ALPHABET).unwrap(); + let encoded = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()).encode(decoded); + + let parsed = I2pDestination::from_str(&encoded).unwrap(); + + assert!(encoded.ends_with("==")); + assert_eq!(parsed.to_string(), format!("{encoded}.i2p")); + } + + #[test] + fn it_should_reject_a_destination_when_the_certificate_length_does_not_match() { + let destination = "A".repeat(520); + + let error = I2pDestination::from_str(&destination).unwrap_err(); + + assert_eq!( + error, + ParseI2pDestinationError::InvalidCertificateLength { declared: 0, actual: 3 } + ); + } +} diff --git a/packages/primitives/src/lib.rs b/packages/primitives/src/lib.rs index e2b23c179..5403d6bea 100644 --- a/packages/primitives/src/lib.rs +++ b/packages/primitives/src/lib.rs @@ -7,6 +7,7 @@ pub mod announce; pub mod configuration_instance_id; pub mod driver; +pub mod i2p; pub mod mode; pub mod number_of_bytes; pub mod pagination; @@ -27,9 +28,10 @@ use std::collections::BTreeMap; pub use announce::{AnnounceData, AnnounceEvent, AnnouncePolicy}; pub use configuration_instance_id::ConfigurationInstanceId; pub use driver::Driver; +pub use i2p::{I2pDestination, I2pPeerAddress}; pub use mode::PrivateMode; pub use number_of_bytes::NumberOfBytes; -pub use peer::{I2pDestination, I2pPeerAddress, PeerAddress}; +pub use peer::PeerAddress; pub use policy::TrackerPolicy; pub use scrape::ScrapeData; pub use service_role::ServiceRole; diff --git a/packages/primitives/src/peer.rs b/packages/primitives/src/peer.rs index 16050a5e0..0bb11cb09 100644 --- a/packages/primitives/src/peer.rs +++ b/packages/primitives/src/peer.rs @@ -28,95 +28,14 @@ use std::ops::{Deref, DerefMut}; use std::str::FromStr; use std::sync::Arc; -use base64::Engine; -use base64::alphabet::Alphabet; -use base64::engine::{GeneralPurpose, GeneralPurposeConfig}; use serde::Serialize; -use sha2::{Digest, Sha256}; use thiserror::Error; use torrust_clock::DurationSinceUnixEpoch; -use crate::{AnnounceEvent, NumberOfBytes, PeerId}; +use crate::{AnnounceEvent, I2pPeerAddress, NumberOfBytes, PeerId}; pub type PeerAnnouncement = Peer; -const I2P_BASE64_ALPHABET: &str = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-~"; -const I2P_SUFFIX: &str = ".i2p"; -const MIN_I2P_DESTINATION_BYTES: usize = 387; -const I2P_CERTIFICATE_LENGTH_OFFSET: usize = 385; - -/// A validated I2P Base64 Destination. -#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] -pub struct I2pDestination { - value: Box, - hash: [u8; 32], -} - -impl I2pDestination { - /// Returns the SHA-256 hash of the decoded binary Destination. - #[must_use] - pub const fn hash(&self) -> &[u8; 32] { - &self.hash - } -} - -impl fmt::Display for I2pDestination { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.write_str(&self.value) - } -} - -impl FromStr for I2pDestination { - type Err = ParseI2pDestinationError; - - fn from_str(value: &str) -> Result { - let encoded = value.strip_suffix(I2P_SUFFIX).unwrap_or(value); - let alphabet = - Alphabet::new(I2P_BASE64_ALPHABET).expect("the I2P Base64 alphabet must contain 64 unique ASCII characters"); - let engine = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()); - let decoded = engine.decode(encoded).map_err(|_| ParseI2pDestinationError::InvalidBase64)?; - - if decoded.len() < MIN_I2P_DESTINATION_BYTES { - return Err(ParseI2pDestinationError::TooShort { actual: decoded.len() }); - } - - let certificate_payload_length = usize::from(u16::from_be_bytes([ - decoded[I2P_CERTIFICATE_LENGTH_OFFSET], - decoded[I2P_CERTIFICATE_LENGTH_OFFSET + 1], - ])); - let expected_length = MIN_I2P_DESTINATION_BYTES + certificate_payload_length; - - if decoded.len() != expected_length { - return Err(ParseI2pDestinationError::InvalidCertificateLength { - declared: certificate_payload_length, - actual: decoded.len() - MIN_I2P_DESTINATION_BYTES, - }); - } - - Ok(Self { - value: format!("{encoded}{I2P_SUFFIX}").into_boxed_str(), - hash: Sha256::digest(decoded).into(), - }) - } -} - -/// Error returned when parsing an I2P Destination. -#[derive(Debug, Error, PartialEq, Eq)] -pub enum ParseI2pDestinationError { - #[error("the I2P Destination is not valid I2P Base64")] - InvalidBase64, - #[error("the decoded I2P Destination must contain at least {MIN_I2P_DESTINATION_BYTES} bytes, got {actual}")] - TooShort { actual: usize }, - #[error("the I2P certificate declares a {declared}-byte payload, but the Destination contains {actual} payload bytes")] - InvalidCertificateLength { declared: usize, actual: usize }, -} - -/// An I2P peer address. I2P routes by Destination and has no peer port. -#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] -pub struct I2pPeerAddress { - pub destination: I2pDestination, -} - /// A peer endpoint on either the public Internet or I2P. #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub enum PeerAddress { @@ -125,6 +44,10 @@ pub enum PeerAddress { } impl PeerAddress { + /// Returns the clearnet port, or the conventional placeholder port `1` for I2P. + /// + /// I2P routes by Destination rather than by port. The placeholder is needed + /// by non-compact tracker responses and legacy peer representations. #[must_use] pub const fn port(&self) -> u16 { match self { @@ -135,6 +58,7 @@ impl PeerAddress { } } + /// Returns the peer's clearnet IP address, or `None` for an I2P peer. #[must_use] pub const fn ip(&self) -> Option { match self { @@ -143,11 +67,13 @@ impl PeerAddress { } } + /// Returns whether this is an I2P peer address. #[must_use] pub const fn is_i2p(&self) -> bool { matches!(self, Self::I2p(_)) } + /// Returns the peer's clearnet socket address, or `None` for an I2P peer. #[must_use] pub const fn socket_addr(&self) -> Option { match self { @@ -426,11 +352,14 @@ impl Peer { } } + /// Returns the peer's clearnet IP address, or `None` for an I2P peer. + #[must_use] pub fn ip(&self) -> Option { self.peer_addr.ip() } - pub fn change_ip(&mut self, new_ip: &IpAddr) { + /// Replaces the IP of a clearnet peer and leaves an I2P peer unchanged. + pub fn set_clearnet_ip(&mut self, new_ip: &IpAddr) { if let PeerAddress::Clearnet(address) = &mut self.peer_addr { address.set_ip(*new_ip); } @@ -785,71 +714,6 @@ pub mod fixture { #[cfg(test)] pub mod test { - mod i2p_destination { - use std::str::FromStr; - - use base64::Engine; - use base64::alphabet::Alphabet; - use base64::engine::{GeneralPurpose, GeneralPurposeConfig}; - - use crate::peer::{I2pDestination, ParseI2pDestinationError}; - - #[test] - fn it_should_parse_a_valid_i2p_base64_destination() { - let destination = "A".repeat(516); - - let parsed = I2pDestination::from_str(&destination).unwrap(); - - assert_eq!(parsed.to_string(), format!("{destination}.i2p")); - assert_eq!(parsed.hash().len(), 32); - } - - #[test] - fn it_should_reject_an_i2p_destination_with_invalid_base64() { - let destination = format!("{}.i2p", "!".repeat(516)); - - let error = I2pDestination::from_str(&destination).unwrap_err(); - - assert_eq!(error, ParseI2pDestinationError::InvalidBase64); - } - - #[test] - fn it_should_reject_an_i2p_destination_shorter_than_the_minimum_length() { - let destination = "A".repeat(512); - - let error = I2pDestination::from_str(&destination).unwrap_err(); - - assert_eq!(error, ParseI2pDestinationError::TooShort { actual: 384 }); - } - - #[test] - fn it_should_parse_a_long_padded_destination_when_the_certificate_length_matches() { - let certificate_payload_length = 91_u16; - let mut decoded = vec![0; 387 + usize::from(certificate_payload_length)]; - decoded[384] = 5; - decoded[385..387].copy_from_slice(&certificate_payload_length.to_be_bytes()); - let alphabet = Alphabet::new("ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-~").unwrap(); - let encoded = GeneralPurpose::new(&alphabet, GeneralPurposeConfig::new()).encode(decoded); - - let parsed = I2pDestination::from_str(&encoded).unwrap(); - - assert!(encoded.ends_with("==")); - assert_eq!(parsed.to_string(), format!("{encoded}.i2p")); - } - - #[test] - fn it_should_reject_a_destination_when_the_certificate_length_does_not_match() { - let destination = "A".repeat(520); - - let error = I2pDestination::from_str(&destination).unwrap_err(); - - assert_eq!( - error, - ParseI2pDestinationError::InvalidCertificateLength { declared: 0, actual: 3 } - ); - } - } - mod peer { use crate::peer::fixture::PeerBuilder; diff --git a/packages/rest-api-runtime-adapter/src/v1/adapters/torrent.rs b/packages/rest-api-runtime-adapter/src/v1/adapters/torrent.rs index 0b304af1f..010f95b90 100644 --- a/packages/rest-api-runtime-adapter/src/v1/adapters/torrent.rs +++ b/packages/rest-api-runtime-adapter/src/v1/adapters/torrent.rs @@ -32,7 +32,7 @@ impl TorrentQueryPort for TrackerTorrentQueryAdapter { async fn get_torrent_info(&self, info_hash: &InfoHash) -> Option { services::get_torrent_info(&self.in_memory_torrent_repository, info_hash) .await - .map(conversion::from_domain_info) + .map(|info| conversion::from_domain_info(&info)) } async fn get_torrents_page(&self, pagination: &Pagination) -> Vec { diff --git a/packages/rest-api-runtime-adapter/src/v1/conversion.rs b/packages/rest-api-runtime-adapter/src/v1/conversion.rs index 78abfa8ea..d0fc7a216 100644 --- a/packages/rest-api-runtime-adapter/src/v1/conversion.rs +++ b/packages/rest-api-runtime-adapter/src/v1/conversion.rs @@ -9,7 +9,7 @@ use torrust_tracker_rest_api_protocol::v1::context::torrent::resources::torrent: /// Convert a domain [`domain_peer::Peer`] into a protocol [`protocol_peer::Peer`]. #[must_use] -pub fn from_domain_peer(value: domain_peer::Peer) -> protocol_peer::Peer { +pub fn from_domain_peer(value: &domain_peer::Peer) -> protocol_peer::Peer { #[allow(deprecated)] protocol_peer::Peer { peer_id: from_domain_peer_id(value.peer_id), @@ -35,8 +35,11 @@ pub fn from_domain_peer_id(peer_id: PeerId) -> protocol_peer::Id { /// Convert a domain [`Info`] into a protocol [`Torrent`]. #[must_use] -pub fn from_domain_info(info: Info) -> Torrent { - let peers: Option> = info.peers.map(|peers| peers.into_iter().map(from_domain_peer).collect()); +pub fn from_domain_info(info: &Info) -> Torrent { + let peers: Option> = info + .peers + .as_deref() + .map(|peers| peers.iter().map(from_domain_peer).collect()); Torrent { info_hash: info.info_hash.to_string(), @@ -92,7 +95,7 @@ mod tests { #[test] fn torrent_resource_should_be_converted_from_torrent_info() { assert_eq!( - from_domain_info(Info { + from_domain_info(&Info { info_hash: InfoHash::from_str("9e0217d0fa71c87332cd8bf9dbeabcb2c2cf3c4d").unwrap(), // DevSkim: ignore DS173237 seeders: 1, completed: 2, @@ -104,7 +107,7 @@ mod tests { seeders: 1, completed: 2, leechers: 3, - peers: Some(vec![from_domain_peer(sample_peer())]), + peers: Some(vec![from_domain_peer(&sample_peer())]), } ); } diff --git a/packages/swarm-coordination-registry/src/swarm/coordinator.rs b/packages/swarm-coordination-registry/src/swarm/coordinator.rs index 4c53d8e24..318506789 100644 --- a/packages/swarm-coordination-registry/src/swarm/coordinator.rs +++ b/packages/swarm-coordination-registry/src/swarm/coordinator.rs @@ -442,7 +442,7 @@ mod tests { } #[tokio::test] - async fn it_should_not_return_clearnet_peers_to_an_i2p_peer() { + async fn it_should_not_return_peers_from_a_different_network() { let mut swarm = Coordinator::new(&sample_info_hash(), 0, None); let clearnet_peer = PeerBuilder::default().build(); swarm.upsert_peer(clearnet_peer.clone().into()).await; @@ -453,9 +453,11 @@ mod tests { }); swarm.upsert_peer(i2p_peer.clone().into()).await; - let peers = swarm.peers_excluding(&i2p_peer.peer_addr, None); + let peers_for_i2p = swarm.peers_excluding(&i2p_peer.peer_addr, None); + let peers_for_clearnet = swarm.peers_excluding(&clearnet_peer.peer_addr, None); - assert!(peers.is_empty()); + assert!(peers_for_i2p.is_empty()); + assert!(peers_for_clearnet.is_empty()); } #[tokio::test] diff --git a/packages/tracker-core/src/announce_handler.rs b/packages/tracker-core/src/announce_handler.rs index 304522ee1..eda7e259e 100644 --- a/packages/tracker-core/src/announce_handler.rs +++ b/packages/tracker-core/src/announce_handler.rs @@ -163,12 +163,10 @@ impl AnnounceHandler { ) -> Result { self.whitelist_authorization.authorize(info_hash).await?; - if !peer.peer_addr.is_i2p() { - peer.change_ip(&assign_ip_address_to_peer( - remote_client_ip, - self.config.net.external_ip.map(Into::into), - )); - } + peer.set_clearnet_ip(&assign_ip_address_to_peer( + remote_client_ip, + self.config.net.external_ip.map(Into::into), + )); self.in_memory_torrent_repository .handle_announcement(info_hash, peer, self.load_downloads_metric_if_needed(info_hash).await?) diff --git a/project-words.txt b/project-words.txt index 863a4996f..230dff62c 100644 --- a/project-words.txt +++ b/project-words.txt @@ -37,6 +37,7 @@ Graphviz Grcov HDRINCL Hydranode +I2P IPPROTO IPV6 Icelake