diff --git a/Cargo.toml b/Cargo.toml index 1d10463..a19e719 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -59,12 +59,12 @@ indexeddb = ["dep:idb", "serde"] redb = ["dep:redb", "dep:cbor4ii", "serde"] sqlite = ["dep:sqlx", "dep:cbor4ii", "serde"] -webrtc =["dep:libp2p-webrtc", "dep:libp2p-webrtc-websys"] +webrtc = ["dep:libp2p-webrtc", "dep:libp2p-webrtc-websys"] websocket = ["libp2p/websocket", "rcgen", "dep:pem", "libp2p/websocket-websys"] webtransport = ["libp2p/webtransport-websys"] [workspace] -members = ["examples/browser-webrtc", "examples/custom-behaviour-and-context", "examples/custom-transport", "examples/distributed-key-value-store", "examples/file-sharing-request-response", "examples/file-sharing-stream", "examples/floodsub", "examples/gossipsub", "examples/ipfs-kad", "examples/peer-store", "examples/relay", "examples/rendezvous", "examples/stream", "examples/upnp"] +members = ["examples/autorelay", "examples/browser-webrtc", "examples/custom-behaviour-and-context", "examples/custom-transport", "examples/distributed-key-value-store", "examples/file-sharing-request-response", "examples/file-sharing-stream", "examples/floodsub", "examples/gossipsub", "examples/ipfs-kad", "examples/peer-store", "examples/relay", "examples/rendezvous", "examples/stream", "examples/upnp"] [workspace.dependencies] async-rt = "0.1.8" @@ -80,15 +80,15 @@ futures-timer = "3.0.0" futures-timeout = "0.1.3" getrandom = { version = "0.2.15" } getrandom_03 = { version = "0.3.3", package = "getrandom" } -hickory-resolver = "0.25.2" +hickory-resolver = "0.26.1" idb = "0.6.5" indexmap = "2.10.0" -libp2p = { version = "0.56.0" } -libp2p-allow-block-list = "0.6.0" -libp2p-connection-limits = "0.6.0" -libp2p-stream = { version = "=0.4.0-alpha" } -libp2p-webrtc = { version = "=0.9.0-alpha.1", features = ["pem"] } -libp2p-webrtc-websys = "0.4.0" +libp2p = { version = "0.57.0" } +libp2p-allow-block-list = "0.7.0" +libp2p-connection-limits = "0.7.0" +libp2p-stream = { version = "0.5.0-alpha" } +libp2p-webrtc = { version = "=0.10.0-alpha", features = ["pem"] } +libp2p-webrtc-websys = "0.5.0" other-error = "0.1.1" pollable-map = "0.1.7" pem = { version = "3.0.5" } @@ -156,10 +156,18 @@ futures-timer = { workspace = true, features = ["wasm-bindgen"] } getrandom = { workspace = true, features = ["js"] } getrandom_03 = { workspace = true, features = ["wasm_js"] } idb = { workspace = true, optional = true } -libp2p = { features = ["macros", "serde", "wasm-bindgen"], workspace = true } +libp2p = { features = ["macros", "serde"], workspace = true } libp2p-webrtc-websys = { workspace = true, optional = true } send_wrapper = { workspace = true, features = ["futures"] } serde-wasm-bindgen.workspace = true tokio = { default-features = false, features = ["sync", "macros"], workspace = true } wasm-bindgen.workspace = true wasm-bindgen-futures.workspace = true + +[patch.crates-io] +libp2p = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } +libp2p-allow-block-list = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } +libp2p-connection-limits = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } +libp2p-stream = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } +libp2p-webrtc = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } +libp2p-webrtc-websys = { git = "https://github.com/libp2p/rust-libp2p.git", branch = "master" } \ No newline at end of file diff --git a/examples/autorelay/Cargo.toml b/examples/autorelay/Cargo.toml new file mode 100644 index 0000000..387430e --- /dev/null +++ b/examples/autorelay/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "autorelay" +version = "0.1.0" +edition = "2024" + +[dependencies] +tokio.workspace = true +futures.workspace = true +connexa = { path = "../../", default-features = false, features = ["ed25519", "rsa", "tcp", "quic", "yamux", "noise", "relay", "ping", "identify", "testing", "kad", "dns"] } +clap = { version = "4.5.39", features = ["derive"] } +tracing-subscriber = { version = "0.3.19", features = ["env-filter"] } diff --git a/examples/autorelay/src/main.rs b/examples/autorelay/src/main.rs new file mode 100644 index 0000000..770f138 --- /dev/null +++ b/examples/autorelay/src/main.rs @@ -0,0 +1,55 @@ +use connexa::prelude::{DefaultConnexaBuilder, Multiaddr, PeerId}; + +pub const BOOTSTRAP_NODES: &[(&str, &str)] = &[ + ( + "/ip4/104.131.131.82/tcp/4001", + "QmaCpDMGvV2BGHeYERUEnRQAwe3N8SzbUtfsmvsqQLuvuJ", + ), + ( + "/dnsaddr/bootstrap.libp2p.io", + "QmNnooDu7bfjPFoTZYxMNLWUQJyrVwtbZg5gBMjTezGAJN", + ), + ( + "/dnsaddr/bootstrap.libp2p.io", + "QmQCU2EcMqAqQPR2i9bChDtGNJchTbq5TbXJJ16u19uLTa", + ), + ( + "/dnsaddr/bootstrap.libp2p.io", + "QmbLHAnMoJPWSCR5Zhtx6BHJX9KiKNN6tpvbUcqanj75Nb", + ), + ( + "/dnsaddr/bootstrap.libp2p.io", + "QmcZf59bWwK5XFi76CZX8cbJ4BhTzzA3gU1ZjYZcYW3dwt", + ), +]; + +#[tokio::main] +async fn main() -> std::io::Result<()> { + let connexa = DefaultConnexaBuilder::new_identity() + .enable_tcp() + .enable_quic() + .enable_dns() + .with_relay() + .with_autorelay() + .with_ping() + .with_identify() + .with_kademlia() + .build() + .await?; + + for (addr, peer_id) in BOOTSTRAP_NODES { + let peer_id: PeerId = peer_id.parse().expect("valid peer id"); + let addr: Multiaddr = addr.parse().expect("valid addr"); + connexa.dht().add_address(peer_id, addr).await?; + } + + tokio::time::sleep(std::time::Duration::from_secs(10)).await; + + let external_addrs = connexa.swarm().external_addresses().await?; + + for addr in external_addrs { + println!("- {}", addr); + } + + Ok(()) +} diff --git a/examples/upnp/src/main.rs b/examples/upnp/src/main.rs index d1f16ed..af34719 100644 --- a/examples/upnp/src/main.rs +++ b/examples/upnp/src/main.rs @@ -13,9 +13,11 @@ async fn main() -> std::io::Result<()> { println!("New listen address: {addr}") } SwarmEvent::Behaviour(BehaviourEvent::Upnp(event)) => match event { - UpnpEvent::NewExternalAddr(addr) => println!("New external address: {addr}"), - UpnpEvent::ExpiredExternalAddr(addr) => { - println!("Expired external address: {addr}") + UpnpEvent::NewExternalAddr { external_addr, .. } => { + println!("New external address: {external_addr}") + } + UpnpEvent::ExpiredExternalAddr { external_addr, .. } => { + println!("Expired external address: {external_addr}") } UpnpEvent::GatewayNotFound => println!("Gateway not found"), UpnpEvent::NonRoutableGateway => println!("Gateway is not routable"), diff --git a/src/behaviour.rs b/src/behaviour.rs index 0e61aed..82000b1 100644 --- a/src/behaviour.rs +++ b/src/behaviour.rs @@ -1,3 +1,5 @@ +#[cfg(feature = "relay")] +pub mod autorelay; pub mod dummy; pub mod peer_store; #[cfg(feature = "request-response")] @@ -56,6 +58,9 @@ where #[cfg(feature = "relay")] pub relay_client: Toggle, + #[cfg(feature = "relay")] + pub autorelay: Toggle, + #[cfg(not(target_arch = "wasm32"))] #[cfg(feature = "upnp")] pub upnp: Toggle, @@ -270,6 +275,15 @@ where } false => (None, None.into()), }; + #[cfg(feature = "relay")] + let autorelay = protocols + .autorelay + .then(|| { + let config_fn = config.autorelay_config; + let config = config_fn(autorelay::Config::default()); + autorelay::Behaviour::new_with_config(config) + }) + .into(); #[cfg(not(feature = "relay"))] let transport = None::<()>; @@ -361,6 +375,8 @@ where relay, #[cfg(feature = "relay")] relay_client, + #[cfg(feature = "relay")] + autorelay, #[cfg(feature = "stream")] stream, #[cfg(not(target_arch = "wasm32"))] diff --git a/src/behaviour/autorelay.rs b/src/behaviour/autorelay.rs new file mode 100644 index 0000000..3d70316 --- /dev/null +++ b/src/behaviour/autorelay.rs @@ -0,0 +1,829 @@ +// TODO: Replace with builtin autorelay behaviour from libp2p. See https://github.com/libp2p/rust-libp2p/pull/6156 +use std::{ + collections::{HashMap, HashSet, VecDeque}, + num::NonZeroU8, + task::{Context, Poll, Waker}, + time::Duration, +}; + +use crate::behaviour::autorelay::handler::Out; +use crate::multiaddr_ext::MultiaddrExt; +use crate::prelude::swarm::derive_prelude::{ListenerId, PortUse}; +use crate::prelude::swarm::{ + ExternalAddresses, ListenOpts, NewListenAddr, NotifyHandler, + derive_prelude::{ + AddressChange, ConnectionClosed, ConnectionDenied, ConnectionEstablished, ConnectionId, + DialFailure, ExpiredListenAddr, FromSwarm, ListenerClosed, ListenerError, Multiaddr, + NetworkBehaviour, THandler, THandlerInEvent, THandlerOutEvent, ToSwarm, + }, + dial_opts::DialOpts, + dummy, +}; +use crate::prelude::transport::Endpoint; +use crate::prelude::{PeerId, Protocol}; +use either::Either; +use web_time::{Instant, SystemTime}; + +mod handler; + +#[derive(Debug)] +pub struct Behaviour { + config: Config, + status: Status, + auto_status_change: bool, + external_addresses: ExternalAddresses, + events: VecDeque::ToSwarm, THandlerInEvent>>, + + connections: HashMap<(PeerId, ConnectionId), Connection>, + + reservations: HashMap, + + external_reservations: HashMap, + + static_relays: HashMap>, + + static_dial_cooldowns: HashMap, + + failure_counts: HashMap, + + previous_relays: VecDeque<(PeerId, Multiaddr, SystemTime)>, + + relays_available: bool, + + waker: Option, +} + +impl Default for Behaviour { + fn default() -> Self { + Self { + config: Config::default(), + status: Status::Enable, + auto_status_change: true, + external_addresses: ExternalAddresses::default(), + events: VecDeque::new(), + connections: HashMap::new(), + reservations: HashMap::new(), + external_reservations: HashMap::new(), + static_relays: HashMap::new(), + static_dial_cooldowns: HashMap::new(), + failure_counts: HashMap::new(), + previous_relays: VecDeque::new(), + relays_available: false, + waker: None, + } + } +} + +#[derive(Default, Debug, Clone, Copy, PartialEq, Eq)] +pub enum Status { + #[default] + Enable, + Disable, +} + +#[derive(Debug)] +struct Connection { + address: Multiaddr, + relay_status: RelayStatus, +} + +impl Connection { + /// Mark relayed connection as not supported + pub(crate) fn disqualify_connection_if_relayed(&mut self) { + if self.address.is_relayed() { + self.relay_status = RelayStatus::NotSupported; + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum RelayStatus { + Supported { status: ReservationStatus }, + NotSupported, + Pending, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum ReservationStatus { + Idle, + Pending { id: ListenerId }, + Active { id: ListenerId }, + Blacklisted, +} + +#[derive(Debug)] +pub struct Config { + max_reservations: NonZeroU8, + failure_cooldown: Duration, + failure_cooldown_max: Duration, + max_previous_relays: usize, + static_relays: HashMap>, +} + +impl Default for Config { + fn default() -> Self { + Self { + max_reservations: NonZeroU8::new(2).unwrap(), + failure_cooldown: Duration::from_secs(30), + failure_cooldown_max: Duration::from_secs(10 * 60), + max_previous_relays: 16, + static_relays: HashMap::new(), + } + } +} + +impl Config { + pub fn set_max_reservations(mut self, max_reservations: NonZeroU8) -> Self { + self.max_reservations = max_reservations; + self + } + + pub fn set_failure_cooldown(mut self, duration: Duration) -> Self { + self.failure_cooldown = duration; + self + } + + pub fn set_failure_cooldown_max(mut self, duration: Duration) -> Self { + self.failure_cooldown_max = duration; + self + } + + pub fn set_max_previous_relays(mut self, max: usize) -> Self { + self.max_previous_relays = max; + self + } + + pub fn add_static_relay(mut self, peer_id: PeerId, addresses: Vec) -> Self { + let entry = self.static_relays.entry(peer_id).or_default(); + for addr in addresses { + if !entry.contains(&addr) { + entry.push(addr); + } + } + self + } +} + +#[derive(Debug)] +#[non_exhaustive] +pub enum Event { + /// The status of the local node has changed. + StatusChanged { status: Status }, + /// No connected peer supports the HOP protocol. + NoRelaysAvailable, + /// At least one connected peer supports the HOP protocol. + RelaysAvailable, +} + +impl Behaviour { + pub fn new_with_config(mut config: Config) -> Self { + let initial_static_relays = std::mem::take(&mut config.static_relays); + let mut behaviour = Self { + config, + ..Default::default() + }; + for (peer_id, addresses) in initial_static_relays { + for address in addresses { + behaviour.add_static_relay(peer_id, address); + } + } + behaviour + } + + /// Sets the autorelay status. + pub fn set_status(&mut self, status: Option) { + match status { + Some(status) => { + self.auto_status_change = false; + if self.status != status { + self.status = status; + self.events + .push_back(ToSwarm::GenerateEvent(Event::StatusChanged { status })); + if status == Status::Enable { + self.meet_reservation_target(); + } + } + } + None => { + self.auto_status_change = true; + self.determine_status_from_external_addresses(); + } + } + + if let Some(waker) = self.waker.take() { + waker.wake(); + } + } + + /// Register a peer as a static relay. + /// + /// This will dial and establish a connection to the peer if it doesn't already have a direct + /// connection. + /// Note that peers that are through a relay cannot be used as a static peer + pub fn add_static_relay(&mut self, peer_id: PeerId, address: Multiaddr) { + if address.is_relayed() { + tracing::warn!(%peer_id, %address, "static relay address is relayed. ignoring."); + return; + } + + let entry = self.static_relays.entry(peer_id).or_default(); + if entry.contains(&address) { + tracing::warn!(%peer_id, %address, "static relay address already exist"); + } else { + entry.push(address); + } + let combined = entry.clone(); + + if self.is_peer_idle(&peer_id) { + self.evict_for_static_peer(peer_id); + } + + if !self.queue_static_dial(peer_id, combined) { + self.meet_reservation_target(); + } + + if let Some(waker) = self.waker.take() { + waker.wake(); + } + } + + /// Remove peer as a static relay. + /// This will not close any connections or terminate any existing reservation with the relay + pub fn remove_static_relay(&mut self, peer_id: &PeerId) -> bool { + self.static_dial_cooldowns.remove(peer_id); + self.static_relays.remove(peer_id).is_some() + } + + pub fn static_relays(&self) -> impl Iterator { + self.static_relays + .iter() + .map(|(peer, addrs)| (peer, addrs.as_slice())) + } + + pub fn previous_relays(&self) -> impl Iterator { + self.previous_relays + .iter() + .map(|(peer, addr, ts)| (peer, addr, ts)) + } + + fn static_dial_in_cooldown(&self, peer_id: &PeerId) -> bool { + self.static_dial_cooldowns + .get(peer_id) + .is_some_and(|deadline| *deadline > Instant::now()) + } + + fn queue_static_dial(&mut self, peer_id: PeerId, addresses: Vec) -> bool { + if addresses.is_empty() + || self.has_direct_connection(&peer_id) + || self.static_dial_in_cooldown(&peer_id) + { + return false; + } + let opts = DialOpts::peer_id(peer_id).addresses(addresses).build(); + self.events.push_back(ToSwarm::Dial { opts }); + true + } + + fn record_previous_relay(&mut self, peer_id: PeerId, address: Multiaddr) { + let max = self.config.max_previous_relays; + if max == 0 { + return; + } + self.previous_relays.retain(|(p, _, _)| *p != peer_id); + if self.previous_relays.len() >= max { + self.previous_relays.pop_front(); + } + self.previous_relays + .push_back((peer_id, address, SystemTime::now())); + } + + fn forget_previous_relay(&mut self, peer_id: &PeerId) { + self.previous_relays.retain(|(p, _, _)| p != peer_id); + } + + fn record_failure(&mut self, peer_id: PeerId) -> Duration { + let attempts = self.failure_counts.entry(peer_id).or_insert(0); + *attempts = attempts.saturating_add(1); + let exponent = attempts.saturating_sub(1).min(20); + let scale = 1u32 << exponent; + self.config + .failure_cooldown + .saturating_mul(scale) + .min(self.config.failure_cooldown_max) + } + + fn clear_failure(&mut self, peer_id: &PeerId) { + self.failure_counts.remove(peer_id); + } + + fn determine_status_from_external_addresses(&mut self) { + let has_public_addr = self + .external_addresses + .iter() + .any(|addr| !addr.is_relayed()); + + let new_status = match has_public_addr { + true => Status::Disable, + false => Status::Enable, + }; + if new_status != self.status { + self.status = new_status; + self.events + .push_back(ToSwarm::GenerateEvent(Event::StatusChanged { + status: new_status, + })); + match new_status { + Status::Enable => self.meet_reservation_target(), + Status::Disable => self.remove_all_reservations(), + } + } + } + + fn is_peer_idle(&self, peer_id: &PeerId) -> bool { + self.connections.iter().any(|((pid, _), info)| { + pid == peer_id + && info.relay_status + == RelayStatus::Supported { + status: ReservationStatus::Idle, + } + }) + } + + fn has_direct_connection(&self, peer_id: &PeerId) -> bool { + self.connections + .iter() + .any(|((pid, _), info)| pid == peer_id && !info.address.is_relayed()) + } + + fn evict_for_static_peer(&mut self, new_static: PeerId) { + let covered = self.covered_peers(); + if covered.contains(&new_static) { + return; + } + let max = self.config.max_reservations.get() as usize; + if covered.len() < max { + return; + } + + if let Some(listener_id) = self + .reservations + .iter() + .find(|(_, (peer_id, _))| !self.static_relays.contains_key(peer_id)) + .map(|(listener_id, _)| *listener_id) + { + self.events + .push_back(ToSwarm::RemoveListener { id: listener_id }); + } + } + + fn select_connection_for_reservation(&mut self, peer_id: PeerId, connection_id: ConnectionId) { + let info = self + .connections + .get_mut(&(peer_id, connection_id)) + .expect("connection is present"); + + if info.relay_status + != (RelayStatus::Supported { + status: ReservationStatus::Idle, + }) + { + return; + } + + let addr_with_peer_id = match info.address.clone().with_p2p(peer_id) { + Ok(addr) => addr, + Err(addr) => { + tracing::warn!(%addr, "address unexpectedly contains a different peer id than the connection; marking relay connection ineligible"); + info.relay_status = RelayStatus::NotSupported; + return; + } + }; + + let opts = ListenOpts::new(addr_with_peer_id.with(Protocol::P2pCircuit)); + let id = opts.listener_id(); + + info.relay_status = RelayStatus::Supported { + status: ReservationStatus::Pending { id }, + }; + self.reservations.insert(id, (peer_id, connection_id)); + self.events.push_back(ToSwarm::ListenOn { opts }); + } + + /// Removes all existing reservations. + fn remove_all_reservations(&mut self) { + let relay_listeners = self + .reservations + .iter() + .map(|(id, (peer_id, conn_id))| (*id, *peer_id, *conn_id)) + .collect::>(); + + for (listener_id, peer_id, connection_id) in relay_listeners { + let Some(connection) = self.connections.get_mut(&(peer_id, connection_id)) else { + continue; + }; + + if !matches!( + connection.relay_status, + RelayStatus::Supported { + status: ReservationStatus::Active { id } | ReservationStatus::Pending { id } + } if id == listener_id + ) { + continue; + } + + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Idle, + }; + + self.events + .push_back(ToSwarm::RemoveListener { id: listener_id }); + } + } + + fn disable_reservation(&mut self, id: ListenerId, failed: bool) { + if self.external_reservations.remove(&id).is_some() { + self.meet_reservation_target(); + return; + } + + let Some((peer_id, connection_id)) = self.reservations.remove(&id) else { + return; + }; + + let Some(address) = self + .connections + .get(&(peer_id, connection_id)) + .filter(|info| { + matches!( + info.relay_status, + RelayStatus::Supported { + status: ReservationStatus::Active { .. } + | ReservationStatus::Pending { .. } + } + ) + }) + .map(|info| info.address.clone()) + else { + self.meet_reservation_target(); + return; + }; + + let blacklist_duration = failed.then(|| self.record_failure(peer_id)); + + let connection = self + .connections + .get_mut(&(peer_id, connection_id)) + .expect("connection is tracked"); + match blacklist_duration { + Some(duration) => { + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Blacklisted, + }; + self.events.push_back(ToSwarm::NotifyHandler { + peer_id, + handler: NotifyHandler::One(connection_id), + event: Either::Left(handler::In::Blacklist { duration }), + }); + } + None => { + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Idle, + }; + } + } + + self.record_previous_relay(peer_id, address); + self.meet_reservation_target(); + } + + fn covered_peers(&self) -> HashSet { + self.reservations + .values() + .map(|(peer_id, _)| *peer_id) + .chain(self.external_reservations.values().copied()) + .collect() + } + + /// Meet the reservation target by selecting connections to establish a reservation. + fn meet_reservation_target(&mut self) { + if self.status == Status::Disable { + return; + } + + let max = self.config.max_reservations.get() as usize; + let covered = self.covered_peers(); + let budget = max.saturating_sub(covered.len()); + if budget == 0 { + return; + } + + let mut static_candidates = BTreeMap::new(); + let mut candidates: BTreeMap<_, ConnectionId> = BTreeMap::new(); + for ((peer_id, connection_id), info) in self.connections.iter() { + if covered.contains(peer_id) { + continue; + } + if info.relay_status + != (RelayStatus::Supported { + status: ReservationStatus::Idle, + }) + { + continue; + } + let bucket = if self.static_relays.contains_key(peer_id) { + &mut static_candidates + } else { + &mut candidates + }; + bucket + .entry(*peer_id) + .and_modify(|existing| *existing = (*existing).min(*connection_id)) + .or_insert(*connection_id); + } + + let selected_candidates: Vec<(PeerId, ConnectionId)> = static_candidates + .into_iter() + .chain(candidates) + .take(budget) + .collect(); + + for (peer_id, connection_id) in selected_candidates { + self.select_connection_for_reservation(peer_id, connection_id); + } + + debug_assert!(self.covered_peers().len() <= max); + } + + fn update_relay_availability(&mut self) { + let has_hop_peer = self + .connections + .values() + .any(|info| matches!(info.relay_status, RelayStatus::Supported { .. })); + + match (has_hop_peer, self.relays_available) { + (true, false) => { + self.relays_available = true; + self.events + .push_back(ToSwarm::GenerateEvent(Event::RelaysAvailable)); + } + (false, true) => { + self.relays_available = false; + self.events + .push_back(ToSwarm::GenerateEvent(Event::NoRelaysAvailable)); + } + _ => {} + } + } +} + +impl NetworkBehaviour for Behaviour { + type ConnectionHandler = Either; + type ToSwarm = Event; + + fn handle_established_inbound_connection( + &mut self, + _connection_id: ConnectionId, + _peer: PeerId, + local_addr: &Multiaddr, + _remote_addr: &Multiaddr, + ) -> Result, ConnectionDenied> { + if local_addr.is_relayed() { + Ok(Either::Right(dummy::ConnectionHandler)) + } else { + Ok(Either::Left(handler::Handler::default())) + } + } + + fn handle_established_outbound_connection( + &mut self, + _connection_id: ConnectionId, + _peer: PeerId, + addr: &Multiaddr, + _role_override: Endpoint, + _port_use: PortUse, + ) -> Result, ConnectionDenied> { + if addr.is_relayed() { + Ok(Either::Right(dummy::ConnectionHandler)) + } else { + Ok(Either::Left(handler::Handler::default())) + } + } + + fn on_swarm_event(&mut self, event: FromSwarm) { + let change = self.external_addresses.on_swarm_event(&event); + + if self.auto_status_change && change { + self.determine_status_from_external_addresses(); + } + + match event { + FromSwarm::ConnectionEstablished(ConnectionEstablished { + peer_id, + endpoint, + connection_id, + .. + }) => { + let remote_addr = endpoint.get_remote_address().clone(); + + let mut connection = Connection { + address: remote_addr, + relay_status: RelayStatus::Pending, + }; + + connection.disqualify_connection_if_relayed(); + + self.connections + .insert((peer_id, connection_id), connection); + + if self.static_relays.contains_key(&peer_id) { + self.static_dial_cooldowns.remove(&peer_id); + } + } + FromSwarm::ConnectionClosed(ConnectionClosed { + peer_id, + connection_id, + .. + }) => { + let Some(connection) = self.connections.remove(&(peer_id, connection_id)) else { + return; + }; + + if !self.connections.keys().any(|(pid, _)| *pid == peer_id) { + self.clear_failure(&peer_id); + } + + let had_reservation = matches!( + connection.relay_status, + RelayStatus::Supported { + status: ReservationStatus::Active { .. } + | ReservationStatus::Pending { .. } + | ReservationStatus::Blacklisted + } + ); + + if let RelayStatus::Supported { + status: ReservationStatus::Active { id } | ReservationStatus::Pending { id }, + } = connection.relay_status + { + self.reservations.remove(&id); + self.meet_reservation_target(); + } + + if had_reservation { + self.record_previous_relay(peer_id, connection.address); + } + + if let Some(addresses) = self.static_relays.get(&peer_id).cloned() { + self.queue_static_dial(peer_id, addresses); + } + + self.update_relay_availability(); + } + FromSwarm::AddressChange(AddressChange { + peer_id, + connection_id, + old: _, + new, + }) => { + let Some(connection) = self.connections.get_mut(&(peer_id, connection_id)) else { + return; + }; + + let new_addr = new.get_remote_address(); + + connection.address = new_addr.clone(); + } + FromSwarm::NewListenAddr(NewListenAddr { listener_id, addr }) => { + if !addr.is_relayed() { + return; + } + + if let Some((peer_id, connection_id)) = self.reservations.get(&listener_id).copied() + { + let Some(connection) = self.connections.get_mut(&(peer_id, connection_id)) + else { + return; + }; + + if matches!( + connection.relay_status, + RelayStatus::Supported { + status: ReservationStatus::Pending { id } + } if id == listener_id + ) { + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Active { id: listener_id }, + }; + self.forget_previous_relay(&peer_id); + self.clear_failure(&peer_id); + } + return; + } + + if let Some(relay_peer_id) = addr.relay_peer_id() { + self.external_reservations + .insert(listener_id, relay_peer_id); + } + } + FromSwarm::ExpiredListenAddr(ExpiredListenAddr { listener_id, .. }) => { + self.disable_reservation(listener_id, false); + } + FromSwarm::ListenerError(ListenerError { listener_id, .. }) => { + self.disable_reservation(listener_id, true); + } + FromSwarm::ListenerClosed(ListenerClosed { + listener_id, + reason, + .. + }) => { + self.disable_reservation(listener_id, reason.is_err()); + } + FromSwarm::DialFailure(DialFailure { + peer_id: Some(peer_id), + error, + .. + }) if self.static_relays.contains_key(&peer_id) => { + tracing::warn!(%peer_id, %error, "dial to static relay failed"); + self.static_dial_cooldowns + .insert(peer_id, Instant::now() + self.config.failure_cooldown); + } + _ => {} + } + } + + fn on_connection_handler_event( + &mut self, + peer_id: PeerId, + connection_id: ConnectionId, + event: THandlerOutEvent, + ) { + let Either::Left(event) = event; + + let Some(connection) = self.connections.get_mut(&(peer_id, connection_id)) else { + return; + }; + + match event { + Out::Supported => { + if matches!( + connection.relay_status, + RelayStatus::Pending | RelayStatus::NotSupported + ) { + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Idle, + }; + if self.static_relays.contains_key(&peer_id) { + self.evict_for_static_peer(peer_id); + } + self.meet_reservation_target(); + self.update_relay_availability(); + } + } + Out::Unsupported => { + let drop_listener = match connection.relay_status { + RelayStatus::Supported { + status: ReservationStatus::Pending { id } | ReservationStatus::Active { id }, + } => Some(id), + _ => None, + }; + let lost_address = drop_listener.map(|_| connection.address.clone()); + connection.relay_status = RelayStatus::NotSupported; + if let Some(id) = drop_listener { + self.reservations.remove(&id); + self.events.push_back(ToSwarm::RemoveListener { id }); + self.meet_reservation_target(); + } + if let Some(address) = lost_address { + self.record_previous_relay(peer_id, address); + } + self.update_relay_availability(); + } + Out::BlacklistExpired => { + if matches!( + connection.relay_status, + RelayStatus::Supported { + status: ReservationStatus::Blacklisted + } + ) { + connection.relay_status = RelayStatus::Supported { + status: ReservationStatus::Idle, + }; + self.meet_reservation_target(); + } + } + } + } + + fn poll( + &mut self, + cx: &mut Context<'_>, + ) -> Poll>> { + if let Some(event) = self.events.pop_front() { + return Poll::Ready(event); + } + + self.waker = Some(cx.waker().clone()); + + Poll::Pending + } +} diff --git a/src/behaviour/autorelay/handler.rs b/src/behaviour/autorelay/handler.rs new file mode 100644 index 0000000..169e1ff --- /dev/null +++ b/src/behaviour/autorelay/handler.rs @@ -0,0 +1,129 @@ +use std::{ + collections::VecDeque, + task::{Context, Poll}, + time::Duration, +}; + +use crate::prelude::swarm::handler::ConnectionEvent; +use crate::prelude::swarm::{ + ConnectionHandler, ConnectionHandlerEvent, SubstreamProtocol, SupportedProtocols, +}; +use crate::prelude::transport::upgrade::DeniedUpgrade; +use futures::FutureExt; +use futures_timer::Delay; +use libp2p::relay::HOP_PROTOCOL_NAME; + +#[derive(Default, Debug)] +pub struct Handler { + events: VecDeque< + ConnectionHandlerEvent< + ::OutboundProtocol, + ::OutboundOpenInfo, + ::ToBehaviour, + >, + >, + + supported: bool, + + supported_protocol: SupportedProtocols, + + blacklist_timer: Option, +} + +#[derive(Debug, Copy, Clone)] +pub enum In { + Blacklist { duration: Duration }, +} + +#[derive(Debug, Copy, Clone)] +pub enum Out { + Supported, + Unsupported, + BlacklistExpired, +} + +#[allow(deprecated)] +impl ConnectionHandler for Handler { + type FromBehaviour = In; + type ToBehaviour = Out; + type InboundProtocol = DeniedUpgrade; + type OutboundProtocol = DeniedUpgrade; + type InboundOpenInfo = (); + type OutboundOpenInfo = (); + + fn listen_protocol(&self) -> SubstreamProtocol { + SubstreamProtocol::new(DeniedUpgrade, ()) + } + + fn connection_keep_alive(&self) -> bool { + false + } + + fn on_behaviour_event(&mut self, event: Self::FromBehaviour) { + match event { + In::Blacklist { duration } => { + self.blacklist_timer = Some(Delay::new(duration)); + } + } + } + + fn on_connection_event( + &mut self, + event: ConnectionEvent< + Self::InboundProtocol, + Self::OutboundProtocol, + Self::InboundOpenInfo, + Self::OutboundOpenInfo, + >, + ) { + if let ConnectionEvent::RemoteProtocolsChange(protocol) = event { + let change = self.supported_protocol.on_protocols_change(protocol); + if change { + let valid = self + .supported_protocol + .iter() + .any(|proto| HOP_PROTOCOL_NAME.eq(proto)); + + match (valid, self.supported) { + (true, false) => { + self.supported = true; + self.events + .push_back(ConnectionHandlerEvent::NotifyBehaviour(Out::Supported)); + } + (false, true) => { + self.supported = false; + self.blacklist_timer = None; + self.events + .push_back(ConnectionHandlerEvent::NotifyBehaviour(Out::Unsupported)); + } + (true, true) => {} + _ => {} + } + } + } + } + + fn poll( + &mut self, + cx: &mut Context<'_>, + ) -> Poll< + ConnectionHandlerEvent, + > { + if let Some(event) = self.events.pop_front() { + return Poll::Ready(event); + } + + if let Some(timer) = self.blacklist_timer.as_mut() + && timer.poll_unpin(cx).is_ready() + { + self.blacklist_timer = None; + if self.supported { + return Poll::Ready(ConnectionHandlerEvent::NotifyBehaviour( + Out::BlacklistExpired, + )); + } + } + + Poll::Pending + } +} diff --git a/src/behaviour/request_response/codec.rs b/src/behaviour/request_response/codec.rs index fc775de..ec04907 100644 --- a/src/behaviour/request_response/codec.rs +++ b/src/behaviour/request_response/codec.rs @@ -1,4 +1,3 @@ -use async_trait::async_trait; use bytes::Bytes; use futures::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; use libp2p::StreamProtocol; @@ -18,7 +17,6 @@ impl Codec { } } -#[async_trait] impl libp2p::request_response::Codec for Codec { type Protocol = StreamProtocol; type Request = Bytes; diff --git a/src/builder.rs b/src/builder.rs index 82294b3..3e8d3f3 100644 --- a/src/builder.rs +++ b/src/builder.rs @@ -124,6 +124,9 @@ pub(crate) struct Config { pub autonat_v2_client_config: Box AutonatV2ClientConfig>, #[cfg(feature = "relay")] pub relay_server_config: Box RelayServerConfig>, + #[cfg(feature = "relay")] + pub autorelay_config: + Box behaviour::autorelay::Config>, #[cfg(feature = "identify")] pub identify_config: (String, Box IdentifyConfig>), #[cfg(feature = "request-response")] @@ -156,6 +159,8 @@ impl Default for Config { autonat_v2_client_config: Box::new(|config| config), #[cfg(feature = "relay")] relay_server_config: Box::new(|config| config), + #[cfg(feature = "relay")] + autorelay_config: Box::new(|config| config), #[cfg(feature = "identify")] identify_config: (String::from("/ipfs/id"), Box::new(|config| config)), #[cfg(feature = "request-response")] @@ -180,6 +185,8 @@ pub(crate) struct Protocols { pub(crate) relay_client: bool, #[cfg(feature = "relay")] pub(crate) relay_server: bool, + #[cfg(feature = "relay")] + pub(crate) autorelay: bool, #[cfg(feature = "dcutr")] #[cfg(not(target_arch = "wasm32"))] pub(crate) dcutr: bool, @@ -441,6 +448,23 @@ where self } + /// Enable autorelay + #[cfg(feature = "relay")] + pub fn with_autorelay(self) -> Self { + self.with_autorelay_with_config(|config| config) + } + + /// Enable autorelay + #[cfg(feature = "relay")] + pub fn with_autorelay_with_config(mut self, f: F) -> Self + where + F: FnOnce(behaviour::autorelay::Config) -> behaviour::autorelay::Config + 'static, + { + self.config.autorelay_config = Box::new(f); + self.protocols.autorelay = true; + self + } + /// Enables DCuTR #[cfg(all(feature = "relay", feature = "dcutr"))] #[cfg(not(target_arch = "wasm32"))] diff --git a/src/builder/transport.rs b/src/builder/transport.rs index 8564281..e5a6a74 100644 --- a/src/builder/transport.rs +++ b/src/builder/transport.rs @@ -156,6 +156,8 @@ pub enum DnsResolver { /// Cloudflare DNS Resolver #[default] Cloudflare, + /// Quad9 DNS Resolver + Quad9, /// Local DNS Resolver Local, /// No DNS Resolver @@ -167,12 +169,22 @@ pub enum DnsResolver { impl From for (ResolverConfig, ResolverOpts) { fn from(value: DnsResolver) -> Self { match value { - DnsResolver::Google => (ResolverConfig::google(), Default::default()), - DnsResolver::Cloudflare => (ResolverConfig::cloudflare(), Default::default()), + DnsResolver::Google => ( + ResolverConfig::udp_and_tcp(&hickory_resolver::config::GOOGLE), + Default::default(), + ), + DnsResolver::Cloudflare => ( + ResolverConfig::udp_and_tcp(&hickory_resolver::config::CLOUDFLARE), + Default::default(), + ), + DnsResolver::Quad9 => ( + ResolverConfig::udp_and_tcp(&hickory_resolver::config::QUAD9), + Default::default(), + ), DnsResolver::Local => { hickory_resolver::system_conf::read_system_conf().unwrap_or_default() } - DnsResolver::None => (ResolverConfig::new(), Default::default()), + DnsResolver::None => (ResolverConfig::default(), Default::default()), } } } diff --git a/src/handle.rs b/src/handle.rs index 359fed0..187e0d1 100644 --- a/src/handle.rs +++ b/src/handle.rs @@ -8,6 +8,11 @@ pub(crate) mod floodsub; #[cfg(feature = "gossipsub")] pub(crate) mod gossipsub; mod peer_store; +#[cfg(feature = "relay")] +mod relay; +#[cfg(not(target_arch = "wasm32"))] +#[cfg(feature = "relay")] +mod relay_server; #[cfg(feature = "rendezvous")] pub(crate) mod rendezvous; #[cfg(feature = "request-response")] @@ -28,6 +33,11 @@ use crate::handle::floodsub::ConnexaFloodsub; #[cfg(feature = "gossipsub")] use crate::handle::gossipsub::ConnexaGossipsub; use crate::handle::peer_store::ConnexaPeerstore; +#[cfg(feature = "relay")] +use crate::handle::relay::ConnexaRelay; +#[cfg(feature = "relay")] +#[cfg(not(target_arch = "wasm32"))] +use crate::handle::relay_server::ConnexaRelayServer; #[cfg(feature = "rendezvous")] use crate::handle::rendezvous::ConnexaRendezvous; #[cfg(feature = "request-response")] @@ -137,6 +147,19 @@ where ConnexaRendezvous::new(self) } + /// Returns a handle for relay functions + #[cfg(feature = "relay")] + pub fn relay(&self) -> ConnexaRelay<'_, T, K> { + ConnexaRelay::new(self) + } + + /// Returns a handle for relay server functions + #[cfg(not(target_arch = "wasm32"))] + #[cfg(feature = "relay")] + pub fn relay_server(&self) -> ConnexaRelayServer<'_, T, K> { + ConnexaRelayServer::new(self) + } + /// Returns a handle to manage peer whitelist functionality pub fn whitelist(&self) -> ConnexaWhitelist<'_, T, K> { ConnexaWhitelist::new(self) diff --git a/src/handle/relay.rs b/src/handle/relay.rs new file mode 100644 index 0000000..fd20792 --- /dev/null +++ b/src/handle/relay.rs @@ -0,0 +1,95 @@ +use crate::handle::Connexa; +use crate::prelude::{Multiaddr, PeerId}; +use crate::types::AutoRelayCommand; +use futures::channel::oneshot; +use std::io; + +pub struct ConnexaRelay<'a, T, K> { + connexa: &'a Connexa, +} + +impl<'a, T, K> ConnexaRelay<'a, T, K> +where + T: Send + Sync + 'static, +{ + pub(crate) fn new(connexa: &'a Connexa) -> Self { + Self { connexa } + } + + pub async fn add_static_relay(&self, peer_id: PeerId, addr: Multiaddr) -> io::Result { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send( + AutoRelayCommand::AddStaticRelay { + peer_id, + relay_addr: addr, + resp: tx, + } + .into(), + ) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn remove_static_relay(&self, peer_id: PeerId) -> io::Result { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::RemoveStaticRelay { peer_id, resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn list_static_relays(&self) -> io::Result)>> { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::ListStaticRelays { resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn get_static_relay(&self, peer_id: PeerId) -> io::Result> { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::GetStaticRelay { peer_id, resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn enable_auto_relay(&self) -> io::Result<()> { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::EnableAutoRelay { resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn disable_auto_relay(&self) -> io::Result<()> { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::DisableAutoRelay { resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } + + pub async fn disable_relays(&self) -> io::Result<()> { + let (tx, rx) = oneshot::channel(); + self.connexa + .to_task + .clone() + .send(AutoRelayCommand::DisableRelays { resp: tx }.into()) + .await?; + rx.await.map_err(io::Error::other)? + } +} diff --git a/src/handle/relay_server.rs b/src/handle/relay_server.rs new file mode 100644 index 0000000..78d1fcb --- /dev/null +++ b/src/handle/relay_server.rs @@ -0,0 +1,31 @@ +use crate::handle::Connexa; +use crate::types::RelayServerCommand; +use libp2p::relay::Status as RelayServerStatus; + +#[derive(Copy, Clone)] +pub struct ConnexaRelayServer<'a, T = (), K = crate::keystore::store::memory::MemoryKeystore> { + connexa: &'a Connexa, +} + +impl<'a, T, K> ConnexaRelayServer<'a, T, K> +where + T: Send + Sync + 'static, +{ + pub(crate) fn new(connexa: &'a Connexa) -> Self { + Self { connexa } + } + + pub async fn change_status( + &self, + status: impl Into>, + ) -> std::io::Result<()> { + let (tx, rx) = futures::channel::oneshot::channel(); + let status = status.into(); + self.connexa + .to_task + .clone() + .send(RelayServerCommand::StatusChanged { status, resp: tx }.into()) + .await?; + rx.await.map_err(std::io::Error::other)? + } +} diff --git a/src/multiaddr_ext.rs b/src/multiaddr_ext.rs index b576f02..f371eb9 100644 --- a/src/multiaddr_ext.rs +++ b/src/multiaddr_ext.rs @@ -1,3 +1,4 @@ +use crate::prelude::PeerId; use libp2p::Multiaddr; use libp2p::multiaddr::Protocol; @@ -11,6 +12,7 @@ pub trait MultiaddrExt { fn is_private(&self) -> bool; fn is_unspecified(&self) -> bool; + fn relay_peer_id(&self) -> Option; } impl MultiaddrExt for Multiaddr { @@ -47,6 +49,18 @@ impl MultiaddrExt for Multiaddr { _ => false, }) } + + fn relay_peer_id(&self) -> Option { + let mut last_p2p = None; + for proto in self.iter() { + match proto { + Protocol::P2p(peer) => last_p2p = Some(peer), + Protocol::P2pCircuit => return last_p2p, + _ => {} + } + } + None + } } #[cfg(test)] diff --git a/src/task.rs b/src/task.rs index b285702..18483d0 100644 --- a/src/task.rs +++ b/src/task.rs @@ -52,7 +52,7 @@ use futures::channel::{mpsc, oneshot}; use futures::future::BoxFuture; use futures::{FutureExt, StreamExt}; use futures_timer::Delay; -use indexmap::IndexMap; +use indexmap::{IndexMap, IndexSet}; #[cfg(feature = "gossipsub")] use libp2p::gossipsub::{MessageAcceptance, MessageId}; #[cfg(feature = "kad")] @@ -573,6 +573,10 @@ where Command::Autonat(autonat_command) => self.process_autonat_v1_command(autonat_command), #[cfg(feature = "kad")] Command::Dht(dht_command) => self.process_kademlia_command(dht_command), + #[cfg(feature = "relay")] + Command::AutoRelay(autorelay_command) => { + self.process_autorelay_commands(autorelay_command) + } #[cfg(feature = "stream")] Command::Stream(stream_command) => self.process_stream_command(stream_command), #[cfg(feature = "request-response")] @@ -583,6 +587,11 @@ where Command::Rendezvous(rendezvous_command) => { self.process_rendezvous_command(rendezvous_command) } + #[cfg(not(target_arch = "wasm32"))] + #[cfg(feature = "relay")] + Command::RelayServer(relay_server_command) => { + self.process_relay_server_command(relay_server_command) + } Command::Custom(custom_command) => { (self.custom_task_callback)( swarm, diff --git a/src/task/ping.rs b/src/task/ping.rs index 3222b40..e915f83 100644 --- a/src/task/ping.rs +++ b/src/task/ping.rs @@ -20,6 +20,18 @@ where match result { Ok(duration) => { tracing::info!("ping to {} at {} took {:?}", peer, connection, duration); + + // #[cfg(feature = "relay")] + // if let Some(autorelay) = self + // .swarm + // .as_mut() + // .expect("swarm valid") + // .behaviour_mut() + // .autorelay + // .as_mut() + // { + // autorelay.set_peer_ping(peer, connection, duration); + // } } Err(e) => { // TODO: Possibly disconnect peer since if there is an error? diff --git a/src/task/relay.rs b/src/task/relay.rs index 219b797..913c47d 100644 --- a/src/task/relay.rs +++ b/src/task/relay.rs @@ -1,9 +1,16 @@ +use crate::behaviour::autorelay; use crate::behaviour::peer_store::store::Store; use crate::task::ConnexaTask; +use crate::types::AutoRelayCommand; +#[cfg(not(target_arch = "wasm32"))] +use crate::types::RelayServerCommand; use libp2p::relay::{Event as RelayServerEvent, client::Event as RelayClientEvent}; use libp2p::swarm::NetworkBehaviour; use std::fmt::Debug; +#[allow(dead_code)] +pub const RELAY_NAMESPACE: &[u8] = b"/libp2p/relay"; + impl ConnexaTask where X: Default + Send + 'static, @@ -11,6 +18,89 @@ where C::ToSwarm: Debug, S: Store, { + pub fn process_autorelay_commands(&mut self, command: AutoRelayCommand) { + let swarm = self.swarm.as_mut().expect("swarm is still valid"); + match command { + AutoRelayCommand::AddStaticRelay { + peer_id, + relay_addr, + resp, + } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + let _ = resp.send(Ok(autorelay.add_static_relay(peer_id, relay_addr))); + } + AutoRelayCommand::RemoveStaticRelay { peer_id, resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + let _ = resp.send(Ok(autorelay.remove_static_relay(&peer_id))); + } + AutoRelayCommand::ListStaticRelays { resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + let list = autorelay + .static_relays() + .map(|(peer_id, addr)| (*peer_id, addr.to_vec())) + .collect::>(); + let _ = resp.send(Ok(list)); + } + AutoRelayCommand::GetStaticRelay { peer_id, resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + let addr = autorelay + .static_relays() + .find(|(p, _)| **p == peer_id) + .map(|(_, addr)| addr.to_vec()) + .ok_or_else(|| { + std::io::Error::new(std::io::ErrorKind::NotFound, "static relay not found") + }); + let _ = resp.send(addr); + } + AutoRelayCommand::EnableAutoRelay { resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + autorelay.set_status(Some(autorelay::Status::Enable)); + + let _ = resp.send(Ok(())); + } + AutoRelayCommand::DisableAutoRelay { resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + autorelay.set_status(Some(autorelay::Status::Disable)); + + let _ = resp.send(Ok(())); + } + AutoRelayCommand::DisableRelays { resp } => { + let Some(autorelay) = swarm.behaviour_mut().autorelay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other("autorelay is not enabled"))); + return; + }; + + autorelay.remove_all_reservations(); + + let _ = resp.send(Ok(())); + } + } + } + pub fn process_relay_client_event(&mut self, event: RelayClientEvent) { match event { RelayClientEvent::ReservationReqAccepted { @@ -32,6 +122,7 @@ where } } + #[cfg(not(target_arch = "wasm32"))] pub fn process_relay_server_event(&mut self, event: RelayServerEvent) { match event { RelayServerEvent::ReservationReqAccepted { @@ -69,7 +160,37 @@ where } => { tracing::warn!(%src_peer_id, %dst_peer_id, ?error, "relay server circuit closed"); } + RelayServerEvent::StatusChanged { status } => { + tracing::info!(?status, "relay server status changed"); + } _ => {} } } } + +#[cfg(not(target_arch = "wasm32"))] +impl ConnexaTask +where + X: Default + Send + 'static, + C: Send, + C::ToSwarm: Debug, + S: Store, +{ + pub fn process_relay_server_command(&mut self, command: RelayServerCommand) { + let swarm = self.swarm.as_mut().expect("swarm is active"); + match command { + RelayServerCommand::StatusChanged { status, resp } => { + let Some(relay) = swarm.behaviour_mut().relay.as_mut() else { + let _ = resp.send(Err(std::io::Error::other( + "relay server protocol is not enabled", + ))); + return; + }; + + relay.set_status(status); + + let _ = resp.send(Ok(())); + } + } + } +} diff --git a/src/task/swarm.rs b/src/task/swarm.rs index f596b77..f75309b 100644 --- a/src/task/swarm.rs +++ b/src/task/swarm.rs @@ -291,6 +291,7 @@ where }; match event { + #[cfg(not(target_arch = "wasm32"))] #[cfg(feature = "relay")] BehaviourEvent::Relay(event) => self.process_relay_server_event(event), #[cfg(feature = "relay")] diff --git a/src/task/upnp.rs b/src/task/upnp.rs index 838c38f..5ed3835 100644 --- a/src/task/upnp.rs +++ b/src/task/upnp.rs @@ -13,11 +13,21 @@ where { pub fn process_upnp_event(&mut self, event: UpnpEvent) { match event { - UpnpEvent::NewExternalAddr(addr) => { - tracing::info!(?addr, "upnp external address discovered"); + UpnpEvent::NewExternalAddr { + external_addr, + local_addr, + } => { + tracing::info!( + ?external_addr, + ?local_addr, + "upnp external address discovered" + ); } - UpnpEvent::ExpiredExternalAddr(addr) => { - tracing::info!(?addr, "upnp external address expired"); + UpnpEvent::ExpiredExternalAddr { + external_addr, + local_addr, + } => { + tracing::info!(?external_addr, ?local_addr, "upnp external address expired"); } UpnpEvent::GatewayNotFound => { tracing::warn!("upnp gateway not found"); diff --git a/src/types.rs b/src/types.rs index d9d868f..5259924 100644 --- a/src/types.rs +++ b/src/types.rs @@ -46,6 +46,11 @@ pub enum Command { Rendezvous(RendezvousCommand), #[cfg(feature = "autonat")] Autonat(AutonatCommand), + #[cfg(feature = "relay")] + AutoRelay(AutoRelayCommand), + #[cfg(not(target_arch = "wasm32"))] + #[cfg(feature = "relay")] + RelayServer(RelayServerCommand), Whitelist(WhitelistCommand), Blacklist(BlacklistCommand), ConnectionLimits(ConnectionLimitsCommand), @@ -108,6 +113,21 @@ impl From for Command { } } +#[cfg(feature = "relay")] +impl From for Command { + fn from(cmd: AutoRelayCommand) -> Self { + Command::AutoRelay(cmd) + } +} + +#[cfg(not(target_arch = "wasm32"))] +#[cfg(feature = "relay")] +impl From for Command { + fn from(cmd: RelayServerCommand) -> Self { + Command::RelayServer(cmd) + } +} + impl From for Command { fn from(cmd: WhitelistCommand) -> Self { Command::Whitelist(cmd) @@ -432,6 +452,36 @@ pub enum DHTCommand { }, } +#[cfg(feature = "relay")] +#[derive(Debug)] +pub enum AutoRelayCommand { + AddStaticRelay { + peer_id: PeerId, + relay_addr: Multiaddr, + resp: oneshot::Sender>, + }, + RemoveStaticRelay { + peer_id: PeerId, + resp: oneshot::Sender>, + }, + DisableRelays { + resp: oneshot::Sender>, + }, + ListStaticRelays { + resp: oneshot::Sender)>>>, + }, + GetStaticRelay { + peer_id: PeerId, + resp: oneshot::Sender>>, + }, + EnableAutoRelay { + resp: oneshot::Sender>, + }, + DisableAutoRelay { + resp: oneshot::Sender>, + }, +} + #[cfg(feature = "request-response")] type ResponseStream = BoxStream<'static, (PeerId, ConnexaResult)>; #[cfg(feature = "request-response")] @@ -505,6 +555,16 @@ pub enum RendezvousCommand { }, } +#[cfg(not(target_arch = "wasm32"))] +#[cfg(feature = "relay")] +#[derive(Debug)] +pub enum RelayServerCommand { + StatusChanged { + status: Option, + resp: oneshot::Sender>, + }, +} + #[cfg(feature = "kad")] #[derive(Clone, Debug)] pub enum DHTEvent {