From d21a521cbe03f034bcdd99497e91f9edd9c2c1d8 Mon Sep 17 00:00:00 2001 From: Oscar Beaumont Date: Thu, 9 Mar 2023 11:27:45 +0800 Subject: [PATCH] [ENG-407] Spacedrop backend (#598) * format Rust * Spacedrop a string * Praise thee Clippy, lord of the Rust * add protobuf to Mac and Linux CI * plz GH Actions have Chocolatey --- .github/scripts/setup-system.ps1 | 3 + .github/scripts/setup-system.sh | 5 +- core/src/lib.rs | 3 +- core/src/library/mod.rs | 1 + core/src/location/manager/mod.rs | 4 +- core/src/p2p/p2p_manager.rs | 97 ++++- crates/p2p/src/event.rs | 56 +-- crates/p2p/src/manager.rs | 208 +++++---- crates/p2p/src/manager_stream.rs | 342 +++++++-------- crates/p2p/src/mdns.rs | 386 ++++++++--------- crates/p2p/src/peer.rs | 44 +- crates/p2p/src/spaceblock/mod.rs | 54 +-- crates/p2p/src/spacetime/behaviour.rs | 472 +++++++++++---------- crates/p2p/src/spacetime/connection.rs | 207 ++++----- crates/p2p/src/spacetime/libp2p.rs | 6 +- crates/p2p/src/spacetime/message.rs | 8 +- crates/p2p/src/spacetime/proto_inbound.rs | 53 +-- crates/p2p/src/spacetime/proto_outbound.rs | 64 +-- crates/p2p/src/spacetime/stream.rs | 58 +-- crates/p2p/src/utils/async_fn.rs | 14 +- crates/p2p/src/utils/keypair.rs | 60 +-- crates/p2p/src/utils/metadata.rs | 8 +- crates/p2p/src/utils/multiaddr.rs | 62 +-- 23 files changed, 1149 insertions(+), 1066 deletions(-) diff --git a/.github/scripts/setup-system.ps1 b/.github/scripts/setup-system.ps1 index 57c00d31e..3655e5150 100644 --- a/.github/scripts/setup-system.ps1 +++ b/.github/scripts/setup-system.ps1 @@ -133,6 +133,9 @@ if ($env:CI -eq $True) { } +Write-Host +Write-Host "Install protobuf compiler..." -ForegroundColor Yellow +choco install protoc Write-Host Write-Host "Downloading the latest ffmpeg build..." -ForegroundColor Yellow diff --git a/.github/scripts/setup-system.sh b/.github/scripts/setup-system.sh index 8f89c6ad8..85bc369cf 100755 --- a/.github/scripts/setup-system.sh +++ b/.github/scripts/setup-system.sh @@ -95,9 +95,10 @@ if [[ "$OSTYPE" == "linux-gnu"* ]]; then DEBIAN_TAURI_DEPS="libwebkit2gtk-4.0-dev build-essential curl wget libssl-dev libgtk-3-dev libayatana-appindicator3-dev librsvg2-dev" # Tauri dependencies DEBIAN_FFMPEG_DEPS="libavcodec-dev libavdevice-dev libavfilter-dev libavformat-dev libavutil-dev libswscale-dev libswresample-dev ffmpeg" # FFmpeg dependencies DEBIAN_BINDGEN_DEPS="pkg-config clang" # Bindgen dependencies - it's used by a dependency of Spacedrive + PROTOBUF="protobuf-compiler" # Protobuf compiler sudo apt-get -y update - sudo apt-get -y install ${SPACEDRIVE_CUSTOM_APT_FLAGS:-} $DEBIAN_TAURI_DEPS $DEBIAN_FFMPEG_DEPS $DEBIAN_BINDGEN_DEPS + sudo apt-get -y install ${SPACEDRIVE_CUSTOM_APT_FLAGS:-} $DEBIAN_TAURI_DEPS $DEBIAN_FFMPEG_DEPS $DEBIAN_BINDGEN_DEPS $PROTOBUF elif command -v pacman >/dev/null; then echo "Detected pacman!" echo "Installing dependencies with pacman..." @@ -145,6 +146,8 @@ elif [[ "$OSTYPE" == "darwin"* ]]; then brew unlink -q ffmpeg || true brew install -q "spacedriveapp/deps/ffmpeg@$FFMPEG_VERSION" + brew install protobuf + echo "FFmpeg version $FFMPEG_VERSION has been installed and is now being used on your system." fi else diff --git a/core/src/lib.rs b/core/src/lib.rs index 5af6fc5e6..1549a208b 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -41,6 +41,7 @@ pub struct Node { config: Arc, library_manager: Arc, jobs: Arc, + #[allow(unused)] // TODO: Remove `allow(unused)` once integrated p2p: Arc, event_bus: (broadcast::Sender, broadcast::Receiver), secure_temp_keystore: Arc, @@ -173,7 +174,7 @@ impl Node { } }); - let p2p = P2PManager::new(config.clone(), library_manager.clone()).await; + let p2p = P2PManager::new(config.clone()).await; let router = api::mount(); let node = Node { diff --git a/core/src/library/mod.rs b/core/src/library/mod.rs index a93237e0f..396a39c0d 100644 --- a/core/src/library/mod.rs +++ b/core/src/library/mod.rs @@ -1,4 +1,5 @@ mod config; +#[allow(clippy::module_inception)] mod library; mod manager; diff --git a/core/src/location/manager/mod.rs b/core/src/location/manager/mod.rs index c103d81d0..bd2c5a702 100644 --- a/core/src/location/manager/mod.rs +++ b/core/src/location/manager/mod.rs @@ -164,7 +164,7 @@ impl LocationManager { self.location_management_tx .send(LocationManagementMessage { location_id, - library: library, + library, action, response_tx: tx, }) @@ -192,7 +192,7 @@ impl LocationManager { self.watcher_management_tx .send(WatcherManagementMessage { location_id, - library: library, + library, action, response_tx: tx, }) diff --git a/core/src/p2p/p2p_manager.rs b/core/src/p2p/p2p_manager.rs index 9e86c4d25..590180e46 100644 --- a/core/src/p2p/p2p_manager.rs +++ b/core/src/p2p/p2p_manager.rs @@ -1,13 +1,15 @@ -use std::sync::Arc; +use std::{path::PathBuf, sync::Arc, time::Instant}; -use sd_p2p::{Event, Manager}; +use sd_p2p::{Event, Manager, PeerId}; use sd_sync::CRDTOperation; -use tokio::io::AsyncReadExt; +use tokio::{ + fs::File, + io::{AsyncReadExt, AsyncWriteExt, BufReader}, +}; use tracing::{debug, error, info}; use uuid::Uuid; use crate::{ - library::LibraryManager, node::NodeConfigManager, p2p::{OperatingSystem, SPACEDRIVE_APP_ID}, }; @@ -19,10 +21,7 @@ pub struct P2PManager { } impl P2PManager { - pub async fn new( - node_config: Arc, - library_manager: Arc, - ) -> Arc { + pub async fn new(node_config: Arc) -> Arc { let (config, keypair) = { let config = node_config.get().await; ( @@ -33,7 +32,7 @@ impl P2PManager { email: config.p2p_email.clone(), img_url: config.p2p_img_url.clone(), }, - config.keypair.clone(), + config.keypair, ) }; // TODO: Update this throughout the application lifecycle @@ -71,12 +70,30 @@ impl P2PManager { debug!("Received ping from peer '{}'", event.peer_id); } Header::Spacedrop => { - todo!(); + let file_length = event.stream.read_u8().await.unwrap(); + + // TODO: Ask the user if they wanna reject/accept it + + info!("Received Spacedrop from peer '{}' with file length '{file_length}'", event.peer_id); + + let mut s = String::new(); + event.stream.read_to_string(&mut s).await.unwrap(); + + // let mut buf = Vec::with_capacity(file_length as usize); // TODO: DOS attack + // loop { + // let n = event.stream.read(&mut buf).await.unwrap();; + // } + + // // TODO: Store this to a file on disk. + println!( + "Recieved file content '{}' through Spacedrop!", + s // String::from_utf8(buf).unwrap() + ); } Header::Sync(library_id) => { let buf_len = event.stream.read_u8().await.unwrap(); - let mut buf = Vec::with_capacity(buf_len as usize); // TODO: Designed for easily being able to be DOS the current Node + let mut buf = vec![0; buf_len as usize]; // TODO: Designed for easily being able to be DOS the current Node event.stream.read_exact(&mut buf).await.unwrap(); let mut buf: &[u8] = &buf; @@ -84,7 +101,7 @@ impl P2PManager { rmp_serde::from_read(&mut buf).unwrap(); // TODO: Handle this @Brendan - println!("Receieved sync events for library '{library_id}': {output:?}"); + println!("Received sync events for library '{library_id}': {output:?}"); // TODO(@Oscar): Remember we can't do a response here cause it's a broadcast. Encode that into type system! } @@ -115,9 +132,31 @@ impl P2PManager { } }); + // TODO: Probs remove this + tokio::spawn({ + let this = this.clone(); + async move { + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + let mut connected = this + .manager + .get_connected_peers() + .await + .unwrap() + .into_iter(); + if let Some(peer_id) = connected.next() { + info!("Starting Spacedrop to peer '{}'", peer_id); + this.big_bad_spacedrop(peer_id, PathBuf::from("./demo.txt")) + .await; + } else { + info!("No clients found so skipping Spacedrop demo!"); + } + } + }); + this } + #[allow(unused)] // TODO: Remove `allow(unused)` once integrated pub async fn broadcast_sync_events(&self, library_id: Uuid, event: Vec) { let mut head_buf = Header::Sync(library_id).to_bytes(); let mut buf = rmp_serde::to_vec_named(&event).unwrap(); // TODO: Error handling @@ -130,4 +169,38 @@ impl P2PManager { pub async fn ping(&self) { self.manager.broadcast(Header::Ping.to_bytes()).await; } + + pub async fn big_bad_spacedrop(&self, peer_id: PeerId, path: PathBuf) { + let mut stream = self.manager.stream(peer_id).await.unwrap(); // TODO: handle providing incorrect peer id + + let file = File::open(path).await.unwrap(); + let file_length = file.metadata().await.unwrap().len(); + let mut reader = BufReader::new(file); + + stream + .write_all(&Header::Spacedrop.to_bytes()) // TODO: Proper Spaceblock Header + .await + .unwrap(); + + stream + .write_u8(file_length as u8) // TODO: This is obviously gonna be an int overflow. Fix that. Use `u64` in proper Spaceblock Header + .await + .unwrap(); + + debug!("Starting Spacedrop to peer '{peer_id}'"); + let i = Instant::now(); + + let mut buffer = Vec::new(); + reader.read_to_end(&mut buffer).await.unwrap(); + println!("READ {:?}", buffer); + + stream.write_all(&buffer).await.unwrap(); + + // io::copy(&mut reader, &mut stream).await.unwrap(); // TODO: Use Spaceblock protocol! + + debug!( + "Finished Spacedrop to peer '{peer_id}' after '{:?}", + i.elapsed() + ); + } } diff --git a/crates/p2p/src/event.rs b/crates/p2p/src/event.rs index dbb6db3cd..1382c0350 100644 --- a/crates/p2p/src/event.rs +++ b/crates/p2p/src/event.rs @@ -12,39 +12,39 @@ use super::PeerId; #[cfg_attr(feature = "specta", derive(specta::Type))] #[cfg_attr(feature = "serde", serde(tag = "type"))] pub enum Event { - /// add a network interface on this node to listen for - AddListenAddr(SocketAddr), - /// remove a network interface from this node so that we don't listen to it - RemoveListenAddr(SocketAddr), - /// discovered peer on your local network - PeerDiscovered(DiscoveredPeer), - /// a discovered peer has disappeared from the network - PeerExpired { - id: PeerId, - // Will be none if we receive the expire event without having ever seen a discover event. - metadata: Option, - }, - /// communication was established with a peer. - /// Theere could actually be multiple connections under the hood but we smooth it over in this API. - PeerConnected(ConnectedPeer), - /// communication was lost with a peer. - PeerDisconnected(PeerId), - /// the peer has opened a new substream - #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] - PeerMessage(PeerMessageEvent), + /// add a network interface on this node to listen for + AddListenAddr(SocketAddr), + /// remove a network interface from this node so that we don't listen to it + RemoveListenAddr(SocketAddr), + /// discovered peer on your local network + PeerDiscovered(DiscoveredPeer), + /// a discovered peer has disappeared from the network + PeerExpired { + id: PeerId, + // Will be none if we receive the expire event without having ever seen a discover event. + metadata: Option, + }, + /// communication was established with a peer. + /// Theere could actually be multiple connections under the hood but we smooth it over in this API. + PeerConnected(ConnectedPeer), + /// communication was lost with a peer. + PeerDisconnected(PeerId), + /// the peer has opened a new substream + #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] + PeerMessage(PeerMessageEvent), } #[derive(Debug)] pub struct PeerMessageEvent { - pub peer_id: PeerId, - pub manager: Arc>, - pub stream: SpaceTimeStream, - // Prevent manual creation by end-user - pub(crate) _priv: (), + pub peer_id: PeerId, + pub manager: Arc>, + pub stream: SpaceTimeStream, + // Prevent manual creation by end-user + pub(crate) _priv: (), } impl From> for Event { - fn from(event: PeerMessageEvent) -> Self { - Self::PeerMessage(event) - } + fn from(event: PeerMessageEvent) -> Self { + Self::PeerMessage(event) + } } diff --git a/crates/p2p/src/manager.rs b/crates/p2p/src/manager.rs index fedbe00b0..e082cdf9c 100644 --- a/crates/p2p/src/manager.rs +++ b/crates/p2p/src/manager.rs @@ -1,7 +1,7 @@ use std::{ - collections::{HashMap, HashSet}, - net::SocketAddr, - sync::{atomic::AtomicBool, Arc}, + collections::{HashMap, HashSet}, + net::SocketAddr, + sync::{atomic::AtomicBool, Arc}, }; use libp2p::{core::muxing::StreamMuxerBox, quic, Swarm, Transport}; @@ -10,132 +10,130 @@ use tokio::sync::{mpsc, oneshot, RwLock}; use tracing::{debug, error, warn}; use crate::{ - spacetime::{SpaceTime, SpaceTimeStream}, - AsyncFn, DiscoveredPeer, Keypair, ManagerStream, ManagerStreamAction, Mdns, Metadata, PeerId, + spacetime::{SpaceTime, SpaceTimeStream}, + AsyncFn, DiscoveredPeer, Keypair, ManagerStream, ManagerStreamAction, Mdns, Metadata, PeerId, }; /// Is the core component of the P2P system that holds the state and delegates actions to the other components #[derive(Debug)] pub struct Manager { - pub(crate) peer_id: PeerId, - pub(crate) listen_addrs: RwLock>, - pub(crate) discovered: RwLock>>, - pub(crate) application_name: &'static [u8], - event_stream_tx: mpsc::Sender>, + pub(crate) peer_id: PeerId, + pub(crate) listen_addrs: RwLock>, + pub(crate) discovered: RwLock>>, + pub(crate) application_name: &'static [u8], + event_stream_tx: mpsc::Sender>, } impl Manager { - /// create a new P2P manager. Please do your best to make the callback closures as fast as possible because they will slow the P2P event loop! - pub async fn new( - application_name: &'static str, - keypair: &Keypair, - fn_get_metadata: TMetadataFn, - ) -> Result<(Arc, ManagerStream), ManagerError> - where - TMetadataFn: AsyncFn, - { - application_name - .chars() - .all(|c| char::is_alphanumeric(c) || c == '-') - .then_some(()) - .ok_or(ManagerError::InvalidAppName)?; + /// create a new P2P manager. Please do your best to make the callback closures as fast as possible because they will slow the P2P event loop! + pub async fn new( + application_name: &'static str, + keypair: &Keypair, + fn_get_metadata: TMetadataFn, + ) -> Result<(Arc, ManagerStream), ManagerError> + where + TMetadataFn: AsyncFn, + { + application_name + .chars() + .all(|c| char::is_alphanumeric(c) || c == '-') + .then_some(()) + .ok_or(ManagerError::InvalidAppName)?; - let (event_stream_tx, event_stream_rx) = mpsc::channel(1024); - let this = Arc::new(Self { - // Look this is bad but it's hard to avoid. Technically a memory leak but it's a small amount of memory and is should done on startup on the P2P system. - application_name: Box::leak(Box::new( - format!("/{}/spacetime/1.0.0", application_name) - .as_bytes() - .to_vec(), - )), - peer_id: PeerId(keypair.public().to_peer_id()), - listen_addrs: RwLock::new(Default::default()), - discovered: RwLock::new(Default::default()), - event_stream_tx, - }); + let (event_stream_tx, event_stream_rx) = mpsc::channel(1024); + let this = Arc::new(Self { + // Look this is bad but it's hard to avoid. Technically a memory leak but it's a small amount of memory and is should done on startup on the P2P system. + application_name: Box::leak(Box::new( + format!("/{}/spacetime/1.0.0", application_name) + .as_bytes() + .to_vec(), + )), + peer_id: PeerId(keypair.public().to_peer_id()), + listen_addrs: RwLock::new(Default::default()), + discovered: RwLock::new(Default::default()), + event_stream_tx, + }); - let mut swarm = Swarm::with_tokio_executor( - quic::GenTransport::::new(quic::Config::new(keypair.inner())) - .map(|(p, c), _| (p, StreamMuxerBox::new(c))) - .boxed(), - SpaceTime::new(this.clone()), - keypair.public().to_peer_id(), - ); - { - let listener_id = swarm + let mut swarm = Swarm::with_tokio_executor( + quic::GenTransport::::new(quic::Config::new(keypair.inner())) + .map(|(p, c), _| (p, StreamMuxerBox::new(c))) + .boxed(), + SpaceTime::new(this.clone()), + keypair.public().to_peer_id(), + ); + { + let listener_id = swarm .listen_on("/ip4/0.0.0.0/udp/0/quic-v1".parse().expect("Error passing libp2p multiaddr. This value is hardcoded so this should be impossible.")) .unwrap(); - debug!("created ipv4 listener with id '{:?}'", listener_id); - } - { - let listener_id = swarm + debug!("created ipv4 listener with id '{:?}'", listener_id); + } + { + let listener_id = swarm .listen_on("/ip6/::/udp/0/quic-v1".parse().expect("Error passing libp2p multiaddr. This value is hardcoded so this should be impossible.")) .unwrap(); - debug!("created ipv4 listener with id '{:?}'", listener_id); - } + debug!("created ipv4 listener with id '{:?}'", listener_id); + } - Ok(( - this.clone(), - ManagerStream { - manager: this.clone(), - event_stream_rx, - swarm, - mdns: Mdns::new(this, application_name, fn_get_metadata).unwrap(), - is_advertisement_queued: AtomicBool::new(false), - }, - )) - } + Ok(( + this.clone(), + ManagerStream { + manager: this.clone(), + event_stream_rx, + swarm, + mdns: Mdns::new(this, application_name, fn_get_metadata).unwrap(), + is_advertisement_queued: AtomicBool::new(false), + }, + )) + } - pub(crate) async fn emit(&self, event: ManagerStreamAction) { - match self.event_stream_tx.send(event).await { - Ok(_) => {} - Err(err) => warn!("error emitting event: {}", err), - } - } + pub(crate) async fn emit(&self, event: ManagerStreamAction) { + match self.event_stream_tx.send(event).await { + Ok(_) => {} + Err(err) => warn!("error emitting event: {}", err), + } + } - pub fn peer_id(&self) -> PeerId { - self.peer_id - } + pub fn peer_id(&self) -> PeerId { + self.peer_id + } - pub async fn listen_addrs(&self) -> HashSet { - self.listen_addrs.read().await.clone() - } + pub async fn listen_addrs(&self) -> HashSet { + self.listen_addrs.read().await.clone() + } - pub async fn get_discovered_peers(&self) -> Vec> { - self.discovered.read().await.values().cloned().collect() - } + pub async fn get_discovered_peers(&self) -> Vec> { + self.discovered.read().await.values().cloned().collect() + } - pub async fn get_connected_peers(&self) -> Result, ()> { - let (tx, rx) = oneshot::channel(); - self.emit(ManagerStreamAction::GetConnectedPeers(tx)).await; - rx.await.map_err(|_| { - warn!("failed to get connected peers 3 times, returning error"); - () - }) - } + pub async fn get_connected_peers(&self) -> Result, ()> { + let (tx, rx) = oneshot::channel(); + self.emit(ManagerStreamAction::GetConnectedPeers(tx)).await; + rx.await.map_err(|_| { + warn!("failed to get connected peers 3 times, returning error"); + }) + } - pub async fn stream(&self, peer_id: PeerId) -> Result { - // TODO: With this system you can send to any random peer id. Can I reduce that by requiring `.connect(peer_id).unwrap().send(data)` or something like that. - let (tx, rx) = oneshot::channel(); - self.emit(ManagerStreamAction::StartStream(peer_id, tx)) - .await; - rx.await.map_err(|_| { - warn!("failed to queue establishing stream to peer '{peer_id}'!"); - () - }) - } + pub async fn stream(&self, peer_id: PeerId) -> Result { + // TODO: With this system you can send to any random peer id. Can I reduce that by requiring `.connect(peer_id).unwrap().send(data)` or something like that. + let (tx, rx) = oneshot::channel(); + self.emit(ManagerStreamAction::StartStream(peer_id, tx)) + .await; + rx.await.map_err(|_| { + warn!("failed to queue establishing stream to peer '{peer_id}'!"); + }) + } - pub async fn broadcast(&self, data: Vec) { - self.emit(ManagerStreamAction::BroadcastData(data)).await; - } + pub async fn broadcast(&self, data: Vec) { + self.emit(ManagerStreamAction::BroadcastData(data)).await; + } } #[derive(Error, Debug)] pub enum ManagerError { - #[error( - "the application name you application provided is invalid. Ensure it is alphanumeric!" - )] - InvalidAppName, - #[error("error with mdns discovery: {0}")] - Mdns(#[from] mdns_sd::Error), + #[error( + "the application name you application provided is invalid. Ensure it is alphanumeric!" + )] + InvalidAppName, + #[error("error with mdns discovery: {0}")] + Mdns(#[from] mdns_sd::Error), } diff --git a/crates/p2p/src/manager_stream.rs b/crates/p2p/src/manager_stream.rs index 9ad76a57d..77311a8ac 100644 --- a/crates/p2p/src/manager_stream.rs +++ b/crates/p2p/src/manager_stream.rs @@ -1,200 +1,200 @@ use std::{ - net::SocketAddr, - sync::{ - atomic::{AtomicBool, Ordering}, - Arc, - }, + net::SocketAddr, + sync::{ + atomic::{AtomicBool, Ordering}, + Arc, + }, }; use libp2p::{ - futures::StreamExt, - swarm::{ - dial_opts::{DialOpts, PeerCondition}, - NetworkBehaviourAction, NotifyHandler, SwarmEvent, - }, - Multiaddr, Swarm, + futures::StreamExt, + swarm::{ + dial_opts::{DialOpts, PeerCondition}, + NetworkBehaviourAction, NotifyHandler, SwarmEvent, + }, + Multiaddr, Swarm, }; use tokio::sync::{mpsc, oneshot}; use tracing::{debug, error, warn}; use crate::{ - quic_multiaddr_to_socketaddr, socketaddr_to_quic_multiaddr, - spacetime::{OutboundRequest, SpaceTime, SpaceTimeStream}, - AsyncFn, Event, Manager, Mdns, Metadata, PeerId, + quic_multiaddr_to_socketaddr, socketaddr_to_quic_multiaddr, + spacetime::{OutboundRequest, SpaceTime, SpaceTimeStream}, + AsyncFn, Event, Manager, Mdns, Metadata, PeerId, }; /// TODO pub enum ManagerStreamAction { - /// Events are returned to the application via the `ManagerStream::next` method. - Event(Event), - /// TODO - GetConnectedPeers(oneshot::Sender>), - /// Tell the [`libp2p::Swarm`](libp2p::Swarm) to establish a new connection to a peer. - Dial { - peer_id: PeerId, - addresses: Vec, - }, - /// TODO - StartStream(PeerId, oneshot::Sender), - /// TODO - BroadcastData(Vec), + /// Events are returned to the application via the `ManagerStream::next` method. + Event(Event), + /// TODO + GetConnectedPeers(oneshot::Sender>), + /// Tell the [`libp2p::Swarm`](libp2p::Swarm) to establish a new connection to a peer. + Dial { + peer_id: PeerId, + addresses: Vec, + }, + /// TODO + StartStream(PeerId, oneshot::Sender), + /// TODO + BroadcastData(Vec), } impl From> for ManagerStreamAction { - fn from(event: Event) -> Self { - Self::Event(event) - } + fn from(event: Event) -> Self { + Self::Event(event) + } } /// TODO pub struct ManagerStream where - TMetadata: Metadata, - TMetadataFn: AsyncFn, + TMetadata: Metadata, + TMetadataFn: AsyncFn, { - pub(crate) manager: Arc>, - pub(crate) event_stream_rx: mpsc::Receiver>, - pub(crate) swarm: Swarm>, - pub(crate) mdns: Mdns, - pub(crate) is_advertisement_queued: AtomicBool, + pub(crate) manager: Arc>, + pub(crate) event_stream_rx: mpsc::Receiver>, + pub(crate) swarm: Swarm>, + pub(crate) mdns: Mdns, + pub(crate) is_advertisement_queued: AtomicBool, } impl ManagerStream where - TMetadata: Metadata, - TMetadataFn: AsyncFn, + TMetadata: Metadata, + TMetadataFn: AsyncFn, { - // Your application should keep polling this until `None` is received or the P2P system will be halted. - pub async fn next(&mut self) -> Option> { - // We loop polling internal services until an event comes in that needs to be sent to the parent application. - loop { - tokio::select! { - _ = self.mdns.poll() => {}, - event = self.event_stream_rx.recv() => { - // If the sender has shut down we return `None` to also shut down too. - match event? { - ManagerStreamAction::Event(event) => return Some(event), - ManagerStreamAction::GetConnectedPeers(response) => { - response.send(self.swarm.behaviour().connected_peers.values().map(|p| p.peer_id.clone()).collect::>()).map_err(|_| error!("Error sending response to `GetConnectedPeers` request! Sending was dropped!")).ok(); - }, - ManagerStreamAction::Dial { peer_id, addresses } => { - match self.swarm.dial( - DialOpts::peer_id(peer_id.0) - .condition(PeerCondition::Disconnected) - .addresses( - addresses - .iter() - .map(|addr| socketaddr_to_quic_multiaddr(addr)) - .collect(), - ) - .extend_addresses_through_behaviour() - .build(), - ) { - Ok(_) => {} - Err(err) => warn!( - "error dialing peer '{}' with addresses '{:?}': {}", - peer_id, addresses, err - ), - } - } - ManagerStreamAction::StartStream(peer_id, rx) => { - self.swarm.behaviour_mut().pending_events - .push_back(NetworkBehaviourAction::NotifyHandler { - peer_id: peer_id.0, - handler: NotifyHandler::Any, - event: OutboundRequest::Stream(rx), - }); - } - ManagerStreamAction::BroadcastData(data) => { - let swarm = self.swarm.behaviour_mut(); - for peer in swarm.connected_peers.values() { - swarm.pending_events - .push_back(NetworkBehaviourAction::NotifyHandler { - peer_id: peer.peer_id.0, - handler: NotifyHandler::Any, - event: OutboundRequest::Data(data.clone()), - }); - } - } - } - } - event = self.swarm.select_next_some() => { - match event { - SwarmEvent::Behaviour(()) => {}, - SwarmEvent::ConnectionEstablished { .. } => {}, - SwarmEvent::ConnectionClosed { .. } => {}, - SwarmEvent::IncomingConnection { local_addr, .. } => debug!("incoming connection from '{}'", local_addr), - SwarmEvent::IncomingConnectionError { local_addr, error, .. } => warn!("handshake error with incoming connection from '{}': {}", local_addr, error), - SwarmEvent::OutgoingConnectionError { peer_id, error } => warn!("error establishing connection with '{:?}': {}", peer_id, error), - SwarmEvent::BannedPeer { peer_id, .. } => warn!("banned peer '{}' attempted to connection and was rejected", peer_id), - SwarmEvent::NewListenAddr{ address, .. } => { - match quic_multiaddr_to_socketaddr(address) { - Ok(addr) => { - debug!("listen address added: {}", addr); - self.manager.listen_addrs.write().await.insert(addr); - if !self.is_advertisement_queued.load(Ordering::Relaxed) { - self.is_advertisement_queued.store(true, Ordering::Relaxed); - self.mdns.advertise(); - } - self.manager.emit(Event::AddListenAddr(addr).into()).await; - self.mdns.advertise(); - }, - Err(err) => { - warn!("error passing listen address: {}", err); - continue; - } - } - }, - SwarmEvent::ExpiredListenAddr { address, .. } => { - match Self::unregister_addr(&self.manager, &self.mdns, &self.is_advertisement_queued, address).await { - Ok(_) => {}, - Err(err) => { - warn!("error passing listen address: {}", err); - continue; - } - } - } - SwarmEvent::ListenerClosed { listener_id, addresses, reason } => { - debug!("listener '{:?}' was closed due to: {:?}", listener_id, reason); - for address in addresses { - match Self::unregister_addr(&self.manager, &self.mdns, &self.is_advertisement_queued, address).await { - Ok(_) => {}, - Err(err) => { - warn!("error passing listen address: {}", err); - continue; - } - } - } - } - SwarmEvent::ListenerError { listener_id, error } => warn!("listener '{:?}' reported a non-fatal error: {}", listener_id, error), - SwarmEvent::Dialing(_peer_id) => {}, - } - } - } - } - } + // Your application should keep polling this until `None` is received or the P2P system will be halted. + pub async fn next(&mut self) -> Option> { + // We loop polling internal services until an event comes in that needs to be sent to the parent application. + loop { + tokio::select! { + _ = self.mdns.poll() => {}, + event = self.event_stream_rx.recv() => { + // If the sender has shut down we return `None` to also shut down too. + match event? { + ManagerStreamAction::Event(event) => return Some(event), + ManagerStreamAction::GetConnectedPeers(response) => { + response.send(self.swarm.behaviour().connected_peers.values().map(|p| p.peer_id).collect::>()).map_err(|_| error!("Error sending response to `GetConnectedPeers` request! Sending was dropped!")).ok(); + }, + ManagerStreamAction::Dial { peer_id, addresses } => { + match self.swarm.dial( + DialOpts::peer_id(peer_id.0) + .condition(PeerCondition::Disconnected) + .addresses( + addresses + .iter() + .map(socketaddr_to_quic_multiaddr) + .collect(), + ) + .extend_addresses_through_behaviour() + .build(), + ) { + Ok(_) => {} + Err(err) => warn!( + "error dialing peer '{}' with addresses '{:?}': {}", + peer_id, addresses, err + ), + } + } + ManagerStreamAction::StartStream(peer_id, rx) => { + self.swarm.behaviour_mut().pending_events + .push_back(NetworkBehaviourAction::NotifyHandler { + peer_id: peer_id.0, + handler: NotifyHandler::Any, + event: OutboundRequest::Stream(rx), + }); + } + ManagerStreamAction::BroadcastData(data) => { + let swarm = self.swarm.behaviour_mut(); + for peer in swarm.connected_peers.values() { + swarm.pending_events + .push_back(NetworkBehaviourAction::NotifyHandler { + peer_id: peer.peer_id.0, + handler: NotifyHandler::Any, + event: OutboundRequest::Data(data.clone()), + }); + } + } + } + } + event = self.swarm.select_next_some() => { + match event { + SwarmEvent::Behaviour(()) => {}, + SwarmEvent::ConnectionEstablished { .. } => {}, + SwarmEvent::ConnectionClosed { .. } => {}, + SwarmEvent::IncomingConnection { local_addr, .. } => debug!("incoming connection from '{}'", local_addr), + SwarmEvent::IncomingConnectionError { local_addr, error, .. } => warn!("handshake error with incoming connection from '{}': {}", local_addr, error), + SwarmEvent::OutgoingConnectionError { peer_id, error } => warn!("error establishing connection with '{:?}': {}", peer_id, error), + SwarmEvent::BannedPeer { peer_id, .. } => warn!("banned peer '{}' attempted to connection and was rejected", peer_id), + SwarmEvent::NewListenAddr{ address, .. } => { + match quic_multiaddr_to_socketaddr(address) { + Ok(addr) => { + debug!("listen address added: {}", addr); + self.manager.listen_addrs.write().await.insert(addr); + if !self.is_advertisement_queued.load(Ordering::Relaxed) { + self.is_advertisement_queued.store(true, Ordering::Relaxed); + self.mdns.advertise(); + } + self.manager.emit(Event::AddListenAddr(addr).into()).await; + self.mdns.advertise(); + }, + Err(err) => { + warn!("error passing listen address: {}", err); + continue; + } + } + }, + SwarmEvent::ExpiredListenAddr { address, .. } => { + match Self::unregister_addr(&self.manager, &self.mdns, &self.is_advertisement_queued, address).await { + Ok(_) => {}, + Err(err) => { + warn!("error passing listen address: {}", err); + continue; + } + } + } + SwarmEvent::ListenerClosed { listener_id, addresses, reason } => { + debug!("listener '{:?}' was closed due to: {:?}", listener_id, reason); + for address in addresses { + match Self::unregister_addr(&self.manager, &self.mdns, &self.is_advertisement_queued, address).await { + Ok(_) => {}, + Err(err) => { + warn!("error passing listen address: {}", err); + continue; + } + } + } + } + SwarmEvent::ListenerError { listener_id, error } => warn!("listener '{:?}' reported a non-fatal error: {}", listener_id, error), + SwarmEvent::Dialing(_peer_id) => {}, + } + } + } + } + } - async fn unregister_addr( - manager: &Arc>, - mdns: &Mdns, - is_advertisement_queued: &AtomicBool, - address: Multiaddr, - ) -> Result<(), String> { - match quic_multiaddr_to_socketaddr(address) { - Ok(addr) => { - debug!("listen address removed: {}", addr); - manager.listen_addrs.write().await.remove(&addr); - let _ = mdns.unregister_mdns(); - if !is_advertisement_queued.load(Ordering::Relaxed) { - is_advertisement_queued.store(true, Ordering::Relaxed); - mdns.advertise(); - } - manager - .emit(ManagerStreamAction::Event(Event::RemoveListenAddr(addr))) - .await; - Ok(()) - } - Err(err) => Err(err), - } - } + async fn unregister_addr( + manager: &Arc>, + mdns: &Mdns, + is_advertisement_queued: &AtomicBool, + address: Multiaddr, + ) -> Result<(), String> { + match quic_multiaddr_to_socketaddr(address) { + Ok(addr) => { + debug!("listen address removed: {}", addr); + manager.listen_addrs.write().await.remove(&addr); + let _ = mdns.unregister_mdns(); + if !is_advertisement_queued.load(Ordering::Relaxed) { + is_advertisement_queued.store(true, Ordering::Relaxed); + mdns.advertise(); + } + manager + .emit(ManagerStreamAction::Event(Event::RemoveListenAddr(addr))) + .await; + Ok(()) + } + Err(err) => Err(err), + } + } } diff --git a/crates/p2p/src/mdns.rs b/crates/p2p/src/mdns.rs index 1146217e0..8dbf9416b 100644 --- a/crates/p2p/src/mdns.rs +++ b/crates/p2p/src/mdns.rs @@ -1,10 +1,10 @@ use std::{ - collections::HashMap, - net::{IpAddr, SocketAddr}, - pin::Pin, - str::FromStr, - sync::Arc, - time::Duration, + collections::HashMap, + net::{IpAddr, SocketAddr}, + pin::Pin, + str::FromStr, + sync::Arc, + time::Duration, }; use mdns_sd::{ServiceDaemon, ServiceEvent, ServiceInfo}; @@ -19,214 +19,214 @@ const MDNS_READVERTISEMENT_INTERVAL: Duration = Duration::from_secs(60); // Ever /// TODO pub struct Mdns where - TMetadata: Metadata, - TMetadataFn: AsyncFn, + TMetadata: Metadata, + TMetadataFn: AsyncFn, { - manager: Arc>, - fn_get_metadata: TMetadataFn, - mdns_daemon: ServiceDaemon, - mdns_service_receiver: flume::Receiver, - service_name: String, - next_mdns_advertisement: Pin>, + manager: Arc>, + fn_get_metadata: TMetadataFn, + mdns_daemon: ServiceDaemon, + mdns_service_receiver: flume::Receiver, + service_name: String, + next_mdns_advertisement: Pin>, } impl Mdns where - TMetadata: Metadata, - TMetadataFn: AsyncFn, + TMetadata: Metadata, + TMetadataFn: AsyncFn, { - pub fn new( - manager: Arc>, - application_name: &'static str, - fn_get_metadata: TMetadataFn, - ) -> Result - where - TMetadataFn: AsyncFn, - { - let mdns_daemon = ServiceDaemon::new()?; - let service_name = format!("_{}._udp.local.", application_name); - let mdns_service_receiver = mdns_daemon.browse(&service_name)?; + pub fn new( + manager: Arc>, + application_name: &'static str, + fn_get_metadata: TMetadataFn, + ) -> Result + where + TMetadataFn: AsyncFn, + { + let mdns_daemon = ServiceDaemon::new()?; + let service_name = format!("_{}._udp.local.", application_name); + let mdns_service_receiver = mdns_daemon.browse(&service_name)?; - let this = Self { - manager, - fn_get_metadata, - mdns_daemon, - mdns_service_receiver, - service_name, - next_mdns_advertisement: Box::pin(sleep_until( - Instant::now() + MDNS_READVERTISEMENT_INTERVAL, - )), - }; - this.advertise(); - Ok(this) - } + let this = Self { + manager, + fn_get_metadata, + mdns_daemon, + mdns_service_receiver, + service_name, + next_mdns_advertisement: Box::pin(sleep_until( + Instant::now() + MDNS_READVERTISEMENT_INTERVAL, + )), + }; + this.advertise(); + Ok(this) + } - pub fn unregister_mdns(&self) -> mdns_sd::Result> { - self.mdns_daemon - .unregister(&format!("{}.{}", self.manager.peer_id, self.service_name)) - } + pub fn unregister_mdns(&self) -> mdns_sd::Result> { + self.mdns_daemon + .unregister(&format!("{}.{}", self.manager.peer_id, self.service_name)) + } - /// Do an mdns advertisement to the network - pub fn advertise(&self) { - // TODO: Instead of spawning maybe do this as part of the polling loop to avoid needing persitent reference to manager. - let manager = self.manager.clone(); - let service_name = self.service_name.clone(); - // let fn_get_metadata = self.fn_get_metadata.clone(); - let mdns_daemon = self.mdns_daemon.clone(); + /// Do an mdns advertisement to the network + pub fn advertise(&self) { + // TODO: Instead of spawning maybe do this as part of the polling loop to avoid needing persitent reference to manager. + let manager = self.manager.clone(); + let service_name = self.service_name.clone(); + // let fn_get_metadata = self.fn_get_metadata.clone(); + let mdns_daemon = self.mdns_daemon.clone(); - let metadata_fut = (self.fn_get_metadata)(); + let metadata_fut = (self.fn_get_metadata)(); - tokio::spawn(async move { - let metadata = metadata_fut.await.to_hashmap(); - let peer_id = manager.peer_id.0.to_base58(); + tokio::spawn(async move { + let metadata = metadata_fut.await.to_hashmap(); + let peer_id = manager.peer_id.0.to_base58(); - // This is in simple terms converts from `Vec<(ip, port)>` to `Vec<(Vec, port)>` - let mut services = HashMap::::new(); - for addr in manager.listen_addrs.read().await.iter() { - let addr = match addr { - SocketAddr::V4(addr) => addr, - // TODO: Our mdns library doesn't support Ipv6. This code has the infra to support it so once this issue is fixed upstream we can just flip it on. - // Refer to issue: https://github.com/keepsimple1/mdns-sd/issues/61 - SocketAddr::V6(_) => continue, - }; + // This is in simple terms converts from `Vec<(ip, port)>` to `Vec<(Vec, port)>` + let mut services = HashMap::::new(); + for addr in manager.listen_addrs.read().await.iter() { + let addr = match addr { + SocketAddr::V4(addr) => addr, + // TODO: Our mdns library doesn't support Ipv6. This code has the infra to support it so once this issue is fixed upstream we can just flip it on. + // Refer to issue: https://github.com/keepsimple1/mdns-sd/issues/61 + SocketAddr::V6(_) => continue, + }; - if let Some(mut service) = services.remove(&addr.port()) { - service.insert_ipv4addr(*addr.ip()); - services.insert(addr.port(), service); - } else { - let service = match ServiceInfo::new( - &service_name, - &peer_id, - &format!("{}.", peer_id), - *addr.ip(), - addr.port(), - Some(metadata.clone()), // TODO: Prevent the user defining a value that overflows a DNS record - ) { - Ok(service) => service, - Err(err) => { - warn!("error creating mdns service info: {}", err); - continue; - } - }; - services.insert(addr.port(), service); - } - } + if let Some(mut service) = services.remove(&addr.port()) { + service.insert_ipv4addr(*addr.ip()); + services.insert(addr.port(), service); + } else { + let service = match ServiceInfo::new( + &service_name, + &peer_id, + &format!("{}.", peer_id), + *addr.ip(), + addr.port(), + Some(metadata.clone()), // TODO: Prevent the user defining a value that overflows a DNS record + ) { + Ok(service) => service, + Err(err) => { + warn!("error creating mdns service info: {}", err); + continue; + } + }; + services.insert(addr.port(), service); + } + } - for (_, service) in services.into_iter() { - debug!("advertising mdns service: {:?}", service); - match mdns_daemon.register(service) { - Ok(_) => {} - Err(err) => warn!("error registering mdns service: {}", err), - } - } - }); - } + for (_, service) in services.into_iter() { + debug!("advertising mdns service: {:?}", service); + match mdns_daemon.register(service) { + Ok(_) => {} + Err(err) => warn!("error registering mdns service: {}", err), + } + } + }); + } - // TODO: if the channel's sender is dropped will this cause the `tokio::select` in the `manager.rs` to infinitely loop? - pub async fn poll(&mut self) { - tokio::select! { - _ = &mut self.next_mdns_advertisement => { - self.advertise(); - self.next_mdns_advertisement = Box::pin(sleep_until(Instant::now() + MDNS_READVERTISEMENT_INTERVAL)); - } - event = self.mdns_service_receiver.recv_async() => { - let event = event.unwrap(); // TODO: Error handling - match event { - ServiceEvent::SearchStarted(_) => {} - ServiceEvent::ServiceFound(_, _) => {} - ServiceEvent::ServiceResolved(info) => { - let raw_peer_id = info - .get_fullname() - .replace(&format!(".{}", self.service_name), ""); + // TODO: if the channel's sender is dropped will this cause the `tokio::select` in the `manager.rs` to infinitely loop? + pub async fn poll(&mut self) { + tokio::select! { + _ = &mut self.next_mdns_advertisement => { + self.advertise(); + self.next_mdns_advertisement = Box::pin(sleep_until(Instant::now() + MDNS_READVERTISEMENT_INTERVAL)); + } + event = self.mdns_service_receiver.recv_async() => { + let event = event.unwrap(); // TODO: Error handling + match event { + ServiceEvent::SearchStarted(_) => {} + ServiceEvent::ServiceFound(_, _) => {} + ServiceEvent::ServiceResolved(info) => { + let raw_peer_id = info + .get_fullname() + .replace(&format!(".{}", self.service_name), ""); - match PeerId::from_str(&raw_peer_id) { - Ok(peer_id) => { - // Prevent discovery of the current peer. - if peer_id == self.manager.peer_id { - return; - } + match PeerId::from_str(&raw_peer_id) { + Ok(peer_id) => { + // Prevent discovery of the current peer. + if peer_id == self.manager.peer_id { + return; + } - match TMetadata::from_hashmap( - &info - .get_properties() - .iter() - .map(|v| (v.key().to_owned(), v.val().to_owned())) - .collect(), - ) { - Ok(metadata) => { - let peer = { - let mut discovered_peers = - self.manager.discovered.write().await; + match TMetadata::from_hashmap( + &info + .get_properties() + .iter() + .map(|v| (v.key().to_owned(), v.val().to_owned())) + .collect(), + ) { + Ok(metadata) => { + let peer = { + let mut discovered_peers = + self.manager.discovered.write().await; - let peer = if let Some(peer) = discovered_peers.remove(&peer_id) { + let peer = if let Some(peer) = discovered_peers.remove(&peer_id) { - peer - } else { - DiscoveredPeer { - manager: self.manager.clone(), - peer_id: peer_id, - metadata, - addresses: info - .get_addresses() - .iter() - .map(|addr| { - SocketAddr::new( - IpAddr::V4(addr.clone()), - info.get_port(), - ) - }) - .collect(), - } - }; + peer + } else { + DiscoveredPeer { + manager: self.manager.clone(), + peer_id, + metadata, + addresses: info + .get_addresses() + .iter() + .map(|addr| { + SocketAddr::new( + IpAddr::V4(*addr), + info.get_port(), + ) + }) + .collect(), + } + }; - discovered_peers.insert(peer_id, peer.clone()); - peer - }; - self.manager.emit(Event::PeerDiscovered(peer).into()).await; - } - Err(err) => { - error!("error parsing metadata for peer '{}': {}", raw_peer_id, err) - } - } - } - Err(_) => warn!( - "resolved peer advertising itself with an invalid peer_id '{}'", - raw_peer_id - ), - } - } - ServiceEvent::ServiceRemoved(_, fullname) => { - let raw_peer_id = fullname.replace(&format!(".{}", self.service_name), ""); + discovered_peers.insert(peer_id, peer.clone()); + peer + }; + self.manager.emit(Event::PeerDiscovered(peer).into()).await; + } + Err(err) => { + error!("error parsing metadata for peer '{}': {}", raw_peer_id, err) + } + } + } + Err(_) => warn!( + "resolved peer advertising itself with an invalid peer_id '{}'", + raw_peer_id + ), + } + } + ServiceEvent::ServiceRemoved(_, fullname) => { + let raw_peer_id = fullname.replace(&format!(".{}", self.service_name), ""); - match PeerId::from_str(&raw_peer_id) { - Ok(peer_id) => { - // Prevent discovery of the current peer. - if peer_id == self.manager.peer_id { - return; - } + match PeerId::from_str(&raw_peer_id) { + Ok(peer_id) => { + // Prevent discovery of the current peer. + if peer_id == self.manager.peer_id { + return; + } - { - let mut discovered_peers = - self.manager.discovered.write().await; - let peer = discovered_peers.remove(&peer_id); + { + let mut discovered_peers = + self.manager.discovered.write().await; + let peer = discovered_peers.remove(&peer_id); - self.manager - .emit(Event::PeerExpired { - id: peer_id, - metadata: peer.map(|p| p.metadata), - }.into()) - .await; - } - } - Err(_) => warn!( - "resolved peer de-advertising itself with an invalid peer_id '{}'", - raw_peer_id - ), - } - } - ServiceEvent::SearchStopped(_) => {} - } - } - } - } + self.manager + .emit(Event::PeerExpired { + id: peer_id, + metadata: peer.map(|p| p.metadata), + }.into()) + .await; + } + } + Err(_) => warn!( + "resolved peer de-advertising itself with an invalid peer_id '{}'", + raw_peer_id + ), + } + } + ServiceEvent::SearchStopped(_) => {} + } + } + } + } } diff --git a/crates/p2p/src/peer.rs b/crates/p2p/src/peer.rs index cabbe5719..a6f0dbc01 100644 --- a/crates/p2p/src/peer.rs +++ b/crates/p2p/src/peer.rs @@ -10,26 +10,26 @@ use crate::{Manager, ManagerStreamAction, Metadata, PeerId}; #[cfg_attr(feature = "serde", derive(serde::Serialize))] #[cfg_attr(feature = "specta", derive(specta::Type))] pub struct DiscoveredPeer { - #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] - pub(crate) manager: Arc>, - /// get the peer id of the discovered peer - pub peer_id: PeerId, - /// get the metadata of the discovered peer - pub metadata: TMetadata, - /// get the addresses of the discovered peer - pub addresses: Vec, + #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] + pub(crate) manager: Arc>, + /// get the peer id of the discovered peer + pub peer_id: PeerId, + /// get the metadata of the discovered peer + pub metadata: TMetadata, + /// get the addresses of the discovered peer + pub addresses: Vec, } impl DiscoveredPeer { - /// dial will queue an event to start a connection with the peer - pub async fn dial(self) { - self.manager - .emit(ManagerStreamAction::Dial { - peer_id: self.peer_id, - addresses: self.addresses, - }) - .await; - } + /// dial will queue an event to start a connection with the peer + pub async fn dial(self) { + self.manager + .emit(ManagerStreamAction::Dial { + peer_id: self.peer_id, + addresses: self.addresses, + }) + .await; + } } /// Represents a connected peer. @@ -38,9 +38,9 @@ impl DiscoveredPeer { #[cfg_attr(feature = "serde", derive(serde::Serialize))] #[cfg_attr(feature = "specta", derive(specta::Type))] pub struct ConnectedPeer { - /// get the peer id of the discovered peer - pub peer_id: PeerId, - /// list of connections between the peer and the local node - #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] - pub connections: HashMap, // TODO: Probs use `thinvec` style thing here + /// get the peer id of the discovered peer + pub peer_id: PeerId, + /// list of connections between the peer and the local node + #[cfg_attr(any(feature = "serde", feature = "specta"), serde(skip))] + pub connections: HashMap, // TODO: Probs use `thinvec` style thing here } diff --git a/crates/p2p/src/spaceblock/mod.rs b/crates/p2p/src/spaceblock/mod.rs index b816c73c2..24b15ce71 100644 --- a/crates/p2p/src/spaceblock/mod.rs +++ b/crates/p2p/src/spaceblock/mod.rs @@ -4,55 +4,55 @@ #![allow(unused)] // TODO: This module is still in heavy development! use std::{ - marker::PhantomData, - path::{Path, PathBuf}, + marker::PhantomData, + path::{Path, PathBuf}, }; /// TODO pub struct BlockSize(i64); impl BlockSize { - // TODO: Validating `BlockSize` are multiple of 2, i think + // TODO: Validating `BlockSize` are multiple of 2, i think - pub fn from_size(size: u64) -> Self { - // TODO: Something like: https://docs.syncthing.net/specs/bep-v1.html#selection-of-block-size - Self(131072) // 128 KiB - } + pub fn from_size(size: u64) -> Self { + // TODO: Something like: https://docs.syncthing.net/specs/bep-v1.html#selection-of-block-size + Self(131072) // 128 KiB + } } /// TODO pub struct TransferRequest<'a> { - name: &'a str, - size: u64, - // TODO: Include file permissions - block_size: u64, + name: &'a str, + size: u64, + // TODO: Include file permissions + block_size: u64, } /// TODO pub struct Block<'a> { - // TODO: File content, checksum, source location so it can be resent! - offset: i64, - size: i64, - data: &'a [u8], - // TODO: Checksum? + // TODO: File content, checksum, source location so it can be resent! + offset: i64, + size: i64, + data: &'a [u8], + // TODO: Checksum? } /// TODO pub struct Transfer<'a> { - // buf: &'a mut [u8], - phantom: PhantomData<&'a ()>, + // buf: &'a mut [u8], + phantom: PhantomData<&'a ()>, } impl<'a> Transfer<'a> { - // TODO: Allow the user to cancel a tranfer - // TODO: Handle if the stream is dropped + // TODO: Allow the user to cancel a tranfer + // TODO: Handle if the stream is dropped - pub fn from_file(path: impl AsRef) -> Self { - // let size = std::fs::metadata(path.as_ref()).unwrap().len(); + pub fn from_file(path: impl AsRef) -> Self { + // let size = std::fs::metadata(path.as_ref()).unwrap().len(); - Self { - // buf: &mut Vec::with_capacity(42069), - phantom: PhantomData, - } - } + Self { + // buf: &mut Vec::with_capacity(42069), + phantom: PhantomData, + } + } } diff --git a/crates/p2p/src/spacetime/behaviour.rs b/crates/p2p/src/spacetime/behaviour.rs index 32f0b792e..89bc95b9c 100644 --- a/crates/p2p/src/spacetime/behaviour.rs +++ b/crates/p2p/src/spacetime/behaviour.rs @@ -1,17 +1,17 @@ use std::{ - collections::{HashMap, VecDeque}, - sync::Arc, - task::{Context, Poll}, + collections::{HashMap, VecDeque}, + sync::Arc, + task::{Context, Poll}, }; use libp2p::{ - core::{ConnectedPoint, Endpoint}, - swarm::{ - derive_prelude::{ConnectionEstablished, ConnectionId, FromSwarm}, - ConnectionClosed, ConnectionDenied, ConnectionHandler, NetworkBehaviour, - NetworkBehaviourAction, PollParameters, THandler, THandlerInEvent, - }, - Multiaddr, + core::{ConnectedPoint, Endpoint}, + swarm::{ + derive_prelude::{ConnectionEstablished, ConnectionId, FromSwarm}, + ConnectionClosed, ConnectionDenied, ConnectionHandler, NetworkBehaviour, + NetworkBehaviourAction, PollParameters, THandler, THandlerInEvent, + }, + Multiaddr, }; use thiserror::Error; use tracing::{debug, warn}; @@ -33,254 +33,256 @@ pub enum OutboundFailure {} /// SpaceTime is a [`NetworkBehaviour`](libp2p::NetworkBehaviour) that implements the SpaceTime protocol. /// This protocol sits under the application to abstract many complexities of 2 way connections and deals with authentication, chucking, etc. pub struct SpaceTime { - pub(crate) manager: Arc>, - pub(crate) pending_events: VecDeque< - NetworkBehaviourAction<::OutEvent, THandlerInEvent>, - >, - // For future me's sake, DON't try and refactor this to use shared state (for the nth time), it doesn't fit into libp2p's synchronous trait and polling model!!! - pub(crate) connected_peers: HashMap, + pub(crate) manager: Arc>, + pub(crate) pending_events: VecDeque< + NetworkBehaviourAction<::OutEvent, THandlerInEvent>, + >, + // For future me's sake, DON't try and refactor this to use shared state (for the nth time), it doesn't fit into libp2p's synchronous trait and polling model!!! + pub(crate) connected_peers: HashMap, } impl SpaceTime { - /// intialise the fabric of space time - pub fn new(manager: Arc>) -> Self { - Self { - manager, - pending_events: VecDeque::new(), - connected_peers: HashMap::new(), - } - } + /// intialise the fabric of space time + pub fn new(manager: Arc>) -> Self { + Self { + manager, + pending_events: VecDeque::new(), + connected_peers: HashMap::new(), + } + } } impl NetworkBehaviour for SpaceTime { - type ConnectionHandler = SpaceTimeConnection; - type OutEvent = (); + type ConnectionHandler = SpaceTimeConnection; + type OutEvent = (); - fn handle_established_inbound_connection( - &mut self, - _connection_id: ConnectionId, - peer_id: libp2p::PeerId, - _local_addr: &Multiaddr, - _remote_addr: &Multiaddr, - ) -> Result, ConnectionDenied> { - Ok(SpaceTimeConnection::new( - PeerId(peer_id), - self.manager.clone(), - )) - } + fn handle_established_inbound_connection( + &mut self, + _connection_id: ConnectionId, + peer_id: libp2p::PeerId, + _local_addr: &Multiaddr, + _remote_addr: &Multiaddr, + ) -> Result, ConnectionDenied> { + Ok(SpaceTimeConnection::new( + PeerId(peer_id), + self.manager.clone(), + )) + } - // TODO: Are we even using the response to this?? - // TODO: Do we need to load from state or can be just pass through the `addresses` arg? - fn handle_pending_outbound_connection( - &mut self, - _connection_id: ConnectionId, - maybe_peer: Option, - _addresses: &[Multiaddr], - _effective_role: Endpoint, - ) -> Result, ConnectionDenied> { - if let Some(peer_id) = maybe_peer { - let mut addresses = Vec::new(); - if let Some(connection) = self.connected_peers.get(&PeerId(peer_id)) { - addresses.extend( - connection - .connections - .iter() - .filter_map(|(_, cp)| match cp { - ConnectedPoint::Dialer { address, .. } => Some(address.clone()), - ConnectedPoint::Listener { - local_addr, - send_back_addr, - } => { - println!( - "TODO: Handle this case! ({} -> {})", - local_addr, send_back_addr - ); - todo!(); - } - }), - ) - } - Ok(addresses) - } else { - Ok(vec![]) - } - } + // TODO: Are we even using the response to this?? + // TODO: Do we need to load from state or can be just pass through the `addresses` arg? + fn handle_pending_outbound_connection( + &mut self, + _connection_id: ConnectionId, + maybe_peer: Option, + _addresses: &[Multiaddr], + _effective_role: Endpoint, + ) -> Result, ConnectionDenied> { + if let Some(peer_id) = maybe_peer { + let mut addresses = Vec::new(); + if let Some(connection) = self.connected_peers.get(&PeerId(peer_id)) { + #[allow(clippy::unnecessary_filter_map)] + // TODO: Clippy is getting annoyed cause of the `todo!()`. Remove this once that is fixed. + addresses.extend( + connection + .connections + .iter() + .filter_map(|(_, cp)| match cp { + ConnectedPoint::Dialer { address, .. } => Some(address.clone()), + ConnectedPoint::Listener { + local_addr, + send_back_addr, + } => { + println!( + "TODO: Handle this case! ({} -> {})", + local_addr, send_back_addr + ); + todo!(); + } + }), + ) + } + Ok(addresses) + } else { + Ok(vec![]) + } + } - fn handle_established_outbound_connection( - &mut self, - _connection_id: ConnectionId, - peer_id: libp2p::PeerId, - _addr: &Multiaddr, - _role_override: Endpoint, - ) -> Result, ConnectionDenied> { - Ok(SpaceTimeConnection::new( - PeerId(peer_id), - self.manager.clone(), - )) - } + fn handle_established_outbound_connection( + &mut self, + _connection_id: ConnectionId, + peer_id: libp2p::PeerId, + _addr: &Multiaddr, + _role_override: Endpoint, + ) -> Result, ConnectionDenied> { + Ok(SpaceTimeConnection::new( + PeerId(peer_id), + self.manager.clone(), + )) + } - fn on_swarm_event(&mut self, event: FromSwarm) { - match event { - FromSwarm::ConnectionEstablished(ConnectionEstablished { - peer_id, - connection_id, - endpoint, - other_established, - .. - }) => { - let address = match endpoint { - ConnectedPoint::Dialer { address, .. } => Some(address.clone()), - ConnectedPoint::Listener { .. } => None, - }; - debug!( + fn on_swarm_event(&mut self, event: FromSwarm) { + match event { + FromSwarm::ConnectionEstablished(ConnectionEstablished { + peer_id, + connection_id, + endpoint, + other_established, + .. + }) => { + let address = match endpoint { + ConnectedPoint::Dialer { address, .. } => Some(address.clone()), + ConnectedPoint::Listener { .. } => None, + }; + debug!( "connection established with peer '{}' found at '{:?}'; peer has {} active connections", peer_id, address, other_established ); - let peer_id = PeerId(peer_id); - let endpoint = endpoint.clone(); - let conn = match self.connected_peers.get_mut(&peer_id) { - Some(peer) => { - peer.connections.insert(connection_id, endpoint); + let peer_id = PeerId(peer_id); + let endpoint = endpoint.clone(); + let conn = match self.connected_peers.get_mut(&peer_id) { + Some(peer) => { + peer.connections.insert(connection_id, endpoint); - peer.clone() - } - None => { - self.connected_peers.insert( - peer_id, - ConnectedPeer { - peer_id, - connections: HashMap::from([(connection_id, endpoint)]), - }, - ); - self.connected_peers - .get(&peer_id) - .expect("We legit have a mutable reference") - .clone() - } - }; + peer.clone() + } + None => { + self.connected_peers.insert( + peer_id, + ConnectedPeer { + peer_id, + connections: HashMap::from([(connection_id, endpoint)]), + }, + ); + self.connected_peers + .get(&peer_id) + .expect("We legit have a mutable reference") + .clone() + } + }; - // TODO: Move this block onto into `connection.rs` -> will probs be required for the ConnectionEstablishmentPayload stuff - { - debug!("sending establishment request to peer '{}'", peer_id); - if other_established == 0 { - let manager = self.manager.clone(); - tokio::spawn(async move { - manager - .emit(ManagerStreamAction::Event(Event::PeerConnected(conn))) - .await; - }); - } - } - } - FromSwarm::ConnectionClosed(ConnectionClosed { - peer_id, - connection_id, - .. - }) => { - let peer_id = PeerId(peer_id); - match self.connected_peers.get_mut(&peer_id) { - Some(peer) => { - if peer.connections.len() == 1 { - let conn = self.connected_peers.remove(&peer_id).expect("Literally impossible. We have a mutable reference to it, no shot it's already been removed."); - debug!("Disconnected from peer '{}'", conn.peer_id); + // TODO: Move this block onto into `connection.rs` -> will probs be required for the ConnectionEstablishmentPayload stuff + { + debug!("sending establishment request to peer '{}'", peer_id); + if other_established == 0 { + let manager = self.manager.clone(); + tokio::spawn(async move { + manager + .emit(ManagerStreamAction::Event(Event::PeerConnected(conn))) + .await; + }); + } + } + } + FromSwarm::ConnectionClosed(ConnectionClosed { + peer_id, + connection_id, + .. + }) => { + let peer_id = PeerId(peer_id); + match self.connected_peers.get_mut(&peer_id) { + Some(peer) => { + if peer.connections.len() == 1 { + let conn = self.connected_peers.remove(&peer_id).expect("Literally impossible. We have a mutable reference to it, no shot it's already been removed."); + debug!("Disconnected from peer '{}'", conn.peer_id); - let manager = self.manager.clone(); - tokio::spawn(async move { - manager - .emit(ManagerStreamAction::Event(Event::PeerDisconnected( - conn.peer_id.clone(), - ))) - .await; - }); - } else { - peer.connections.remove(&connection_id); - } - } - None => { - warn!( + let manager = self.manager.clone(); + tokio::spawn(async move { + manager + .emit(ManagerStreamAction::Event(Event::PeerDisconnected( + conn.peer_id, + ))) + .await; + }); + } else { + peer.connections.remove(&connection_id); + } + } + None => { + warn!( "Received connection closed event for peer '{}' but no connection was found!", peer_id ); - } - } - } - FromSwarm::AddressChange(_event) => { - // TODO: Reenable? - // let new_address = match event.new { - // ConnectedPoint::Dialer { address, .. } => Some(address.clone()), - // ConnectedPoint::Listener { .. } => None, - // }; - // let connected = self - // .manager - // .connections - // .blocking_read() - // .get_mut(&PeerId(event.peer_id)) - // .expect("Address change can only happen on an established connection."); + } + } + } + FromSwarm::AddressChange(_event) => { + // TODO: Reenable? + // let new_address = match event.new { + // ConnectedPoint::Dialer { address, .. } => Some(address.clone()), + // ConnectedPoint::Listener { .. } => None, + // }; + // let connected = self + // .manager + // .connections + // .blocking_read() + // .get_mut(&PeerId(event.peer_id)) + // .expect("Address change can only happen on an established connection."); - // let connection = connected - // .connections - // .iter_mut() - // .find(|c| c.id == event.connection_id) - // .expect("Address change can only happen on an established connection."); - // connection.address = new_address; - } - FromSwarm::DialFailure(event) => { - if let Some(peer_id) = event.peer_id { - debug!("Dialing failure to peer '{}': {:?}", peer_id, event.error); + // let connection = connected + // .connections + // .iter_mut() + // .find(|c| c.id == event.connection_id) + // .expect("Address change can only happen on an established connection."); + // connection.address = new_address; + } + FromSwarm::DialFailure(event) => { + if let Some(peer_id) = event.peer_id { + debug!("Dialing failure to peer '{}': {:?}", peer_id, event.error); - // TODO - // If there are pending outgoing requests when a dial failure occurs, - // it is implied that we are not connected to the peer, since pending - // outgoing requests are drained when a connection is established and - // only created when a peer is not connected when a request is made. - // Thus these requests must be considered failed, even if there is - // another, concurrent dialing attempt ongoing. - // if let Some(pending) = self.pending_outbound_requests.remove(&peer_id) { - // for request in pending { - // self.pending_events - // .push_back(NetworkBehaviourAction::GenerateEvent( - // Event::OutboundFailure { - // peer_id, - // request_id: request.request_id, - // error: OutboundFailure::DialFailure, - // }, - // )); - // } - // } - } - } - FromSwarm::ListenFailure(_) - | FromSwarm::NewListener(_) - | FromSwarm::NewListenAddr(_) - | FromSwarm::ExpiredListenAddr(_) - | FromSwarm::ListenerError(_) - | FromSwarm::ListenerClosed(_) - | FromSwarm::NewExternalAddr(_) - | FromSwarm::ExpiredExternalAddr(_) => {} - } - } + // TODO + // If there are pending outgoing requests when a dial failure occurs, + // it is implied that we are not connected to the peer, since pending + // outgoing requests are drained when a connection is established and + // only created when a peer is not connected when a request is made. + // Thus these requests must be considered failed, even if there is + // another, concurrent dialing attempt ongoing. + // if let Some(pending) = self.pending_outbound_requests.remove(&peer_id) { + // for request in pending { + // self.pending_events + // .push_back(NetworkBehaviourAction::GenerateEvent( + // Event::OutboundFailure { + // peer_id, + // request_id: request.request_id, + // error: OutboundFailure::DialFailure, + // }, + // )); + // } + // } + } + } + FromSwarm::ListenFailure(_) + | FromSwarm::NewListener(_) + | FromSwarm::NewListenAddr(_) + | FromSwarm::ExpiredListenAddr(_) + | FromSwarm::ListenerError(_) + | FromSwarm::ListenerClosed(_) + | FromSwarm::NewExternalAddr(_) + | FromSwarm::ExpiredExternalAddr(_) => {} + } + } - fn on_connection_handler_event( - &mut self, - _peer_id: libp2p::PeerId, - _connection: ConnectionId, - _event: as ConnectionHandler>::OutEvent, - ) { - todo!(); - } + fn on_connection_handler_event( + &mut self, + _peer_id: libp2p::PeerId, + _connection: ConnectionId, + _event: as ConnectionHandler>::OutEvent, + ) { + todo!(); + } - fn poll( - &mut self, - _: &mut Context<'_>, - _: &mut impl PollParameters, - ) -> Poll>> { - if let Some(ev) = self.pending_events.pop_front() { - return Poll::Ready(ev); - } else if self.pending_events.capacity() > EMPTY_QUEUE_SHRINK_THRESHOLD { - self.pending_events.shrink_to_fit(); - } + fn poll( + &mut self, + _: &mut Context<'_>, + _: &mut impl PollParameters, + ) -> Poll>> { + if let Some(ev) = self.pending_events.pop_front() { + return Poll::Ready(ev); + } else if self.pending_events.capacity() > EMPTY_QUEUE_SHRINK_THRESHOLD { + self.pending_events.shrink_to_fit(); + } - Poll::Pending - } + Poll::Pending + } } diff --git a/crates/p2p/src/spacetime/connection.rs b/crates/p2p/src/spacetime/connection.rs index a3ad13797..6a660e878 100644 --- a/crates/p2p/src/spacetime/connection.rs +++ b/crates/p2p/src/spacetime/connection.rs @@ -1,16 +1,16 @@ use libp2p::swarm::{ - handler::{ - ConnectionEvent, ConnectionHandler, ConnectionHandlerEvent, ConnectionHandlerUpgrErr, - KeepAlive, - }, - SubstreamProtocol, + handler::{ + ConnectionEvent, ConnectionHandler, ConnectionHandlerEvent, ConnectionHandlerUpgrErr, + KeepAlive, + }, + SubstreamProtocol, }; use std::{ - collections::VecDeque, - io, - sync::Arc, - task::{Context, Poll}, - time::Duration, + collections::VecDeque, + io, + sync::Arc, + task::{Context, Poll}, + time::Duration, }; use tracing::error; @@ -21,114 +21,115 @@ use super::{InboundProtocol, OutboundProtocol, OutboundRequest, EMPTY_QUEUE_SHRI // TODO: Probs change this based on the ConnectionEstablishmentPayload const SUBSTREAM_TIMEOUT: Duration = Duration::from_secs(10); // TODO: Tune value +#[allow(clippy::type_complexity)] pub struct SpaceTimeConnection { - peer_id: PeerId, - manager: Arc>, - pending_events: VecDeque< - ConnectionHandlerEvent< - OutboundProtocol, - ::OutboundOpenInfo, - ::OutEvent, - ::Error, - >, - >, + peer_id: PeerId, + manager: Arc>, + pending_events: VecDeque< + ConnectionHandlerEvent< + OutboundProtocol, + ::OutboundOpenInfo, + ::OutEvent, + ::Error, + >, + >, } impl SpaceTimeConnection { - pub(super) fn new(peer_id: PeerId, manager: Arc>) -> Self { - Self { - peer_id, - manager, - pending_events: VecDeque::new(), - } - } + pub(super) fn new(peer_id: PeerId, manager: Arc>) -> Self { + Self { + peer_id, + manager, + pending_events: VecDeque::new(), + } + } } // pub enum Connection impl ConnectionHandler for SpaceTimeConnection { - type InEvent = OutboundRequest; - type OutEvent = (); - type Error = ConnectionHandlerUpgrErr; - type InboundProtocol = InboundProtocol; - type OutboundProtocol = OutboundProtocol; - type OutboundOpenInfo = (); - type InboundOpenInfo = (); + type InEvent = OutboundRequest; + type OutEvent = (); + type Error = ConnectionHandlerUpgrErr; + type InboundProtocol = InboundProtocol; + type OutboundProtocol = OutboundProtocol; + type OutboundOpenInfo = (); + type InboundOpenInfo = (); - fn listen_protocol(&self) -> SubstreamProtocol { - SubstreamProtocol::new( - InboundProtocol { - peer_id: self.peer_id, - manager: self.manager.clone(), - }, - (), - ) - .with_timeout(SUBSTREAM_TIMEOUT) - } + fn listen_protocol(&self) -> SubstreamProtocol { + SubstreamProtocol::new( + InboundProtocol { + peer_id: self.peer_id, + manager: self.manager.clone(), + }, + (), + ) + .with_timeout(SUBSTREAM_TIMEOUT) + } - fn on_behaviour_event(&mut self, req: Self::InEvent) { - // TODO: Working keep alives - // self.keep_alive = KeepAlive::Yes; - // self.outbound.push_back(request); + fn on_behaviour_event(&mut self, req: Self::InEvent) { + // TODO: Working keep alives + // self.keep_alive = KeepAlive::Yes; + // self.outbound.push_back(request); - self.pending_events - .push_back(ConnectionHandlerEvent::OutboundSubstreamRequest { - protocol: SubstreamProtocol::new( - OutboundProtocol(self.manager.application_name.clone(), req), - (), - ) // TODO: Use `info` here maybe to pass into about the client. Idk? - .with_timeout(SUBSTREAM_TIMEOUT), - }); - } + self.pending_events + .push_back(ConnectionHandlerEvent::OutboundSubstreamRequest { + protocol: SubstreamProtocol::new( + OutboundProtocol(self.manager.application_name, req), + (), + ) // TODO: Use `info` here maybe to pass into about the client. Idk? + .with_timeout(SUBSTREAM_TIMEOUT), + }); + } - fn connection_keep_alive(&self) -> KeepAlive { - KeepAlive::Yes // TODO: Make this work how the old one did with storing it on `self` and updating on events - } + fn connection_keep_alive(&self) -> KeepAlive { + KeepAlive::Yes // TODO: Make this work how the old one did with storing it on `self` and updating on events + } - fn poll( - &mut self, - _cx: &mut Context<'_>, - ) -> Poll< - ConnectionHandlerEvent< - Self::OutboundProtocol, - Self::OutboundOpenInfo, - Self::OutEvent, - Self::Error, - >, - > { - if let Some(event) = self.pending_events.pop_front() { - return Poll::Ready(event); - } else if self.pending_events.capacity() > EMPTY_QUEUE_SHRINK_THRESHOLD { - self.pending_events.shrink_to_fit(); - } + fn poll( + &mut self, + _cx: &mut Context<'_>, + ) -> Poll< + ConnectionHandlerEvent< + Self::OutboundProtocol, + Self::OutboundOpenInfo, + Self::OutEvent, + Self::Error, + >, + > { + if let Some(event) = self.pending_events.pop_front() { + return Poll::Ready(event); + } else if self.pending_events.capacity() > EMPTY_QUEUE_SHRINK_THRESHOLD { + self.pending_events.shrink_to_fit(); + } - Poll::Pending - } + Poll::Pending + } - // TODO: Which level we doing error handler?. On swarm, on Behavior or here??? - fn on_connection_event( - &mut self, - event: ConnectionEvent< - Self::InboundProtocol, - Self::OutboundProtocol, - Self::InboundOpenInfo, - Self::OutboundOpenInfo, - >, - ) { - match event { - ConnectionEvent::FullyNegotiatedInbound(_) => {} - ConnectionEvent::FullyNegotiatedOutbound(_) => {} - ConnectionEvent::DialUpgradeError(event) => { - error!("DialUpgradeError: {:#?}", event.error); - } - ConnectionEvent::ListenUpgradeError(event) => { - error!("DialUpgradeError: {:#?}", event.error); + // TODO: Which level we doing error handler?. On swarm, on Behavior or here??? + fn on_connection_event( + &mut self, + event: ConnectionEvent< + Self::InboundProtocol, + Self::OutboundProtocol, + Self::InboundOpenInfo, + Self::OutboundOpenInfo, + >, + ) { + match event { + ConnectionEvent::FullyNegotiatedInbound(_) => {} + ConnectionEvent::FullyNegotiatedOutbound(_) => {} + ConnectionEvent::DialUpgradeError(event) => { + error!("DialUpgradeError: {:#?}", event.error); + } + ConnectionEvent::ListenUpgradeError(event) => { + error!("DialUpgradeError: {:#?}", event.error); - // TODO: If `event.error` close connection cause we don't "speak the same language"! - } - ConnectionEvent::AddressChange(_) => { - // TODO: Should we be telling `SpaceTime` to update it's info here or is it also getting this event? - } - } - } + // TODO: If `event.error` close connection cause we don't "speak the same language"! + } + ConnectionEvent::AddressChange(_) => { + // TODO: Should we be telling `SpaceTime` to update it's info here or is it also getting this event? + } + } + } } diff --git a/crates/p2p/src/spacetime/libp2p.rs b/crates/p2p/src/spacetime/libp2p.rs index 39b69c9ba..9a2032c1d 100644 --- a/crates/p2p/src/spacetime/libp2p.rs +++ b/crates/p2p/src/spacetime/libp2p.rs @@ -4,7 +4,7 @@ pub struct SpaceTimeProtocolName(pub &'static [u8]); impl libp2p::core::ProtocolName for SpaceTimeProtocolName { - fn protocol_name(&self) -> &[u8] { - self.0 - } + fn protocol_name(&self) -> &[u8] { + self.0 + } } diff --git a/crates/p2p/src/spacetime/message.rs b/crates/p2p/src/spacetime/message.rs index 47296d07c..fad5a16f3 100644 --- a/crates/p2p/src/spacetime/message.rs +++ b/crates/p2p/src/spacetime/message.rs @@ -3,9 +3,9 @@ use serde::{Deserialize, Serialize}; /// TODO #[derive(Debug, Clone, Serialize, Deserialize)] pub enum SpaceTimeMessage { - /// Establish the connection - Establish, + /// Establish the connection + Establish, - /// Send data on behalf of application - Application(Vec), + /// Send data on behalf of application + Application(Vec), } diff --git a/crates/p2p/src/spacetime/proto_inbound.rs b/crates/p2p/src/spacetime/proto_inbound.rs index ab12ce283..842d62b3d 100644 --- a/crates/p2p/src/spacetime/proto_inbound.rs +++ b/crates/p2p/src/spacetime/proto_inbound.rs @@ -7,38 +7,39 @@ use crate::{Manager, ManagerStreamAction, Metadata, PeerId, PeerMessageEvent}; use super::{SpaceTimeProtocolName, SpaceTimeStream}; pub struct InboundProtocol { - pub(crate) peer_id: PeerId, - pub(crate) manager: Arc>, + pub(crate) peer_id: PeerId, + pub(crate) manager: Arc>, } impl UpgradeInfo for InboundProtocol { - type Info = SpaceTimeProtocolName; - type InfoIter = [Self::Info; 1]; + type Info = SpaceTimeProtocolName; + type InfoIter = [Self::Info; 1]; - fn protocol_info(&self) -> Self::InfoIter { - [SpaceTimeProtocolName(self.manager.application_name)] - } + fn protocol_info(&self) -> Self::InfoIter { + [SpaceTimeProtocolName(self.manager.application_name)] + } } impl InboundUpgrade for InboundProtocol { - type Output = (); - type Error = (); - type Future = Pin> + Send + 'static>>; + type Output = (); + type Error = (); + type Future = Pin> + Send + 'static>>; - fn upgrade_inbound(self, io: NegotiatedSubstream, _: Self::Info) -> Self::Future { - Box::pin(async move { - Ok(self - .manager - .emit(ManagerStreamAction::Event( - PeerMessageEvent { - peer_id: self.peer_id, - manager: self.manager.clone(), - stream: SpaceTimeStream::new(io), - _priv: (), - } - .into(), - )) - .await) - }) - } + fn upgrade_inbound(self, io: NegotiatedSubstream, _: Self::Info) -> Self::Future { + Box::pin(async move { + self.manager + .emit(ManagerStreamAction::Event( + PeerMessageEvent { + peer_id: self.peer_id, + manager: self.manager.clone(), + stream: SpaceTimeStream::new(io), + _priv: (), + } + .into(), + )) + .await; + + Ok(()) + }) + } } diff --git a/crates/p2p/src/spacetime/proto_outbound.rs b/crates/p2p/src/spacetime/proto_outbound.rs index cdd249b08..efb7075c4 100644 --- a/crates/p2p/src/spacetime/proto_outbound.rs +++ b/crates/p2p/src/spacetime/proto_outbound.rs @@ -8,47 +8,47 @@ use super::{SpaceTimeProtocolName, SpaceTimeStream}; #[derive(Debug)] pub enum OutboundRequest { - Data(Vec), - Stream(oneshot::Sender), + Data(Vec), + Stream(oneshot::Sender), } pub struct OutboundProtocol(pub(crate) &'static [u8], pub(crate) OutboundRequest); impl UpgradeInfo for OutboundProtocol { - type Info = SpaceTimeProtocolName; - type InfoIter = [Self::Info; 1]; + type Info = SpaceTimeProtocolName; + type InfoIter = [Self::Info; 1]; - fn protocol_info(&self) -> Self::InfoIter { - [SpaceTimeProtocolName(self.0)] - } + fn protocol_info(&self) -> Self::InfoIter { + [SpaceTimeProtocolName(self.0)] + } } impl OutboundUpgrade for OutboundProtocol { - type Output = (); - type Error = (); - type Future = Ready>; + type Output = (); + type Error = (); + type Future = Ready>; - fn upgrade_outbound(self, io: NegotiatedSubstream, _protocol: Self::Info) -> Self::Future { - let mut stream = SpaceTimeStream::new(io); - match self.1 { - OutboundRequest::Data(data) => { - tokio::spawn(async move { - if let Err(err) = stream.write_all(&data).await { - // TODO: Print the peer which we failed to send to here - error!("Error sending broadcast: {:?}", err); - } - stream.flush().await.unwrap(); - stream.close().await.unwrap(); - // TODO: We close the connection here without waiting for a response. - // TODO: If the other side's user-code doesn't account for that on this specific message they will error. - // TODO: Add an abstraction so the user can't respond to fixed size messages. - }); - } - OutboundRequest::Stream(sender) => { - sender.send(stream).unwrap(); - } - } + fn upgrade_outbound(self, io: NegotiatedSubstream, _protocol: Self::Info) -> Self::Future { + let mut stream = SpaceTimeStream::new(io); + match self.1 { + OutboundRequest::Data(data) => { + tokio::spawn(async move { + if let Err(err) = stream.write_all(&data).await { + // TODO: Print the peer which we failed to send to here + error!("Error sending broadcast: {:?}", err); + } + stream.flush().await.unwrap(); + stream.close().await.unwrap(); + // TODO: We close the connection here without waiting for a response. + // TODO: If the other side's user-code doesn't account for that on this specific message they will error. + // TODO: Add an abstraction so the user can't respond to fixed size messages. + }); + } + OutboundRequest::Stream(sender) => { + sender.send(stream).unwrap(); + } + } - ready(Ok(())) - } + ready(Ok(())) + } } diff --git a/crates/p2p/src/spacetime/stream.rs b/crates/p2p/src/spacetime/stream.rs index 213ab6a01..51924bced 100644 --- a/crates/p2p/src/spacetime/stream.rs +++ b/crates/p2p/src/spacetime/stream.rs @@ -1,7 +1,7 @@ use std::{ - io, - pin::Pin, - task::{Context, Poll}, + io, + pin::Pin, + task::{Context, Poll}, }; use libp2p::{futures::AsyncWriteExt, swarm::NegotiatedSubstream}; @@ -14,23 +14,23 @@ pub struct SpaceTimeStream(Compat); // TODO: Utils for sending msgpack and stuff over the stream. -> Have a max size of reading buffers so we are less susceptible to DoS attacks. impl SpaceTimeStream { - pub fn new(io: NegotiatedSubstream) -> Self { - Self(io.compat()) - } + pub fn new(io: NegotiatedSubstream) -> Self { + Self(io.compat()) + } - pub async fn close(self) -> Result<(), io::Error> { - self.0.into_inner().close().await - } + pub async fn close(self) -> Result<(), io::Error> { + self.0.into_inner().close().await + } } impl AsyncRead for SpaceTimeStream { - fn poll_read( - self: Pin<&mut Self>, - cx: &mut Context<'_>, - buf: &mut ReadBuf<'_>, - ) -> Poll> { - Pin::new(&mut self.get_mut().0).poll_read(cx, buf) - } + fn poll_read( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &mut ReadBuf<'_>, + ) -> Poll> { + Pin::new(&mut self.get_mut().0).poll_read(cx, buf) + } } // impl AsyncBufRead for SpaceTimeStream { @@ -44,19 +44,19 @@ impl AsyncRead for SpaceTimeStream { // } impl AsyncWrite for SpaceTimeStream { - fn poll_write( - self: Pin<&mut Self>, - cx: &mut Context<'_>, - buf: &[u8], - ) -> Poll> { - Pin::new(&mut self.get_mut().0).poll_write(cx, buf) - } + fn poll_write( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + Pin::new(&mut self.get_mut().0).poll_write(cx, buf) + } - fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - Pin::new(&mut self.get_mut().0).poll_flush(cx) - } + fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + Pin::new(&mut self.get_mut().0).poll_flush(cx) + } - fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - Pin::new(&mut self.get_mut().0).poll_shutdown(cx) - } + fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + Pin::new(&mut self.get_mut().0).poll_shutdown(cx) + } } diff --git a/crates/p2p/src/utils/async_fn.rs b/crates/p2p/src/utils/async_fn.rs index 88cbcf261..d1d6129ec 100644 --- a/crates/p2p/src/utils/async_fn.rs +++ b/crates/p2p/src/utils/async_fn.rs @@ -3,17 +3,17 @@ use std::future::Future; // A trait which allows me to represent a closure, it's future return type and the futures return type with only a single generic. pub trait AsyncFn where - Self: Fn() -> Self::Future + Send + Sync + 'static, + Self: Fn() -> Self::Future + Send + Sync + 'static, { - type Output; - type Future: Future::Output> + Send; + type Output; + type Future: Future::Output> + Send; } impl AsyncFn for TFunc where - TFut: Future + Send, - TFunc: Fn() -> TFut + Send + Sync + 'static, + TFut: Future + Send, + TFunc: Fn() -> TFut + Send + Sync + 'static, { - type Output = TOutput; - type Future = TFut; + type Output = TOutput; + type Future = TFut; } diff --git a/crates/p2p/src/utils/keypair.rs b/crates/p2p/src/utils/keypair.rs index 7cd527d1e..978b33f1f 100644 --- a/crates/p2p/src/utils/keypair.rs +++ b/crates/p2p/src/utils/keypair.rs @@ -5,42 +5,42 @@ use serde::{Deserialize, Serialize}; pub struct Keypair(libp2p::identity::Keypair); impl Keypair { - pub fn generate() -> Self { - Self(libp2p::identity::Keypair::generate_ed25519()) - } + pub fn generate() -> Self { + Self(libp2p::identity::Keypair::generate_ed25519()) + } - pub fn public(&self) -> PublicKey { - self.0.public() - } + pub fn public(&self) -> PublicKey { + self.0.public() + } - pub fn inner(&self) -> &libp2p::identity::Keypair { - &self.0 - } + pub fn inner(&self) -> &libp2p::identity::Keypair { + &self.0 + } } impl Serialize for Keypair { - fn serialize(&self, serializer: S) -> Result - where - S: serde::Serializer, - { - match &self.0 { - libp2p::identity::Keypair::Ed25519(keypair) => { - serializer.serialize_bytes(&keypair.encode()) - } - #[allow(unreachable_patterns)] - _ => unreachable!(), - } - } + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + match &self.0 { + libp2p::identity::Keypair::Ed25519(keypair) => { + serializer.serialize_bytes(&keypair.encode()) + } + #[allow(unreachable_patterns)] + _ => unreachable!(), + } + } } impl<'de> Deserialize<'de> for Keypair { - fn deserialize(deserializer: D) -> Result - where - D: serde::Deserializer<'de>, - { - let mut bytes = Vec::::deserialize(deserializer)?; - Ok(Self(libp2p::identity::Keypair::Ed25519( - ed25519::Keypair::decode(bytes.as_mut_slice()).map_err(serde::de::Error::custom)?, - ))) - } + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let mut bytes = Vec::::deserialize(deserializer)?; + Ok(Self(libp2p::identity::Keypair::Ed25519( + ed25519::Keypair::decode(bytes.as_mut_slice()).map_err(serde::de::Error::custom)?, + ))) + } } diff --git a/crates/p2p/src/utils/metadata.rs b/crates/p2p/src/utils/metadata.rs index 4bbf538ce..3373c805a 100644 --- a/crates/p2p/src/utils/metadata.rs +++ b/crates/p2p/src/utils/metadata.rs @@ -2,9 +2,9 @@ use std::collections::HashMap; /// this trait must be implemented for the metadata type to allow it to be converted to MDNS DNS records. pub trait Metadata: Clone + Send + Sync + 'static { - fn to_hashmap(self) -> HashMap; + fn to_hashmap(self) -> HashMap; - fn from_hashmap(data: &HashMap) -> Result - where - Self: Sized; + fn from_hashmap(data: &HashMap) -> Result + where + Self: Sized; } diff --git a/crates/p2p/src/utils/multiaddr.rs b/crates/p2p/src/utils/multiaddr.rs index 4fc46a77f..3efaa576c 100644 --- a/crates/p2p/src/utils/multiaddr.rs +++ b/crates/p2p/src/utils/multiaddr.rs @@ -5,41 +5,41 @@ use libp2p::{multiaddr::Protocol, Multiaddr}; // TODO: Turn these into From/Into impls on a wrapper type pub(crate) fn quic_multiaddr_to_socketaddr(m: Multiaddr) -> Result { - let mut addr_parts = m.iter(); + let mut addr_parts = m.iter(); - let addr = match addr_parts.next() { - Some(Protocol::Ip4(addr)) => IpAddr::V4(addr), - Some(Protocol::Ip6(addr)) => IpAddr::V6(addr), - Some(proto) => { - return Err(format!( - "Invalid multiaddr. Segment 1 found protocol 'Ip4' or 'Ip6' but found '{}'", - proto - )) - } - None => return Err(format!("Invalid multiaddr. Segment 1 missing")), - }; + let addr = match addr_parts.next() { + Some(Protocol::Ip4(addr)) => IpAddr::V4(addr), + Some(Protocol::Ip6(addr)) => IpAddr::V6(addr), + Some(proto) => { + return Err(format!( + "Invalid multiaddr. Segment 1 found protocol 'Ip4' or 'Ip6' but found '{}'", + proto + )) + } + None => return Err("Invalid multiaddr. Segment 1 missing".to_string()), + }; - let port = match addr_parts.next() { - Some(Protocol::Udp(port)) => port, - Some(proto) => { - return Err(format!( - "Invalid multiaddr. Segment 2 expected protocol 'Udp' but found '{}'", - proto - )) - } - None => return Err(format!("Invalid multiaddr. Segment 2 missing")), - }; + let port = match addr_parts.next() { + Some(Protocol::Udp(port)) => port, + Some(proto) => { + return Err(format!( + "Invalid multiaddr. Segment 2 expected protocol 'Udp' but found '{}'", + proto + )) + } + None => return Err("Invalid multiaddr. Segment 2 missing".to_string()), + }; - Ok(SocketAddr::new(addr, port)) + Ok(SocketAddr::new(addr, port)) } pub(crate) fn socketaddr_to_quic_multiaddr(m: &SocketAddr) -> Multiaddr { - let mut addr = Multiaddr::empty(); - match m { - SocketAddr::V4(ip) => addr.push(Protocol::Ip4(*ip.ip())), - SocketAddr::V6(ip) => addr.push(Protocol::Ip6(*ip.ip())), - } - addr.push(Protocol::Udp(m.port())); - addr.push(Protocol::QuicV1); - addr + let mut addr = Multiaddr::empty(); + match m { + SocketAddr::V4(ip) => addr.push(Protocol::Ip4(*ip.ip())), + SocketAddr::V6(ip) => addr.push(Protocol::Ip6(*ip.ip())), + } + addr.push(Protocol::Udp(m.port())); + addr.push(Protocol::QuicV1); + addr }