From 462c3723ae542231155128a2e6aa21106ed057f2 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 26 Sep 2022 17:24:11 +0200 Subject: [PATCH 01/16] create new tcp_config since clone() is not supported anymore --- fuel-p2p/src/config.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/fuel-p2p/src/config.rs b/fuel-p2p/src/config.rs index 1b6432d3f17..92e7ea8f7f7 100644 --- a/fuel-p2p/src/config.rs +++ b/fuel-p2p/src/config.rs @@ -123,7 +123,9 @@ pub(crate) async fn build_transport( ) -> Boxed<(PeerId, StreamMuxerBox)> { let transport = { let tcp = libp2p::tcp::TcpConfig::new().nodelay(true); - let ws_tcp = libp2p::websocket::WsConfig::new(tcp.clone()).or_transport(tcp); + let ws_tcp = + libp2p::websocket::WsConfig::new(libp2p::tcp::TcpConfig::new().nodelay(true)) + .or_transport(tcp); libp2p::dns::DnsConfig::system(ws_tcp).await.unwrap() }; From 042894fd2e53cec53fc42380800282b6c6933013 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 26 Sep 2022 17:55:17 +0200 Subject: [PATCH 02/16] import and generate tcp transport as per new transport updates --- fuel-p2p/src/config.rs | 14 +++++++++++--- fuel-p2p/src/discovery.rs | 8 +++----- fuel-p2p/src/peer_info.rs | 6 ++---- 3 files changed, 16 insertions(+), 12 deletions(-) diff --git a/fuel-p2p/src/config.rs b/fuel-p2p/src/config.rs index 92e7ea8f7f7..9638c99511a 100644 --- a/fuel-p2p/src/config.rs +++ b/fuel-p2p/src/config.rs @@ -9,6 +9,10 @@ use libp2p::{ }, mplex, noise, + tcp::{ + GenTcpConfig, + TcpTransport, + }, yamux, Multiaddr, PeerId, @@ -122,10 +126,14 @@ pub(crate) async fn build_transport( local_keypair: Keypair, ) -> Boxed<(PeerId, StreamMuxerBox)> { let transport = { - let tcp = libp2p::tcp::TcpConfig::new().nodelay(true); + let generate_tcp_transpot = + || TcpTransport::new(GenTcpConfig::new().port_reuse(true).nodelay(true)); + + let tcp = generate_tcp_transpot(); + let ws_tcp = - libp2p::websocket::WsConfig::new(libp2p::tcp::TcpConfig::new().nodelay(true)) - .or_transport(tcp); + libp2p::websocket::WsConfig::new(generate_tcp_transpot()).or_transport(tcp); + libp2p::dns::DnsConfig::system(ws_tcp).await.unwrap() }; diff --git a/fuel-p2p/src/discovery.rs b/fuel-p2p/src/discovery.rs index d438b51c98e..a241844eb2e 100644 --- a/fuel-p2p/src/discovery.rs +++ b/fuel-p2p/src/discovery.rs @@ -4,10 +4,8 @@ use futures_timer::Delay; use ip_network::IpNetwork; use libp2p::{ core::{ - connection::{ - ConnectionId, - ListenerId, - }, + connection::ConnectionId, + transport::ListenerId, ConnectedPoint, }, kad::{ @@ -398,7 +396,7 @@ mod tests { .into_authentic(&keypair) .unwrap(); - let transport = core::transport::MemoryTransport + let transport = core::transport::MemoryTransport::default() .upgrade(core::upgrade::Version::V1) .authenticate(noise::NoiseConfig::xx(noise_keys).into_authenticated()) .multiplex(yamux::YamuxConfig::default()) diff --git a/fuel-p2p/src/peer_info.rs b/fuel-p2p/src/peer_info.rs index b16ce5091f1..00be82c9b80 100644 --- a/fuel-p2p/src/peer_info.rs +++ b/fuel-p2p/src/peer_info.rs @@ -1,11 +1,9 @@ use crate::config::P2PConfig; use libp2p::{ core::{ - connection::{ - ConnectionId, - ListenerId, - }, + connection::ConnectionId, either::EitherOutput, + transport::ListenerId, ConnectedPoint, PublicKey, }, From a4c5fafd531672e7cfce599e4690448e8318067d Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Fri, 30 Sep 2022 10:39:07 +0200 Subject: [PATCH 03/16] move state and logic to service from behaviour --- fuel-p2p/src/behavior.rs | 379 +++++++---------------------------- fuel-p2p/src/orchestrator.rs | 64 +++--- fuel-p2p/src/service.rs | 364 ++++++++++++++++++++++++++------- 3 files changed, 403 insertions(+), 404 deletions(-) diff --git a/fuel-p2p/src/behavior.rs b/fuel-p2p/src/behavior.rs index e12d69d5f7f..10c826e26f6 100644 --- a/fuel-p2p/src/behavior.rs +++ b/fuel-p2p/src/behavior.rs @@ -7,15 +7,8 @@ use crate::{ DiscoveryEvent, }, gossipsub::{ - self, - messages::{ - GossipsubBroadcastRequest, - GossipsubMessage as FuelGossipsubMessage, - }, - topics::{ - GossipTopic, - GossipsubTopics, - }, + build_gossipsub, + topics::GossipTopic, }, peer_info::{ PeerInfo, @@ -24,11 +17,7 @@ use crate::{ }, request_response::messages::{ IntermediateResponse, - OutboundResponse, RequestMessage, - ResponseChannelItem, - ResponseError, - ResponseMessage, }, }; use libp2p::{ @@ -40,7 +29,6 @@ use libp2p::{ Gossipsub, GossipsubEvent, MessageId, - TopicHash, }, request_response::{ ProtocolSupport, @@ -48,64 +36,49 @@ use libp2p::{ RequestResponse, RequestResponseConfig, RequestResponseEvent, - RequestResponseMessage, ResponseChannel, }, - swarm::{ - NetworkBehaviour, - NetworkBehaviourAction, - NetworkBehaviourEventProcess, - PollParameters, - }, + Multiaddr, NetworkBehaviour, PeerId, }; use std::{ - collections::{ - HashMap, - VecDeque, - }, - task::{ - Context, - Poll, - }, -}; -use tracing::{ - debug, - warn, + collections::HashMap, + marker::PhantomData, }; -#[allow(clippy::large_enum_variant)] -#[derive(Debug, Clone)] -pub enum FuelBehaviourEvent { - PeerConnected(PeerId), - PeerDisconnected(PeerId), - PeerIdentified(PeerId), - PeerInfoUpdated(PeerId), - GossipsubMessage { - peer_id: PeerId, - topic_hash: TopicHash, - message: FuelGossipsubMessage, - }, - RequestMessage { - request_id: RequestId, - request_message: RequestMessage, - }, +#[derive(Debug)] +pub struct FuelBehaviourEvent { + event: InnerBehaviourEvent, + codec: PhantomData, } -/// Holds additional Network data for FuelBehavior #[derive(Debug)] -struct NetworkMetadata { - gossipsub_topics: GossipsubTopics, +pub enum InnerBehaviourEvent { + Discovery(DiscoveryEvent), + PeerInfo(PeerInfoEvent), + Gossipsub(GossipsubEvent), + RequestResponse(RequestResponseEvent), +} + +impl FuelBehaviourEvent { + fn new(event: InnerBehaviourEvent) -> Self { + Self { + event, + codec: PhantomData, + } + } +} + +impl From> for InnerBehaviourEvent { + fn from(fuel_event: FuelBehaviourEvent) -> Self { + fuel_event.event + } } /// Handles all p2p protocols needed for Fuel. #[derive(NetworkBehaviour)] -#[behaviour( - out_event = "FuelBehaviourEvent", - poll_method = "poll", - event_process = true -)] +#[behaviour(out_event = "FuelBehaviourEvent")] pub struct FuelBehaviour { /// Node discovery discovery: DiscoveryBehaviour, @@ -118,30 +91,6 @@ pub struct FuelBehaviour { /// RequestResponse protocol request_response: RequestResponse, - - /// Holds the Sender(s) part of the Oneshot Channel from the NetworkOrchestrator - /// Once the ResponseMessage is received from the p2p Network - /// It will send it to the NetworkOrchestrator via its unique Sender - #[behaviour(ignore)] - outbound_requests_table: HashMap, - - /// Holds the ResponseChannel(s) for the inbound requests from the p2p Network - /// Once the Response is prepared by the NetworkOrchestrator - /// It will send it to the specified Peer via its unique ResponseChannel - #[behaviour(ignore)] - inbound_requests_table: HashMap>, - - /// Double-ended queue of FuelBehaviour Events - #[behaviour(ignore)] - events: VecDeque, - - /// NetworkCodec used as for encoding and decoding of Gossipsub messages - #[behaviour(ignore)] - codec: Codec, - - /// Stores additional p2p network info - #[behaviour(ignore)] - network_metadata: NetworkMetadata, } impl FuelBehaviour { @@ -177,29 +126,32 @@ impl FuelBehaviour { req_res_config.set_connection_keep_alive(p2p_config.set_connection_keep_alive); let request_response = - RequestResponse::new(codec.clone(), req_res_protocol, req_res_config); - - let gossipsub_topics = GossipsubTopics::new(&p2p_config.network_name); - let network_metadata = NetworkMetadata { gossipsub_topics }; + RequestResponse::new(codec, req_res_protocol, req_res_config); Self { discovery: discovery_config.finish(), - gossipsub: gossipsub::build_gossipsub(&p2p_config.local_keypair, p2p_config), + gossipsub: build_gossipsub(&p2p_config.local_keypair, p2p_config), peer_info, request_response, - - outbound_requests_table: HashMap::default(), - inbound_requests_table: HashMap::default(), - events: VecDeque::default(), - codec, - network_metadata, } } - // Currently only used in testing hence `allow` - #[allow(dead_code)] - pub fn get_peer_info(&self, peer_id: &PeerId) -> Option<&PeerInfo> { - self.peer_info.get_peer_info(peer_id) + pub fn add_addresses_to_peer_info( + &mut self, + peer_id: &PeerId, + addresses: Vec, + ) { + self.peer_info.insert_peer_addresses(peer_id, addresses); + } + + pub fn add_addresses_to_discovery( + &mut self, + peer_id: &PeerId, + addresses: Vec, + ) { + for address in addresses { + self.discovery.add_address(peer_id, address.clone()); + } } pub fn get_peers(&self) -> &HashMap { @@ -208,16 +160,10 @@ impl FuelBehaviour { pub fn publish_message( &mut self, - message: GossipsubBroadcastRequest, + topic: GossipTopic, + encoded_data: Vec, ) -> Result { - let topic = self - .network_metadata - .gossipsub_topics - .get_gossipsub_topic(&message); - match self.codec.encode(message) { - Ok(encoded_data) => self.gossipsub.publish(topic, encoded_data), - Err(e) => Err(PublishError::TransformFailed(e)), - } + self.gossipsub.publish(topic, encoded_data) } pub fn subscribe_to_topic( @@ -231,227 +177,48 @@ impl FuelBehaviour { &mut self, message_request: RequestMessage, peer_id: PeerId, - channel_item: ResponseChannelItem, ) -> RequestId { - let request_id = self - .request_response - .send_request(&peer_id, message_request); - - self.outbound_requests_table - .insert(request_id, channel_item); - - request_id + self.request_response + .send_request(&peer_id, message_request) } pub fn send_response_msg( &mut self, - request_id: RequestId, - message: OutboundResponse, - ) -> Result<(), ResponseError> { - match ( - self.codec.convert_to_intermediate(&message), - self.inbound_requests_table.remove(&request_id), - ) { - (Ok(message), Some(channel)) => { - if self - .request_response - .send_response(channel, message) - .is_err() - { - debug!("Failed to send ResponseMessage for {:?}", request_id); - return Err(ResponseError::SendingResponseFailed) - } - } - (Ok(_), None) => { - debug!("ResponseChannel for {:?} does not exist!", request_id); - return Err(ResponseError::ResponseChannelDoesNotExist) - } - (Err(e), _) => { - debug!("Failed to convert to IntermediateResponse with {:?}", e); - return Err(ResponseError::ConversionToIntermediateFailed) - } - } - - Ok(()) + channel: ResponseChannel, + message: IntermediateResponse, + ) -> Result<(), IntermediateResponse> { + self.request_response.send_response(channel, message) } - // report events to the swarm - fn poll( - &mut self, - _cx: &mut Context, - _: &mut impl PollParameters, - ) -> Poll< - NetworkBehaviourAction< - ::OutEvent, - ::ConnectionHandler, - >, - > { - match self.events.pop_front() { - Some(event) => Poll::Ready(NetworkBehaviourAction::GenerateEvent(event)), - _ => Poll::Pending, - } - } - - /// Getter for outbound_requests_table - /// Used only in testing in `service.rs` + // Currently only used in testing, but should be useful for the NetworkOrchestrator API #[allow(dead_code)] - pub(super) fn get_outbound_requests_table( - &self, - ) -> &HashMap { - &self.outbound_requests_table + pub fn get_peer_info(&self, peer_id: &PeerId) -> Option<&PeerInfo> { + self.peer_info.get_peer_info(peer_id) } } -impl NetworkBehaviourEventProcess - for FuelBehaviour -{ - fn inject_event(&mut self, event: DiscoveryEvent) { - match event { - DiscoveryEvent::Connected(peer_id, addresses) => { - self.peer_info.insert_peer_addresses(&peer_id, addresses); - - self.events - .push_back(FuelBehaviourEvent::PeerConnected(peer_id)); - } - DiscoveryEvent::Disconnected(peer_id) => self - .events - .push_back(FuelBehaviourEvent::PeerDisconnected(peer_id)), - - _ => {} - } +impl From for FuelBehaviourEvent { + fn from(event: DiscoveryEvent) -> Self { + FuelBehaviourEvent::new(InnerBehaviourEvent::Discovery(event)) } } -impl NetworkBehaviourEventProcess - for FuelBehaviour -{ - fn inject_event(&mut self, event: PeerInfoEvent) { - match event { - PeerInfoEvent::PeerIdentified { peer_id, addresses } => { - for address in addresses { - self.discovery.add_address(&peer_id, address.clone()); - } - - self.events - .push_back(FuelBehaviourEvent::PeerIdentified(peer_id)); - } - - PeerInfoEvent::PeerInfoUpdated { peer_id } => self - .events - .push_back(FuelBehaviourEvent::PeerInfoUpdated(peer_id)), - } +impl From for FuelBehaviourEvent { + fn from(event: PeerInfoEvent) -> Self { + FuelBehaviourEvent::new(InnerBehaviourEvent::PeerInfo(event)) } } -impl NetworkBehaviourEventProcess - for FuelBehaviour -{ - fn inject_event(&mut self, message: GossipsubEvent) { - if let GossipsubEvent::Message { - propagation_source, - message, - message_id: _, - } = message - { - if let Some(correct_topic) = self - .network_metadata - .gossipsub_topics - .get_gossipsub_tag(&message.topic) - { - match self.codec.decode(&message.data, correct_topic) { - Ok(decoded_message) => { - self.events.push_back(FuelBehaviourEvent::GossipsubMessage { - peer_id: propagation_source, - topic_hash: message.topic, - message: decoded_message, - }) - } - Err(err) => { - warn!(target: "fuel-libp2p", "Failed to decode a message: {:?} with error: {:?}", &message.data, err); - } - } - } else { - warn!(target: "fuel-libp2p", "GossipTopicTag does not exist for {:?}", &message.topic); - } - } +impl From for FuelBehaviourEvent { + fn from(event: GossipsubEvent) -> Self { + FuelBehaviourEvent::new(InnerBehaviourEvent::Gossipsub(event)) } } -impl - NetworkBehaviourEventProcess< - RequestResponseEvent, - > for FuelBehaviour +impl From> + for FuelBehaviourEvent { - fn inject_event( - &mut self, - event: RequestResponseEvent, - ) { - match event { - RequestResponseEvent::Message { message, .. } => match message { - RequestResponseMessage::Request { - request, - channel, - request_id, - } => { - self.inbound_requests_table.insert(request_id, channel); - self.events.push_back(FuelBehaviourEvent::RequestMessage { - request_id, - request_message: request, - }) - } - RequestResponseMessage::Response { - request_id, - response, - } => { - match ( - self.outbound_requests_table.remove(&request_id), - self.codec.convert_to_response(&response), - ) { - ( - Some(ResponseChannelItem::ResponseBlock(channel)), - Ok(ResponseMessage::ResponseBlock(block)), - ) => { - if channel.send(block).is_err() { - debug!( - "Failed to send through the channel for {:?}", - request_id - ); - } - } - - (Some(_), Err(e)) => { - debug!("Failed to convert IntermediateResponse into a ResponseMessage {:?} with {:?}", response, e); - } - (None, Ok(_)) => { - debug!("Send channel not found for {:?}", request_id); - } - _ => {} - } - } - }, - RequestResponseEvent::InboundFailure { - peer, - error, - request_id, - } => { - debug!( - "RequestResponse inbound error for peer: {:?} with id: {:?} and error: {:?}", - peer, request_id, error - ); - } - RequestResponseEvent::OutboundFailure { - peer, - error, - request_id, - } => { - debug!( - "RequestResponse outbound error for peer: {:?} with id: {:?} and error: {:?}", - peer, request_id, error - ); - - let _ = self.outbound_requests_table.remove(&request_id); - } - _ => {} - } + fn from(event: RequestResponseEvent) -> Self { + FuelBehaviourEvent::new(InnerBehaviourEvent::RequestResponse(event)) } } diff --git a/fuel-p2p/src/orchestrator.rs b/fuel-p2p/src/orchestrator.rs index 81df8f5783a..92ae8afb0c1 100644 --- a/fuel-p2p/src/orchestrator.rs +++ b/fuel-p2p/src/orchestrator.rs @@ -23,7 +23,7 @@ use tokio::{ use tracing::warn; use crate::{ - behavior::FuelBehaviourEvent, + codecs::bincode::BincodeCodec, config::P2PConfig, gossipsub::messages::{ GossipsubBroadcastRequest, @@ -83,7 +83,11 @@ impl NetworkOrchestrator { } pub async fn run(mut self) -> anyhow::Result { - let mut p2p_service = FuelP2PService::new(self.p2p_config.clone()).await?; + let mut p2p_service = FuelP2PService::new( + self.p2p_config.clone(), + BincodeCodec::new(self.p2p_config.max_block_size), + ) + .await?; loop { tokio::select! { @@ -93,36 +97,34 @@ impl NetworkOrchestrator { } }, p2p_event = p2p_service.next_event() => { - if let FuelP2PEvent::Behaviour(behaviour_event) = p2p_event { - match behaviour_event { - FuelBehaviourEvent::GossipsubMessage { message, .. } => { - match message { - GossipsubMessage::NewTx(tx) => { - let _ = self.tx_transaction.send(TransactionBroadcast::NewTransaction(tx)); - }, - GossipsubMessage::NewBlock(block) => { - let _ = self.tx_block.send(BlockBroadcast::NewBlock(block)); - }, - GossipsubMessage::ConsensusVote(vote) => { - let _ = self.tx_consensus.send(ConsensusBroadcast::NewVote(vote)); - }, + match p2p_event { + FuelP2PEvent::GossipsubMessage { message, .. } => { + match message { + GossipsubMessage::NewTx(tx) => { + let _ = self.tx_transaction.send(TransactionBroadcast::NewTransaction(tx)); + }, + GossipsubMessage::NewBlock(block) => { + let _ = self.tx_block.send(BlockBroadcast::NewBlock(block)); + }, + GossipsubMessage::ConsensusVote(vote) => { + let _ = self.tx_consensus.send(ConsensusBroadcast::NewVote(vote)); + }, + } + }, + FuelP2PEvent::RequestMessage { request_message, request_id } => { + match request_message { + RequestMessage::RequestBlock(block_height) => { + let db = self.db.clone(); + let tx_outbound_response = self.tx_outbound_responses.clone(); + + tokio::spawn(async move { + let res = db.get_sealed_block(block_height).await.map(|block| (OutboundResponse::ResponseBlock(block), request_id)); + let _ = tx_outbound_response.send(res); + }); } - }, - FuelBehaviourEvent::RequestMessage { request_message, request_id } => { - match request_message { - RequestMessage::RequestBlock(block_height) => { - let db = self.db.clone(); - let tx_outbound_response = self.tx_outbound_responses.clone(); - - tokio::spawn(async move { - let res = db.get_sealed_block(block_height).await.map(|block| (OutboundResponse::ResponseBlock(block), request_id)); - let _ = tx_outbound_response.send(res); - }); - } - } - }, - _ => {} - } + } + }, + _ => {} } }, module_request_msg = self.rx_request_event.recv() => { diff --git a/fuel-p2p/src/service.rs b/fuel-p2p/src/service.rs index f21ea795490..1aa9c2cad03 100644 --- a/fuel-p2p/src/service.rs +++ b/fuel-p2p/src/service.rs @@ -2,31 +2,51 @@ use crate::{ behavior::{ FuelBehaviour, FuelBehaviourEvent, + InnerBehaviourEvent, }, - codecs::bincode::BincodeCodec, + codecs::NetworkCodec, config::{ build_transport, P2PConfig, }, - gossipsub::messages::GossipsubBroadcastRequest, - peer_info::PeerInfo, + discovery::DiscoveryEvent, + gossipsub::{ + messages::{ + GossipsubBroadcastRequest, + GossipsubMessage as FuelGossipsubMessage, + }, + topics::GossipsubTopics, + }, + peer_info::{ + PeerInfo, + PeerInfoEvent, + }, request_response::messages::{ + IntermediateResponse, OutboundResponse, RequestError, RequestMessage, ResponseChannelItem, ResponseError, + ResponseMessage, }, }; use futures::prelude::*; use libp2p::{ gossipsub::{ error::PublishError, + GossipsubEvent, MessageId, Topic, + TopicHash, }, multiaddr::Protocol, - request_response::RequestId, + request_response::{ + RequestId, + RequestResponseEvent, + RequestResponseMessage, + ResponseChannel, + }, swarm::SwarmEvent, Multiaddr, PeerId, @@ -34,31 +54,67 @@ use libp2p::{ }; use rand::Rng; use std::collections::HashMap; +use tracing::{ + debug, + warn, +}; /// Listens to the events on the p2p network /// And forwards them to the Orchestrator -pub struct FuelP2PService { +pub struct FuelP2PService { /// Store the local peer id pub local_peer_id: PeerId, /// Swarm handler for FuelBehaviour - swarm: Swarm>, + swarm: Swarm>, + + /// Holds the Sender(s) part of the Oneshot Channel from the NetworkOrchestrator + /// Once the ResponseMessage is received from the p2p Network + /// It will send it to the NetworkOrchestrator via its unique Sender + outbound_requests_table: HashMap, + + /// Holds the ResponseChannel(s) for the inbound requests from the p2p Network + /// Once the Response is prepared by the NetworkOrchestrator + /// It will send it to the specified Peer via its unique ResponseChannel + inbound_requests_table: HashMap>, + + /// NetworkCodec used as for encoding and decoding of Gossipsub messages + network_codec: Codec, + + /// Stores additional p2p network info + network_metadata: NetworkMetadata, } -#[allow(clippy::large_enum_variant)] +/// Holds additional Network data for FuelBehavior +#[derive(Debug)] +struct NetworkMetadata { + gossipsub_topics: GossipsubTopics, +} + +//#[allow(clippy::large_enum_variant)] #[derive(Debug, Clone)] +#[allow(clippy::large_enum_variant)] pub enum FuelP2PEvent { - Behaviour(FuelBehaviourEvent), - NewListenAddr(Multiaddr), + GossipsubMessage { + peer_id: PeerId, + topic_hash: TopicHash, + message: FuelGossipsubMessage, + }, + RequestMessage { + request_id: RequestId, + request_message: RequestMessage, + }, + PeerConnected(PeerId), + PeerDisconnected(PeerId), + PeerInfoUpdated(PeerId), } -impl FuelP2PService { - pub async fn new(config: P2PConfig) -> anyhow::Result { +impl FuelP2PService { + pub async fn new(config: P2PConfig, codec: Codec) -> anyhow::Result { let local_peer_id = PeerId::from(config.local_keypair.public()); // configure and build P2P Service let transport = build_transport(config.local_keypair.clone()).await; - let behaviour = - FuelBehaviour::new(&config, BincodeCodec::new(config.max_block_size)); + let behaviour = FuelBehaviour::new(&config, codec.clone()); let mut swarm = Swarm::new(transport, behaviour, local_peer_id); // set up node's address to listen on @@ -77,9 +133,16 @@ impl FuelP2PService { // start listening at the given address swarm.listen_on(listen_multiaddr)?; + let gossipsub_topics = GossipsubTopics::new(&config.network_name); + let network_metadata = NetworkMetadata { gossipsub_topics }; + Ok(Self { - swarm, local_peer_id, + swarm, + network_codec: codec, + outbound_requests_table: HashMap::default(), + inbound_requests_table: HashMap::default(), + network_metadata, }) } @@ -91,20 +154,17 @@ impl FuelP2PService { &mut self, message: GossipsubBroadcastRequest, ) -> Result { - self.swarm.behaviour_mut().publish_message(message) - } - - pub async fn next_event(&mut self) -> FuelP2PEvent { - loop { - match self.swarm.select_next_some().await { - SwarmEvent::Behaviour(fuel_behaviour) => { - return FuelP2PEvent::Behaviour(fuel_behaviour) - } - SwarmEvent::NewListenAddr { address, .. } => { - return FuelP2PEvent::NewListenAddr(address) - } - _ => {} - } + let topic = self + .network_metadata + .gossipsub_topics + .get_gossipsub_topic(&message); + + match self.network_codec.encode(message) { + Ok(encoded_data) => self + .swarm + .behaviour_mut() + .publish_message(topic, encoded_data), + Err(e) => Err(PublishError::TransformFailed(e)), } } @@ -128,11 +188,15 @@ impl FuelP2PService { } }; - Ok(self.swarm.behaviour_mut().send_request_msg( - message_request, - peer_id, - channel_item, - )) + let request_id = self + .swarm + .behaviour_mut() + .send_request_msg(message_request, peer_id); + + self.outbound_requests_table + .insert(request_id, channel_item); + + Ok(request_id) } /// Sends ResponseMessage to a peer that requested the data @@ -141,19 +205,177 @@ impl FuelP2PService { request_id: RequestId, message: OutboundResponse, ) -> Result<(), ResponseError> { - self.swarm - .behaviour_mut() - .send_response_msg(request_id, message) + match ( + self.network_codec.convert_to_intermediate(&message), + self.inbound_requests_table.remove(&request_id), + ) { + (Ok(message), Some(channel)) => { + if self + .swarm + .behaviour_mut() + .send_response_msg(channel, message) + .is_err() + { + debug!("Failed to send ResponseMessage for {:?}", request_id); + return Err(ResponseError::SendingResponseFailed) + } + } + (Ok(_), None) => { + debug!("ResponseChannel for {:?} does not exist!", request_id); + return Err(ResponseError::ResponseChannelDoesNotExist) + } + (Err(e), _) => { + debug!("Failed to convert to IntermediateResponse with {:?}", e); + return Err(ResponseError::ConversionToIntermediateFailed) + } + } + + Ok(()) + } + + pub async fn next_event(&mut self) -> FuelP2PEvent { + loop { + if let SwarmEvent::Behaviour(fuel_behaviour) = + self.swarm.select_next_some().await + { + if let Some(event) = self.handle_behaviour_event(fuel_behaviour) { + return event + } + } + } + } + + fn handle_behaviour_event( + &mut self, + event: FuelBehaviourEvent, + ) -> Option { + match event.into() { + InnerBehaviourEvent::Discovery(discovery_event) => match discovery_event { + DiscoveryEvent::Connected(peer_id, addresses) => { + self.swarm + .behaviour_mut() + .add_addresses_to_peer_info(&peer_id, addresses); + + return Some(FuelP2PEvent::PeerConnected(peer_id)) + } + DiscoveryEvent::Disconnected(peer_id) => { + return Some(FuelP2PEvent::PeerDisconnected(peer_id)) + } + _ => {} + }, + InnerBehaviourEvent::Gossipsub(gossipsub_event) => { + if let GossipsubEvent::Message { + propagation_source, + message, + .. + } = gossipsub_event + { + if let Some(correct_topic) = self + .network_metadata + .gossipsub_topics + .get_gossipsub_tag(&message.topic) + { + match self.network_codec.decode(&message.data, correct_topic) { + Ok(decoded_message) => { + return Some(FuelP2PEvent::GossipsubMessage { + peer_id: propagation_source, + topic_hash: message.topic, + message: decoded_message, + }) + } + Err(err) => { + warn!(target: "fuel-libp2p", "Failed to decode a message: {:?} with error: {:?}", &message.data, err); + } + } + } else { + warn!(target: "fuel-libp2p", "GossipTopicTag does not exist for {:?}", &message.topic); + } + } + } + + InnerBehaviourEvent::PeerInfo(peer_info_event) => match peer_info_event { + PeerInfoEvent::PeerIdentified { peer_id, addresses } => { + self.swarm + .behaviour_mut() + .add_addresses_to_discovery(&peer_id, addresses); + } + PeerInfoEvent::PeerInfoUpdated { peer_id } => { + return Some(FuelP2PEvent::PeerInfoUpdated(peer_id)) + } + }, + InnerBehaviourEvent::RequestResponse(req_res_event) => match req_res_event { + RequestResponseEvent::Message { message, .. } => match message { + RequestResponseMessage::Request { + request, + channel, + request_id, + } => { + self.inbound_requests_table.insert(request_id, channel); + + return Some(FuelP2PEvent::RequestMessage { + request_id, + request_message: request, + }) + } + RequestResponseMessage::Response { + request_id, + response, + } => { + match ( + self.outbound_requests_table.remove(&request_id), + self.network_codec.convert_to_response(&response), + ) { + ( + Some(ResponseChannelItem::ResponseBlock(channel)), + Ok(ResponseMessage::ResponseBlock(block)), + ) => { + if channel.send(block).is_err() { + debug!( + "Failed to send through the channel for {:?}", + request_id + ); + } + } + + (Some(_), Err(e)) => { + debug!("Failed to convert IntermediateResponse into a ResponseMessage {:?} with {:?}", response, e); + } + (None, Ok(_)) => { + debug!("Send channel not found for {:?}", request_id); + } + _ => {} + } + } + }, + RequestResponseEvent::InboundFailure { + peer, + error, + request_id, + } => { + debug!("RequestResponse inbound error for peer: {:?} with id: {:?} and error: {:?}", peer, request_id, error); + } + RequestResponseEvent::OutboundFailure { + peer, + error, + request_id, + } => { + debug!("RequestResponse outbound error for peer: {:?} with id: {:?} and error: {:?}", peer, request_id, error); + + let _ = self.outbound_requests_table.remove(&request_id); + } + _ => {} + }, + } + + None } } #[cfg(test)] mod tests { - use super::{ - FuelBehaviourEvent, - FuelP2PService, - }; + use super::FuelP2PService; use crate::{ + codecs::bincode::BincodeCodec, config::P2PConfig, gossipsub::{ messages::{ @@ -183,9 +405,11 @@ mod tests { FuelBlock, }, }; + use futures::StreamExt; use libp2p::{ gossipsub::Topic, identity::Keypair, + swarm::SwarmEvent, Multiaddr, PeerId, }; @@ -228,10 +452,16 @@ mod tests { } } - /// helper function for building FuelP2PService - async fn build_fuel_p2p_service(mut p2p_config: P2PConfig) -> FuelP2PService { + /// helper function for building FuelP2PService + async fn build_fuel_p2p_service( + mut p2p_config: P2PConfig, + ) -> FuelP2PService { p2p_config.local_keypair = Keypair::generate_secp256k1(); // change keypair for each Node - FuelP2PService::new(p2p_config).await.unwrap() + let max_block_size = p2p_config.max_block_size; + + FuelP2PService::new(p2p_config, BincodeCodec::new(max_block_size)) + .await + .unwrap() } /// attaches PeerId to the Multiaddr @@ -247,8 +477,8 @@ mod tests { .await; loop { - match fuel_p2p_service.next_event().await { - FuelP2PEvent::NewListenAddr(_address) => { + match fuel_p2p_service.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { .. } => { // listener address registered, we are good to go break } @@ -276,13 +506,13 @@ mod tests { loop { tokio::select! { node_b_event = node_b.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerConnected(_)) = node_b_event { + if let FuelP2PEvent::PeerConnected(_) = node_b_event { // successfully connected to Node B break } tracing::info!("Node B Event: {:?}", node_b_event); }, - node_a_event = node_a.next_event() => { + node_a_event = node_a.swarm.select_next_some() => { tracing::info!("Node A Event: {:?}", node_a_event); } }; @@ -299,8 +529,8 @@ mod tests { P2PConfig::default_with_network("nodes_connected_via_identify"); let mut node_a = build_fuel_p2p_service(p2p_config.clone()).await; - let node_a_address = match node_a.next_event().await { - FuelP2PEvent::NewListenAddr(address) => Some(address), + let node_a_address = match node_a.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { address, .. } => Some(address), _ => None, }; @@ -324,7 +554,7 @@ mod tests { }, node_c_event = node_c.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerConnected(peer_id)) = node_c_event { + if let FuelP2PEvent::PeerConnected(peer_id) = node_c_event { // we have connected to Node B! if peer_id == node_b.local_peer_id { break @@ -345,8 +575,8 @@ mod tests { let mut p2p_config = P2PConfig::default_with_network("peer_info_updates_work"); let mut node_a = build_fuel_p2p_service(p2p_config.clone()).await; - let node_a_address = match node_a.next_event().await { - FuelP2PEvent::NewListenAddr(address) => Some(address), + let node_a_address = match node_a.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { address, .. } => Some(address), _ => None, }; @@ -360,7 +590,7 @@ mod tests { loop { tokio::select! { node_a_event = node_a.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerInfoUpdated(peer_id)) = node_a_event { + if let FuelP2PEvent::PeerInfoUpdated(peer_id) = node_a_event { if let Some(PeerInfo { peer_addresses, latest_ping, client_version, .. }) = node_a.swarm.behaviour().get_peer_info(&peer_id) { // Exits after it verifies that: // 1. Peer Addresses are known @@ -434,8 +664,8 @@ mod tests { p2p_config.topics = topics.clone(); let mut node_a = build_fuel_p2p_service(p2p_config.clone()).await; - let node_a_address = match node_a.next_event().await { - FuelP2PEvent::NewListenAddr(address) => Some(address), + let node_a_address = match node_a.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { address, .. } => Some(address), _ => None, }; @@ -449,7 +679,7 @@ mod tests { loop { tokio::select! { node_a_event = node_a.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerInfoUpdated(peer_id)) = node_a_event { + if let FuelP2PEvent::PeerInfoUpdated(peer_id) = node_a_event { if let Some(PeerInfo { peer_addresses, .. }) = node_a.swarm.behaviour().get_peer_info(&peer_id) { // verifies that we've got at least a single peer address to send message to if !peer_addresses.is_empty() && !message_sent { @@ -463,7 +693,7 @@ mod tests { tracing::info!("Node A Event: {:?}", node_a_event); }, node_b_event = node_b.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::GossipsubMessage { topic_hash, message, .. }) = node_b_event.clone() { + if let FuelP2PEvent::GossipsubMessage { topic_hash, message, .. } = node_b_event.clone() { if topic_hash != selected_topic.hash() { tracing::error!("Wrong topic hash, expected: {} - actual: {}", selected_topic.hash(), topic_hash); panic!("Wrong Topic"); @@ -518,8 +748,8 @@ mod tests { // Node A let mut node_a = build_fuel_p2p_service(p2p_config.clone()).await; - let node_a_address = match node_a.next_event().await { - FuelP2PEvent::NewListenAddr(address) => Some(address), + let node_a_address = match node_a.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { address, .. } => Some(address), _ => None, }; @@ -542,7 +772,7 @@ mod tests { break; } node_a_event = node_a.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerInfoUpdated(peer_id)) = node_a_event { + if let FuelP2PEvent::PeerInfoUpdated(peer_id) = node_a_event { if let Some(PeerInfo { peer_addresses, .. }) = node_a.swarm.behaviour().get_peer_info(&peer_id) { // 0. verifies that we've got at least a single peer address to request message from if !peer_addresses.is_empty() && !request_sent { @@ -575,7 +805,7 @@ mod tests { }, node_b_event = node_b.next_event() => { // 2. Node B receives the RequestMessage from Node A initiated by the NetworkOrchestrator - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::RequestMessage{ request_id, .. }) = node_b_event { + if let FuelP2PEvent::RequestMessage{ request_id, .. } = node_b_event { let block = FuelBlock { header: FuelBlockHeader::default(), transactions: vec![Transaction::default(), Transaction::default(), Transaction::default(), Transaction::default(), Transaction::default()], @@ -609,8 +839,8 @@ mod tests { p2p_config.set_request_timeout = Duration::from_secs(0); let mut node_a = build_fuel_p2p_service(p2p_config.clone()).await; - let node_a_address = match node_a.next_event().await { - FuelP2PEvent::NewListenAddr(address) => Some(address), + let node_a_address = match node_a.swarm.select_next_some().await { + SwarmEvent::NewListenAddr { address, .. } => Some(address), _ => None, }; @@ -629,7 +859,7 @@ mod tests { loop { tokio::select! { node_a_event = node_a.next_event() => { - if let FuelP2PEvent::Behaviour(FuelBehaviourEvent::PeerInfoUpdated(peer_id)) = node_a_event { + if let FuelP2PEvent::PeerInfoUpdated(peer_id) = node_a_event { if let Some(PeerInfo { peer_addresses, .. }) = node_a.swarm.behaviour().get_peer_info(&peer_id) { // 0. verifies that we've got at least a single peer address to request message from if !peer_addresses.is_empty() && !request_sent { @@ -639,14 +869,14 @@ mod tests { let (tx_orchestrator, rx_orchestrator) = oneshot::channel(); // 2a. there should be ZERO pending outbound requests in the table - assert_eq!(node_a.swarm.behaviour().get_outbound_requests_table().len(), 0); + assert_eq!(node_a.outbound_requests_table.len(), 0); // Request successfully sent let requested_block_height = RequestMessage::RequestBlock(0_u64.into()); assert!(node_a.send_request_msg(None, requested_block_height, ResponseChannelItem::ResponseBlock(tx_orchestrator)).is_ok()); // 2b. there should be ONE pending outbound requests in the table - assert_eq!(node_a.swarm.behaviour().get_outbound_requests_table().len(), 1); + assert_eq!(node_a.outbound_requests_table.len(), 1); let tx_test_end = tx_test_end.clone(); @@ -666,7 +896,7 @@ mod tests { // we received a signal to end the test // 4. there should be ZERO pending outbound requests in the table // after the Outbound Request Failed with Timeout - assert_eq!(node_a.swarm.behaviour().get_outbound_requests_table().len(), 0); + assert_eq!(node_a.outbound_requests_table.len(), 0); break; }, // will not receive the request at all From 8f5d192a86f7809aee0c95ccef3636f2980fc0f5 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Fri, 30 Sep 2022 20:15:46 +0200 Subject: [PATCH 04/16] output disconnect event only if this is the last connection --- fuel-p2p/src/discovery.rs | 53 ++++++++++++++-------- fuel-p2p/src/discovery/discovery_config.rs | 5 +- 2 files changed, 37 insertions(+), 21 deletions(-) diff --git a/fuel-p2p/src/discovery.rs b/fuel-p2p/src/discovery.rs index a241844eb2e..78a5fd09da0 100644 --- a/fuel-p2p/src/discovery.rs +++ b/fuel-p2p/src/discovery.rs @@ -119,6 +119,17 @@ impl NetworkBehaviour for DiscoveryBehaviour { self.kademlia.inject_event(peer_id, connection, event); } + fn inject_address_change( + &mut self, + peer_id: &PeerId, + connection_id: &ConnectionId, + old: &ConnectedPoint, + new: &ConnectedPoint, + ) { + self.kademlia + .inject_address_change(peer_id, connection_id, old, new) + } + // gets polled by the swarm fn poll( &mut self, @@ -277,15 +288,8 @@ impl NetworkBehaviour for DiscoveryBehaviour { failed_addresses: Option<&Vec>, other_established: usize, ) { - if self.connected_peers.insert(*peer_id) { - self.kademlia.inject_connection_established( - peer_id, - connection_id, - endpoint, - failed_addresses, - other_established, - ); - + if other_established == 0 { + self.connected_peers.insert(*peer_id); let addresses = self.addresses_of_peer(peer_id); self.events @@ -293,6 +297,14 @@ impl NetworkBehaviour for DiscoveryBehaviour { trace!("Connected to a peer {:?}", peer_id); } + + self.kademlia.inject_connection_established( + peer_id, + connection_id, + endpoint, + failed_addresses, + other_established, + ); } fn inject_connection_closed( @@ -301,22 +313,23 @@ impl NetworkBehaviour for DiscoveryBehaviour { connection_id: &ConnectionId, connection_point: &ConnectedPoint, handler: ::Handler, - other_established: usize, + remaining_established: usize, ) { - if self.connected_peers.remove(peer_id) { - self.kademlia.inject_connection_closed( - peer_id, - connection_id, - connection_point, - handler, - other_established, - ); - + if remaining_established == 0 { + self.connected_peers.remove(peer_id); self.events .push_back(DiscoveryEvent::Disconnected(*peer_id)); trace!("Disconnected from {:?}", peer_id); } + + self.kademlia.inject_connection_closed( + peer_id, + connection_id, + connection_point, + handler, + remaining_established, + ); } fn inject_new_external_addr(&mut self, addr: &Multiaddr) { @@ -396,7 +409,7 @@ mod tests { .into_authentic(&keypair) .unwrap(); - let transport = core::transport::MemoryTransport::default() + let transport = core::transport::MemoryTransport::new() .upgrade(core::upgrade::Version::V1) .authenticate(noise::NoiseConfig::xx(noise_keys).into_authenticated()) .multiplex(yamux::YamuxConfig::default()) diff --git a/fuel-p2p/src/discovery/discovery_config.rs b/fuel-p2p/src/discovery/discovery_config.rs index 359de1ffacd..86d6daf4f84 100644 --- a/fuel-p2p/src/discovery/discovery_config.rs +++ b/fuel-p2p/src/discovery/discovery_config.rs @@ -102,8 +102,11 @@ impl DiscoveryConfig { let memory_store = MemoryStore::new(local_peer_id.to_owned()); let mut kademlia_config = KademliaConfig::default(); let network = format!("/fuel/kad/{}/kad/1.0.0", network_name); - kademlia_config.set_protocol_name(network.as_bytes().to_vec()); + let network_names = network.as_bytes().to_vec(); + kademlia_config + .set_protocol_names(std::iter::once(network_names.into()).collect()); kademlia_config.set_connection_idle_timeout(connection_idle_timeout); + let mut kademlia = Kademlia::with_config(local_peer_id, memory_store, kademlia_config); From 3b8625ebe53095af1d423b3d4d9ae47e07767358 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Fri, 30 Sep 2022 20:16:17 +0200 Subject: [PATCH 05/16] update libp2p --- Cargo.lock | 211 ++++++++++++++++++-------------------------- fuel-p2p/Cargo.toml | 4 +- 2 files changed, 87 insertions(+), 128 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c02703ea7eb..e9f234b43f9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -330,7 +330,7 @@ dependencies = [ "http-types", "httparse", "log", - "pin-project 1.0.12", + "pin-project", ] [[package]] @@ -497,15 +497,6 @@ dependencies = [ "pin-project-lite 0.2.9", ] -[[package]] -name = "atomic" -version = "0.5.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b88d82667eca772c4aa12f0f1348b3ae643424c8876448f3f7bd5787032e234c" -dependencies = [ - "autocfg", -] - [[package]] name = "atomic-waker" version = "1.0.0" @@ -1026,7 +1017,7 @@ version = "3.2.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ea0c8bce528c4be4da13ea6fead8965e95b6073585a2f05204bd8f4119f82a65" dependencies = [ - "heck 0.4.0", + "heck", "proc-macro-error", "proc-macro2", "quote", @@ -1753,7 +1744,7 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "21cdad81446a7f7dc43f6a77409efeb9733d2fa65553efef6018ef257c959b73" dependencies = [ - "heck 0.4.0", + "heck", "proc-macro2", "quote", "syn", @@ -1871,7 +1862,7 @@ dependencies = [ "futures-util", "hex", "once_cell", - "pin-project 1.0.12", + "pin-project", "serde", "serde_json", "thiserror", @@ -2009,7 +2000,7 @@ dependencies = [ "http", "once_cell", "parking_lot 0.11.2", - "pin-project 1.0.12", + "pin-project", "reqwest", "serde", "serde_json", @@ -2814,15 +2805,6 @@ dependencies = [ "fxhash", ] -[[package]] -name = "heck" -version = "0.3.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6d621efb26863f0e9924c6ac577e8275e5e6b77455db64ffa6c65c904e9e132c" -dependencies = [ - "unicode-segmentation", -] - [[package]] name = "heck" version = "0.4.0" @@ -3332,11 +3314,10 @@ dependencies = [ [[package]] name = "libp2p" -version = "0.44.0" +version = "0.48.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "475ce2ac4a9727e53a519f6ee05b38abfcba8f0d39c4d24f103d184e36fd5b0f" +checksum = "94c996fe5bfdba47f5a5af71d48ecbe8cec900b7b97391cc1d3ba1afb0e2d3b6" dependencies = [ - "atomic", "bytes", "futures", "futures-timer", @@ -3361,16 +3342,16 @@ dependencies = [ "libp2p-yamux", "multiaddr", "parking_lot 0.12.1", - "pin-project 1.0.12", + "pin-project", "rand 0.7.3", "smallvec", ] [[package]] name = "libp2p-core" -version = "0.32.1" +version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "db5b02602099fb75cb2d16f9ea860a320d6eb82ce41e95ab680912c454805cd5" +checksum = "b1fff5bd889c82a0aec668f2045edd066f559d4e5c40354e5a4c77ac00caac38" dependencies = [ "asn1_der", "bs58", @@ -3387,11 +3368,10 @@ dependencies = [ "multihash", "multistream-select", "parking_lot 0.12.1", - "pin-project 1.0.12", + "pin-project", "prost", "prost-build", "rand 0.8.5", - "ring", "rw-stream-sink", "sha2 0.10.6", "smallvec", @@ -3403,23 +3383,24 @@ dependencies = [ [[package]] name = "libp2p-dns" -version = "0.32.1" +version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "066e33e854e10b5c93fc650458bf2179c7e0d143db260b0963e44a94859817f1" +checksum = "6cb3c16e3bb2f76c751ae12f0f26e788c89d353babdded40411e7923f01fc978" dependencies = [ "async-std-resolver", "futures", "libp2p-core", "log", + "parking_lot 0.12.1", "smallvec", "trust-dns-resolver", ] [[package]] name = "libp2p-gossipsub" -version = "0.37.0" +version = "0.41.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a90c989a7c0969c2ab63e898da9bc735e3be53fb4f376e9c045ce516bcc9f928" +checksum = "2185aac44b162c95180ae4ddd1f4dfb705217ea1cb8e16bdfc70d31496fd80fa" dependencies = [ "asynchronous-codec", "base64 0.13.0", @@ -3445,10 +3426,11 @@ dependencies = [ [[package]] name = "libp2p-identify" -version = "0.35.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c5ef5a5b57904c7c33d6713ef918d239dc6b7553458f3475d87f8a18e9c651c8" +checksum = "f19440c84b509d69b13f0c9c28caa9bd3a059d25478527e937e86761f25c821e" dependencies = [ + "asynchronous-codec", "futures", "futures-timer", "libp2p-core", @@ -3457,16 +3439,19 @@ dependencies = [ "lru", "prost", "prost-build", + "prost-codec", "smallvec", + "thiserror", + "void", ] [[package]] name = "libp2p-kad" -version = "0.36.0" +version = "0.40.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "564e6bd64d177446399ed835b9451a8825b07929d6daa6a94e6405592974725e" +checksum = "f840579eed6503ec8a7a0c7e919bcd645df11c991be8e54640ff09f7109b8a43" dependencies = [ - "arrayvec 0.5.2", + "arrayvec 0.7.2", "asynchronous-codec", "bytes", "either", @@ -3490,9 +3475,9 @@ dependencies = [ [[package]] name = "libp2p-mdns" -version = "0.36.0" +version = "0.40.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "611ae873c8e280ccfab0d57c7a13cac5644f364529e233114ff07863946058b0" +checksum = "ff531fbceee32be0e39409e985c5536897e4578addb1c702fd4973e2c381fc76" dependencies = [ "async-io", "data-encoding", @@ -3511,9 +3496,9 @@ dependencies = [ [[package]] name = "libp2p-metrics" -version = "0.5.0" +version = "0.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "985be799bb3796e0c136c768208c3c06604a38430571906a13dcfeda225a3b9d" +checksum = "a74ab339e8b5d989e8c1000a78adb5c064a6319245bb22d1e70b415ec18c39b8" dependencies = [ "libp2p-core", "libp2p-gossipsub", @@ -3526,9 +3511,9 @@ dependencies = [ [[package]] name = "libp2p-mplex" -version = "0.32.0" +version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "442eb0c9fff0bf22a34f015724b4143ce01877e079ed0963c722d94c07c72160" +checksum = "ce53169351226ee0eb18ee7bef8d38f308fa8ad7244f986ae776390c0ae8a44d" dependencies = [ "asynchronous-codec", "bytes", @@ -3544,9 +3529,9 @@ dependencies = [ [[package]] name = "libp2p-noise" -version = "0.35.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9dd7e0c94051cda67123be68cf6b65211ba3dde7277be9068412de3e7ffd63ef" +checksum = "7cb0f939a444b06779ce551b3d78ebf13970ac27906ada452fd70abd160b09b8" dependencies = [ "bytes", "curve25519-dalek 3.2.0", @@ -3566,9 +3551,9 @@ dependencies = [ [[package]] name = "libp2p-ping" -version = "0.35.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bf57a3c2e821331dda9fe612d4654d676ab6e33d18d9434a18cced72630df6ad" +checksum = "76a36f78e107bb55330341018874c5168851f455f8bdc3e0cd44e6c84e0a7069" dependencies = [ "futures", "futures-timer", @@ -3582,9 +3567,9 @@ dependencies = [ [[package]] name = "libp2p-request-response" -version = "0.17.0" +version = "0.21.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b5e6a6fc6c9ad95661f46989473b34bd2993d14a4de497ff3b2668a910d4b869" +checksum = "2344aa93dc8b1de90e26091d7cf63044519bbea7b0b03d829f350563a2dd26f7" dependencies = [ "async-trait", "bytes", @@ -3600,9 +3585,9 @@ dependencies = [ [[package]] name = "libp2p-swarm" -version = "0.35.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8f0c69ad9e8f7c5fc50ad5ad9c7c8b57f33716532a2b623197f69f93e374d14c" +checksum = "70ad2db60c06603606b54b58e4247e32efec87a93cb4387be24bf32926c600f2" dependencies = [ "either", "fnv", @@ -3611,7 +3596,7 @@ dependencies = [ "instant", "libp2p-core", "log", - "pin-project 1.0.12", + "pin-project", "rand 0.7.3", "smallvec", "thiserror", @@ -3620,19 +3605,20 @@ dependencies = [ [[package]] name = "libp2p-swarm-derive" -version = "0.27.2" +version = "0.30.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4f693c8c68213034d472cbb93a379c63f4f307d97c06f1c41e4985de481687a5" +checksum = "1f02622b9dd150011b4eeec387f8bd013189a2f27da08ba363e7c6e606d77a48" dependencies = [ + "heck", "quote", "syn", ] [[package]] name = "libp2p-tcp" -version = "0.32.0" +version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "193447aa729c85aac2376828df76d171c1a589c9e6b58fcc7f9d9a020734122c" +checksum = "9675432b4c94b3960f3d2c7e57427b81aea92aab67fd0eebef09e2ae0ff54895" dependencies = [ "async-io", "futures", @@ -3647,15 +3633,16 @@ dependencies = [ [[package]] name = "libp2p-websocket" -version = "0.34.0" +version = "0.38.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c932834c3754501c368d1bf3d0fb458487a642b90fc25df082a3a2f3d3b32e37" +checksum = "de8a9e825cc03f2fc194d2e1622113d7fe18e1c7f4458a582b83140c9b9aea27" dependencies = [ "either", "futures", "futures-rustls", "libp2p-core", "log", + "parking_lot 0.12.1", "quicksink", "rw-stream-sink", "soketto", @@ -3665,9 +3652,9 @@ dependencies = [ [[package]] name = "libp2p-yamux" -version = "0.36.0" +version = "0.40.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "be902ebd89193cd020e89e89107726a38cfc0d16d18f613f4a37d046e92c7517" +checksum = "b74ec8dc042b583f0b2b93d52917f3b374c1e4b1cfa79ee74c7672c41257694c" dependencies = [ "futures", "libp2p-core", @@ -3957,7 +3944,7 @@ dependencies = [ "bytes", "futures", "log", - "pin-project 1.0.12", + "pin-project", "smallvec", "unsigned-varint", ] @@ -4134,15 +4121,6 @@ version = "6.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ff7415e9ae3fff1225851df9e0d9e4e5479f947619774677a63572e55e80eff" -[[package]] -name = "owning_ref" -version = "0.4.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6ff55baddef9e4ad00f88b6c743a2a8062d4c6ade126c2a528644b8e444d52ce" -dependencies = [ - "stable_deref_trait", -] - [[package]] name = "parity-scale-codec" version = "3.2.1" @@ -4362,33 +4340,13 @@ dependencies = [ "uncased", ] -[[package]] -name = "pin-project" -version = "0.4.30" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ef0f924a5ee7ea9cbcea77529dba45f8a9ba9f622419fe3386ca581a3ae9d5a" -dependencies = [ - "pin-project-internal 0.4.30", -] - [[package]] name = "pin-project" version = "1.0.12" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ad29a609b6bcd67fee905812e544992d216af9d755757c05ed2d0e15a74c6ecc" dependencies = [ - "pin-project-internal 1.0.12", -] - -[[package]] -name = "pin-project-internal" -version = "0.4.30" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "851c8d0ce9bebe43790dedfc86614c23494ac9f423dd618d3a61fc693eafe61e" -dependencies = [ - "proc-macro2", - "quote", - "syn", + "pin-project-internal", ] [[package]] @@ -4598,21 +4556,21 @@ dependencies = [ [[package]] name = "prometheus-client" -version = "0.15.1" +version = "0.18.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c9a896938cc6018c64f279888b8c7559d3725210d5db9a3a1ee6bc7188d51d34" +checksum = "3c473049631c233933d6286c88bbb7be30e62ec534cf99a9ae0079211f7fa603" dependencies = [ "dtoa", "itoa", - "owning_ref", + "parking_lot 0.12.1", "prometheus-client-derive-text-encode", ] [[package]] name = "prometheus-client-derive-text-encode" -version = "0.2.0" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8e12d01b9d66ad9eb4529c57666b6263fc1993cb30261d83ead658fdd932652" +checksum = "66a455fbcb954c1a7decf3c586e860fd7889cddf4b8e164be736dbac95a953cd" dependencies = [ "proc-macro2", "quote", @@ -4621,9 +4579,9 @@ dependencies = [ [[package]] name = "prost" -version = "0.9.0" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "444879275cb4fd84958b1a1d5420d15e6fcf7c235fe47f053c9c2a80aceb6001" +checksum = "399c3c31cdec40583bb68f0b18403400d01ec4289c383aa047560439952c4dd7" dependencies = [ "bytes", "prost-derive", @@ -4631,12 +4589,12 @@ dependencies = [ [[package]] name = "prost-build" -version = "0.9.0" +version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62941722fb675d463659e49c4f3fe1fe792ff24fe5bbaa9c08cd3b98a1c354f5" +checksum = "7f835c582e6bd972ba8347313300219fed5bfa52caf175298d860b61ff6069bb" dependencies = [ "bytes", - "heck 0.3.3", + "heck", "itertools", "lazy_static", "log", @@ -4649,11 +4607,24 @@ dependencies = [ "which", ] +[[package]] +name = "prost-codec" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "011ae9ff8359df7915f97302d591cdd9e0e27fbd5a4ddc5bd13b71079bb20987" +dependencies = [ + "asynchronous-codec", + "bytes", + "prost", + "thiserror", + "unsigned-varint", +] + [[package]] name = "prost-derive" -version = "0.9.0" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f9cc1a3263e07e0bf68e96268f37665207b49560d98739662cdfaae215c720fe" +checksum = "7345d5f0e08c0536d7ac7229952590239e77abf0a0100a1b1d890add6ea96364" dependencies = [ "anyhow", "itertools", @@ -4664,9 +4635,9 @@ dependencies = [ [[package]] name = "prost-types" -version = "0.9.0" +version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "534b7a0e836e3c482d2693070f982e39e7611da9695d4d1f5a4b186b51faef0a" +checksum = "4dfaa718ad76a44b3415e6c4d53b17c8f99160dcb3a99b10470fce8ad43f6e3e" dependencies = [ "bytes", "prost", @@ -5104,12 +5075,12 @@ checksum = "97477e48b4cf8603ad5f7aaf897467cf42ab4218a38ef76fb14c2d6773a6d6a8" [[package]] name = "rw-stream-sink" -version = "0.2.1" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4da5fcb054c46f5a5dff833b129285a93d3f0179531735e6c866e8cc307d2020" +checksum = "26338f5e09bb721b85b135ea05af7767c90b52f6de4f087d4f4a3a9d64e7dc04" dependencies = [ "futures", - "pin-project 0.4.30", + "pin-project", "static_assertions", ] @@ -5581,12 +5552,6 @@ dependencies = [ "der", ] -[[package]] -name = "stable_deref_trait" -version = "1.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a8f112729512f8e442d81f95a8a7ddf2b7c6b8a1a6f509a95864142b30cab2d3" - [[package]] name = "standback" version = "0.2.17" @@ -5672,7 +5637,7 @@ version = "0.24.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1e385be0d24f186b4ce2f9982191e7101bb737312ad61c1f2f984f34bcf85d59" dependencies = [ - "heck 0.4.0", + "heck", "proc-macro2", "quote", "rustversion", @@ -6032,7 +5997,7 @@ checksum = "b8fa9be0de6cf49e536ce1851f987bd21a43b771b09473c3549a6c853db37c1c" dependencies = [ "futures-core", "futures-util", - "pin-project 1.0.12", + "pin-project", "pin-project-lite 0.2.9", "tokio", "tokio-util", @@ -6144,7 +6109,7 @@ checksum = "97d095ae15e245a057c8e8451bab9b3ee1e1f68e9ba2b4fbc18d0ac5237835f2" dependencies = [ "futures", "futures-task", - "pin-project 1.0.12", + "pin-project", "tracing", ] @@ -6346,12 +6311,6 @@ dependencies = [ "tinyvec", ] -[[package]] -name = "unicode-segmentation" -version = "1.10.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0fdbf052a0783de01e944a6ce7a8cb939e295b1e7be835a1112c3b9a7f047a5a" - [[package]] name = "unicode-xid" version = "0.2.4" diff --git a/fuel-p2p/Cargo.toml b/fuel-p2p/Cargo.toml index 01232f9b461..32987006e42 100644 --- a/fuel-p2p/Cargo.toml +++ b/fuel-p2p/Cargo.toml @@ -18,8 +18,8 @@ fuel-core-interfaces = { path = "../fuel-core-interfaces", features = ["serde"], futures = "0.3" futures-timer = "3.0" ip_network = "0.4" -libp2p = { version = "0.44", default-features = false, features = [ - "dns-async-std", "gossipsub", "identify", "kad", "mdns", "mplex", "noise", +libp2p = { version = "0.48", default-features = false, features = [ + "dns-async-std", "gossipsub", "identify", "kad", "mdns-async-io", "mplex", "noise", "ping", "request-response", "secp256k1", "tcp-async-io", "yamux", "websocket" ] } rand = "0.8" From 049b1fc907b6caf8818c56c40d521edbdbf02e52 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Sat, 1 Oct 2022 11:52:32 +0200 Subject: [PATCH 06/16] set back the dial concucrrency factor to 1 --- fuel-p2p/src/discovery.rs | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/fuel-p2p/src/discovery.rs b/fuel-p2p/src/discovery.rs index 78a5fd09da0..f105da0c479 100644 --- a/fuel-p2p/src/discovery.rs +++ b/fuel-p2p/src/discovery.rs @@ -382,7 +382,10 @@ mod tests { identity::Keypair, multiaddr::Protocol, noise, - swarm::SwarmEvent, + swarm::{ + SwarmBuilder, + SwarmEvent, + }, yamux, Multiaddr, PeerId, @@ -394,6 +397,7 @@ mod tests { HashSet, VecDeque, }, + num::NonZeroU8, task::Poll, time::Duration, }; @@ -430,7 +434,11 @@ mod tests { }; let listen_addr: Multiaddr = Protocol::Memory(rand::random::()).into(); - let mut swarm = Swarm::new(transport, behaviour, keypair.public().to_peer_id()); + let swarm_builder = + SwarmBuilder::new(transport, behaviour, keypair.public().to_peer_id()) + .dial_concurrency_factor(NonZeroU8::new(1).expect("1 > 0")); + + let mut swarm = swarm_builder.build(); swarm .listen_on(listen_addr.clone()) From 378ec7577480576e9ad0a4ecd6e9f0ce884d54a0 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Sat, 1 Oct 2022 12:42:39 +0200 Subject: [PATCH 07/16] bump ci From 78cd25e79004f74daa041befb86ad1e192942e48 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Sat, 1 Oct 2022 13:06:05 +0200 Subject: [PATCH 08/16] setup protoc in CI --- .github/workflows/ci.yml | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3f01ec41871..d10ed32b9cd 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -109,7 +109,9 @@ jobs: - name: Install Cargo Make uses: davidB/rust-cargo-make@v1 with: - version: "0.36.0" + version: "0.36.0" + - name: Install Protoc + uses: arduino/setup-protoc@v1 - uses: Swatinem/rust-cache@v2 with: key: '${{ matrix.command }} ${{ matrix.args }}' From d33c1dd872458d29ea958dbe8d7878aaa0ef87a8 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 3 Oct 2022 11:26:52 +0200 Subject: [PATCH 09/16] rename events --- fuel-p2p/src/behavior.rs | 32 ++++++++++++++++---------------- fuel-p2p/src/service.rs | 12 ++++++------ 2 files changed, 22 insertions(+), 22 deletions(-) diff --git a/fuel-p2p/src/behavior.rs b/fuel-p2p/src/behavior.rs index 10c826e26f6..1b991fa0b7e 100644 --- a/fuel-p2p/src/behavior.rs +++ b/fuel-p2p/src/behavior.rs @@ -48,21 +48,21 @@ use std::{ }; #[derive(Debug)] -pub struct FuelBehaviourEvent { - event: InnerBehaviourEvent, +pub struct BehaviourEventWrapper { + event: FuelBehaviourEvent, codec: PhantomData, } #[derive(Debug)] -pub enum InnerBehaviourEvent { +pub enum FuelBehaviourEvent { Discovery(DiscoveryEvent), PeerInfo(PeerInfoEvent), Gossipsub(GossipsubEvent), RequestResponse(RequestResponseEvent), } -impl FuelBehaviourEvent { - fn new(event: InnerBehaviourEvent) -> Self { +impl BehaviourEventWrapper { + fn new(event: FuelBehaviourEvent) -> Self { Self { event, codec: PhantomData, @@ -70,15 +70,15 @@ impl FuelBehaviourEvent { } } -impl From> for InnerBehaviourEvent { - fn from(fuel_event: FuelBehaviourEvent) -> Self { +impl From> for FuelBehaviourEvent { + fn from(fuel_event: BehaviourEventWrapper) -> Self { fuel_event.event } } /// Handles all p2p protocols needed for Fuel. #[derive(NetworkBehaviour)] -#[behaviour(out_event = "FuelBehaviourEvent")] +#[behaviour(out_event = "BehaviourEventWrapper")] pub struct FuelBehaviour { /// Node discovery discovery: DiscoveryBehaviour, @@ -197,28 +197,28 @@ impl FuelBehaviour { } } -impl From for FuelBehaviourEvent { +impl From for BehaviourEventWrapper { fn from(event: DiscoveryEvent) -> Self { - FuelBehaviourEvent::new(InnerBehaviourEvent::Discovery(event)) + BehaviourEventWrapper::new(FuelBehaviourEvent::Discovery(event)) } } -impl From for FuelBehaviourEvent { +impl From for BehaviourEventWrapper { fn from(event: PeerInfoEvent) -> Self { - FuelBehaviourEvent::new(InnerBehaviourEvent::PeerInfo(event)) + BehaviourEventWrapper::new(FuelBehaviourEvent::PeerInfo(event)) } } -impl From for FuelBehaviourEvent { +impl From for BehaviourEventWrapper { fn from(event: GossipsubEvent) -> Self { - FuelBehaviourEvent::new(InnerBehaviourEvent::Gossipsub(event)) + BehaviourEventWrapper::new(FuelBehaviourEvent::Gossipsub(event)) } } impl From> - for FuelBehaviourEvent + for BehaviourEventWrapper { fn from(event: RequestResponseEvent) -> Self { - FuelBehaviourEvent::new(InnerBehaviourEvent::RequestResponse(event)) + BehaviourEventWrapper::new(FuelBehaviourEvent::RequestResponse(event)) } } diff --git a/fuel-p2p/src/service.rs b/fuel-p2p/src/service.rs index 1aa9c2cad03..5397bbb8579 100644 --- a/fuel-p2p/src/service.rs +++ b/fuel-p2p/src/service.rs @@ -1,8 +1,8 @@ use crate::{ behavior::{ + BehaviourEventWrapper, FuelBehaviour, FuelBehaviourEvent, - InnerBehaviourEvent, }, codecs::NetworkCodec, config::{ @@ -247,10 +247,10 @@ impl FuelP2PService { fn handle_behaviour_event( &mut self, - event: FuelBehaviourEvent, + event: BehaviourEventWrapper, ) -> Option { match event.into() { - InnerBehaviourEvent::Discovery(discovery_event) => match discovery_event { + FuelBehaviourEvent::Discovery(discovery_event) => match discovery_event { DiscoveryEvent::Connected(peer_id, addresses) => { self.swarm .behaviour_mut() @@ -263,7 +263,7 @@ impl FuelP2PService { } _ => {} }, - InnerBehaviourEvent::Gossipsub(gossipsub_event) => { + FuelBehaviourEvent::Gossipsub(gossipsub_event) => { if let GossipsubEvent::Message { propagation_source, message, @@ -293,7 +293,7 @@ impl FuelP2PService { } } - InnerBehaviourEvent::PeerInfo(peer_info_event) => match peer_info_event { + FuelBehaviourEvent::PeerInfo(peer_info_event) => match peer_info_event { PeerInfoEvent::PeerIdentified { peer_id, addresses } => { self.swarm .behaviour_mut() @@ -303,7 +303,7 @@ impl FuelP2PService { return Some(FuelP2PEvent::PeerInfoUpdated(peer_id)) } }, - InnerBehaviourEvent::RequestResponse(req_res_event) => match req_res_event { + FuelBehaviourEvent::RequestResponse(req_res_event) => match req_res_event { RequestResponseEvent::Message { message, .. } => match message { RequestResponseMessage::Request { request, From 04904741d2255761304016c5a3828f2d2560dff4 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 3 Oct 2022 11:41:11 +0200 Subject: [PATCH 10/16] remove commented out code --- fuel-p2p/src/service.rs | 1 - 1 file changed, 1 deletion(-) diff --git a/fuel-p2p/src/service.rs b/fuel-p2p/src/service.rs index 5397bbb8579..aabc211169d 100644 --- a/fuel-p2p/src/service.rs +++ b/fuel-p2p/src/service.rs @@ -90,7 +90,6 @@ struct NetworkMetadata { gossipsub_topics: GossipsubTopics, } -//#[allow(clippy::large_enum_variant)] #[derive(Debug, Clone)] #[allow(clippy::large_enum_variant)] pub enum FuelP2PEvent { From 9c0368c83d8cb0e5c7bbeda77b8bc621cfe1eec4 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 3 Oct 2022 11:45:23 +0200 Subject: [PATCH 11/16] add github token for protoc --- .github/workflows/ci.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d10ed32b9cd..4691c900b0a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -112,6 +112,8 @@ jobs: version: "0.36.0" - name: Install Protoc uses: arduino/setup-protoc@v1 + with: + token: ${{ secrets.GITHUB_TOKEN }} - uses: Swatinem/rust-cache@v2 with: key: '${{ matrix.command }} ${{ matrix.args }}' From 4e7923cc8e51173701a8942e625594a8939f6826 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Mon, 3 Oct 2022 11:49:48 +0200 Subject: [PATCH 12/16] cleanup network name --- fuel-p2p/src/discovery/discovery_config.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/fuel-p2p/src/discovery/discovery_config.rs b/fuel-p2p/src/discovery/discovery_config.rs index 86d6daf4f84..5bfea1adfac 100644 --- a/fuel-p2p/src/discovery/discovery_config.rs +++ b/fuel-p2p/src/discovery/discovery_config.rs @@ -102,9 +102,8 @@ impl DiscoveryConfig { let memory_store = MemoryStore::new(local_peer_id.to_owned()); let mut kademlia_config = KademliaConfig::default(); let network = format!("/fuel/kad/{}/kad/1.0.0", network_name); - let network_names = network.as_bytes().to_vec(); - kademlia_config - .set_protocol_names(std::iter::once(network_names.into()).collect()); + let network_name = network.as_bytes().to_vec(); + kademlia_config.set_protocol_names(vec![network_name.into()]); kademlia_config.set_connection_idle_timeout(connection_idle_timeout); let mut kademlia = From 130a7613d81fd4dcee3e81f6ede983a91e5e265f Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Tue, 4 Oct 2022 15:22:12 +0200 Subject: [PATCH 13/16] switch to tokio mdns --- Cargo.lock | 2 +- fuel-p2p/Cargo.toml | 2 +- fuel-p2p/src/discovery/mdns.rs | 13 ++++++++----- 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9c5e8d981d0..15b38baa866 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3494,7 +3494,6 @@ version = "0.40.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff531fbceee32be0e39409e985c5536897e4578addb1c702fd4973e2c381fc76" dependencies = [ - "async-io", "data-encoding", "dns-parser", "futures", @@ -3506,6 +3505,7 @@ dependencies = [ "rand 0.8.5", "smallvec", "socket2", + "tokio", "void", ] diff --git a/fuel-p2p/Cargo.toml b/fuel-p2p/Cargo.toml index 4a43f3d2f69..4389bf284b4 100644 --- a/fuel-p2p/Cargo.toml +++ b/fuel-p2p/Cargo.toml @@ -19,7 +19,7 @@ futures = "0.3" futures-timer = "3.0" ip_network = "0.4" libp2p = { version = "0.48", default-features = false, features = [ - "dns-async-std", "gossipsub", "identify", "kad", "mdns-async-io", "mplex", "noise", + "dns-async-std", "gossipsub", "identify", "kad", "mdns-tokio", "mplex", "noise", "ping", "request-response", "secp256k1", "tcp-async-io", "yamux", "websocket" ] } rand = "0.8" diff --git a/fuel-p2p/src/discovery/mdns.rs b/fuel-p2p/src/discovery/mdns.rs index 9f7405c8783..bd74d87e3a5 100644 --- a/fuel-p2p/src/discovery/mdns.rs +++ b/fuel-p2p/src/discovery/mdns.rs @@ -4,9 +4,9 @@ use futures::{ }; use libp2p::{ mdns::{ - Mdns, MdnsConfig, MdnsEvent, + TokioMdns, }, swarm::{ NetworkBehaviour, @@ -25,14 +25,14 @@ use tracing::warn; #[allow(clippy::large_enum_variant)] // Wrapper around mDNS so that `DiscoveryConfig::finish` does not have to be an `async` function pub enum MdnsWrapper { - Instantiating(BoxFuture<'static, std::io::Result>), - Ready(Mdns), + Instantiating(BoxFuture<'static, std::io::Result>), + Ready(TokioMdns), Disabled, } impl Default for MdnsWrapper { fn default() -> Self { - MdnsWrapper::Instantiating(Mdns::new(MdnsConfig::default()).boxed()) + MdnsWrapper::Instantiating(TokioMdns::new(MdnsConfig::default()).boxed()) } } @@ -53,7 +53,10 @@ impl MdnsWrapper { cx: &mut Context<'_>, params: &mut impl PollParameters, ) -> Poll< - NetworkBehaviourAction::ConnectionHandler>, + NetworkBehaviourAction< + MdnsEvent, + ::ConnectionHandler, + >, > { loop { match self { From 78d4301602b2b2a9d6cf2b53d5943177728026b0 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Tue, 4 Oct 2022 16:06:53 +0200 Subject: [PATCH 14/16] move to tokio only features --- Cargo.lock | 51 ++++------------------------------------- fuel-p2p/Cargo.toml | 4 ++-- fuel-p2p/src/config.rs | 12 +++++----- fuel-p2p/src/service.rs | 11 +++++++-- 4 files changed, 21 insertions(+), 57 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 15b38baa866..6977625ccf4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -363,24 +363,6 @@ dependencies = [ "event-listener", ] -[[package]] -name = "async-process" -version = "1.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "02111fd8655a613c25069ea89fc8d9bb89331fa77486eb3bc059ee757cfa481c" -dependencies = [ - "async-io", - "autocfg", - "blocking", - "cfg-if", - "event-listener", - "futures-lite", - "libc", - "once_cell", - "signal-hook", - "winapi", -] - [[package]] name = "async-std" version = "1.12.0" @@ -391,7 +373,6 @@ dependencies = [ "async-global-executor", "async-io", "async-lock", - "async-process", "crossbeam-utils", "futures-channel", "futures-core", @@ -408,21 +389,6 @@ dependencies = [ "wasm-bindgen-futures", ] -[[package]] -name = "async-std-resolver" -version = "0.21.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0f2f8a4a203be3325981310ab243a28e6e4ea55b6519bffce05d41ab60e09ad8" -dependencies = [ - "async-std", - "async-trait", - "futures-io", - "futures-util", - "pin-utils", - "socket2", - "trust-dns-resolver", -] - [[package]] name = "async-stream" version = "0.3.3" @@ -3402,7 +3368,6 @@ version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6cb3c16e3bb2f76c751ae12f0f26e788c89d353babdded40411e7923f01fc978" dependencies = [ - "async-std-resolver", "futures", "libp2p-core", "log", @@ -3635,15 +3600,15 @@ version = "0.36.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9675432b4c94b3960f3d2c7e57427b81aea92aab67fd0eebef09e2ae0ff54895" dependencies = [ - "async-io", "futures", "futures-timer", - "if-watch", + "if-addrs", "ipnet", "libc", "libp2p-core", "log", "socket2", + "tokio", ] [[package]] @@ -5437,16 +5402,6 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "43b2853a4d09f215c24cc5489c992ce46052d359b5109343cbafbf26bc62f8a3" -[[package]] -name = "signal-hook" -version = "0.3.14" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a253b5e89e2698464fc26b545c9edceb338e18a89effeeecfea192c3025be29d" -dependencies = [ - "libc", - "signal-hook-registry", -] - [[package]] name = "signal-hook-registry" version = "1.4.0" @@ -6213,6 +6168,7 @@ dependencies = [ "smallvec", "thiserror", "tinyvec", + "tokio", "url", ] @@ -6232,6 +6188,7 @@ dependencies = [ "resolv-conf", "smallvec", "thiserror", + "tokio", "trust-dns-proto", ] diff --git a/fuel-p2p/Cargo.toml b/fuel-p2p/Cargo.toml index 4389bf284b4..27d822cbc3d 100644 --- a/fuel-p2p/Cargo.toml +++ b/fuel-p2p/Cargo.toml @@ -19,8 +19,8 @@ futures = "0.3" futures-timer = "3.0" ip_network = "0.4" libp2p = { version = "0.48", default-features = false, features = [ - "dns-async-std", "gossipsub", "identify", "kad", "mdns-tokio", "mplex", "noise", - "ping", "request-response", "secp256k1", "tcp-async-io", "yamux", "websocket" + "dns-tokio", "gossipsub", "identify", "kad", "mdns-tokio", "mplex", "noise", + "ping", "request-response", "secp256k1", "tcp-tokio", "yamux", "websocket" ] } rand = "0.8" serde = { version = "1.0", features = ["derive"] } diff --git a/fuel-p2p/src/config.rs b/fuel-p2p/src/config.rs index 9638c99511a..ce707104f89 100644 --- a/fuel-p2p/src/config.rs +++ b/fuel-p2p/src/config.rs @@ -11,7 +11,7 @@ use libp2p::{ noise, tcp::{ GenTcpConfig, - TcpTransport, + TokioTcpTransport, }, yamux, Multiaddr, @@ -126,15 +126,15 @@ pub(crate) async fn build_transport( local_keypair: Keypair, ) -> Boxed<(PeerId, StreamMuxerBox)> { let transport = { - let generate_tcp_transpot = - || TcpTransport::new(GenTcpConfig::new().port_reuse(true).nodelay(true)); + let generate_tcp_transport = + || TokioTcpTransport::new(GenTcpConfig::new().port_reuse(true).nodelay(true)); - let tcp = generate_tcp_transpot(); + let tcp = generate_tcp_transport(); let ws_tcp = - libp2p::websocket::WsConfig::new(generate_tcp_transpot()).or_transport(tcp); + libp2p::websocket::WsConfig::new(generate_tcp_transport()).or_transport(tcp); - libp2p::dns::DnsConfig::system(ws_tcp).await.unwrap() + libp2p::dns::TokioDnsConfig::system(ws_tcp).unwrap() }; let auth_config = { diff --git a/fuel-p2p/src/service.rs b/fuel-p2p/src/service.rs index aabc211169d..eeaf7c3ea96 100644 --- a/fuel-p2p/src/service.rs +++ b/fuel-p2p/src/service.rs @@ -47,7 +47,10 @@ use libp2p::{ RequestResponseMessage, ResponseChannel, }, - swarm::SwarmEvent, + swarm::{ + SwarmBuilder, + SwarmEvent, + }, Multiaddr, PeerId, Swarm, @@ -114,7 +117,11 @@ impl FuelP2PService { // configure and build P2P Service let transport = build_transport(config.local_keypair.clone()).await; let behaviour = FuelBehaviour::new(&config, codec.clone()); - let mut swarm = Swarm::new(transport, behaviour, local_peer_id); + let mut swarm = SwarmBuilder::new(transport, behaviour, local_peer_id) + .executor(Box::new(|fut| { + tokio::spawn(fut); + })) + .build(); // set up node's address to listen on let listen_multiaddr = { From 681a5ebfd5027b54b4fc9dd37414d822e5aa8367 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Tue, 4 Oct 2022 16:08:25 +0200 Subject: [PATCH 15/16] remove async from method --- fuel-p2p/src/config.rs | 4 +--- fuel-p2p/src/orchestrator.rs | 3 +-- fuel-p2p/src/service.rs | 8 +++----- 3 files changed, 5 insertions(+), 10 deletions(-) diff --git a/fuel-p2p/src/config.rs b/fuel-p2p/src/config.rs index ce707104f89..f67e6670e25 100644 --- a/fuel-p2p/src/config.rs +++ b/fuel-p2p/src/config.rs @@ -122,9 +122,7 @@ impl P2PConfig { /// TCP/IP, Websocket /// Noise as encryption layer /// mplex or yamux for multiplexing -pub(crate) async fn build_transport( - local_keypair: Keypair, -) -> Boxed<(PeerId, StreamMuxerBox)> { +pub(crate) fn build_transport(local_keypair: Keypair) -> Boxed<(PeerId, StreamMuxerBox)> { let transport = { let generate_tcp_transport = || TokioTcpTransport::new(GenTcpConfig::new().port_reuse(true).nodelay(true)); diff --git a/fuel-p2p/src/orchestrator.rs b/fuel-p2p/src/orchestrator.rs index 92ae8afb0c1..e77f0521777 100644 --- a/fuel-p2p/src/orchestrator.rs +++ b/fuel-p2p/src/orchestrator.rs @@ -86,8 +86,7 @@ impl NetworkOrchestrator { let mut p2p_service = FuelP2PService::new( self.p2p_config.clone(), BincodeCodec::new(self.p2p_config.max_block_size), - ) - .await?; + )?; loop { tokio::select! { diff --git a/fuel-p2p/src/service.rs b/fuel-p2p/src/service.rs index eeaf7c3ea96..95c8bb20fc1 100644 --- a/fuel-p2p/src/service.rs +++ b/fuel-p2p/src/service.rs @@ -111,11 +111,11 @@ pub enum FuelP2PEvent { } impl FuelP2PService { - pub async fn new(config: P2PConfig, codec: Codec) -> anyhow::Result { + pub fn new(config: P2PConfig, codec: Codec) -> anyhow::Result { let local_peer_id = PeerId::from(config.local_keypair.public()); // configure and build P2P Service - let transport = build_transport(config.local_keypair.clone()).await; + let transport = build_transport(config.local_keypair.clone()); let behaviour = FuelBehaviour::new(&config, codec.clone()); let mut swarm = SwarmBuilder::new(transport, behaviour, local_peer_id) .executor(Box::new(|fut| { @@ -465,9 +465,7 @@ mod tests { p2p_config.local_keypair = Keypair::generate_secp256k1(); // change keypair for each Node let max_block_size = p2p_config.max_block_size; - FuelP2PService::new(p2p_config, BincodeCodec::new(max_block_size)) - .await - .unwrap() + FuelP2PService::new(p2p_config, BincodeCodec::new(max_block_size)).unwrap() } /// attaches PeerId to the Multiaddr From 7001101549c30b6e023caf56f487944881ab1fc9 Mon Sep 17 00:00:00 2001 From: leviathanbeak88 Date: Wed, 5 Oct 2022 14:33:00 +0200 Subject: [PATCH 16/16] document protoc as system requirement --- CONTRIBUTING.md | 51 +++++++++++++++++++++++++------------------------ README.md | 11 +++++------ 2 files changed, 31 insertions(+), 31 deletions(-) diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 7b54411dd92..3312859d312 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -18,9 +18,10 @@ But for `rustfmt`, we use Rust nightly toolchain because it provides more code s To build Fuel Core you'll need to at least have the following installed: -* `git` - version control -* [`rustup`](https://rustup.rs/) - Rust installer and toolchain manager -* [`clang`](http://releases.llvm.org/download.html) - Used to build system libraries (required for rocksdb). +- `git` - version control +- [`rustup`](https://rustup.rs/) - Rust installer and toolchain manager +- [`clang`](http://releases.llvm.org/download.html) - Used to build system libraries (required for rocksdb). +- [`protoc`](https://grpc.io/docs/protoc-installation/) - Used to compile Protocol Buffer files (required by libp2p). See the [README.md](README.md#system-requirements) for platform specific setup steps. @@ -66,7 +67,7 @@ cargo +nightly fmt --all --check cargo clippy --all-targets ``` -The test suite follows the Rust cargo standards. The GraphQL service will be instantiated by +The test suite follows the Rust cargo standards. The GraphQL service will be instantiated by Tower and will emulate a server/client structure. Testing is simply done using Cargo: @@ -98,22 +99,22 @@ cargo build -p fuel-core --no-default-features This is a rough outline of what a contributor's workflow looks like: -* Make sure what you want to contribute is already traced as an issue. - * We may discuss the problem and solution in the issue. -* Create a Git branch from where you want to base your work. This is usually master. -* Write code, add test cases, and commit your work. -* Run tests and make sure all tests pass. -* If the PR contains any breaking changes, add the breaking label to your PR. -* If you are part of the FuelLabs Github org, please open a PR from the repository itself. -* Otherwise, push your changes to a branch in your fork of the repository and submit a pull request. - * Make sure mention the issue, which is created at step 1, in the commit message. -* Your PR will be reviewed and some changes may be requested. - * Once you've made changes, your PR must be re-reviewed and approved. - * If the PR becomes out of date, you can use GitHub's 'update branch' button. - * If there are conflicts, you can merge and resolve them locally. Then push to your PR branch. - Any changes to the branch will require a re-review. -* Our CI system (Github Actions) automatically tests all authorized pull requests. -* Use Github to merge the PR once approved. +- Make sure what you want to contribute is already traced as an issue. + - We may discuss the problem and solution in the issue. +- Create a Git branch from where you want to base your work. This is usually master. +- Write code, add test cases, and commit your work. +- Run tests and make sure all tests pass. +- If the PR contains any breaking changes, add the breaking label to your PR. +- If you are part of the FuelLabs Github org, please open a PR from the repository itself. +- Otherwise, push your changes to a branch in your fork of the repository and submit a pull request. + - Make sure mention the issue, which is created at step 1, in the commit message. +- Your PR will be reviewed and some changes may be requested. + - Once you've made changes, your PR must be re-reviewed and approved. + - If the PR becomes out of date, you can use GitHub's 'update branch' button. + - If there are conflicts, you can merge and resolve them locally. Then push to your PR branch. + Any changes to the branch will require a re-review. +- Our CI system (Github Actions) automatically tests all authorized pull requests. +- Use Github to merge the PR once approved. Thanks for your contributions! @@ -125,11 +126,11 @@ If you are planning something big, for example, relates to multiple components o The Client team actively develops and maintains several dependencies used in Fuel Core, which you may be also interested in: -* [fuel-types](https://github.com/FuelLabs/fuel-types) -* [fuel-merkle](https://github.com/FuelLabs/fuel-merkle) -* [fuel-tx](https://github.com/FuelLabs/fuel-tx) -* [fuel-asm](https://github.com/FuelLabs/fuel-asm) -* [fuel-vm](https://github.com/FuelLabs/fuel-vm) +- [fuel-types](https://github.com/FuelLabs/fuel-types) +- [fuel-merkle](https://github.com/FuelLabs/fuel-merkle) +- [fuel-tx](https://github.com/FuelLabs/fuel-tx) +- [fuel-asm](https://github.com/FuelLabs/fuel-asm) +- [fuel-vm](https://github.com/FuelLabs/fuel-vm) ### Linking issues diff --git a/README.md b/README.md index 2760166a742..2193e30226a 100644 --- a/README.md +++ b/README.md @@ -22,19 +22,20 @@ There are several system requirements including clang. ```bash brew update brew install cmake +brew install protobuf ``` ###### Debian ```bash apt update -apt install -y cmake pkg-config build-essential git clang libclang-dev +apt install -y cmake pkg-config build-essential git clang libclang-dev protobuf-compiler ``` ###### Arch ```bash -pacman -Syu --needed --noconfirm cmake gcc pkgconf git clang +pacman -Syu --needed --noconfirm cmake gcc pkgconf git clang protobuf-compiler ``` ## Building @@ -64,7 +65,6 @@ OPTIONS: ... ``` - For many development puposes it is useful to have a state that won't persist and the `db-type` option can be set to `in-memory` as in the following example. #### Example @@ -94,7 +94,6 @@ On some macOS versions the default file descriptor limit is quite low, which can ulimit -n 10240 ``` - #### Log level The service relies on the environment variable `RUST_LOG`. For more information, check the [EnvFilter examples](https://docs.rs/tracing-subscriber/latest/tracing_subscriber/struct.EnvFilter.html#examples) crate. @@ -125,8 +124,8 @@ The client functionality is available through a service endpoint that expect Gra The transaction executor currently performs instant block production. Changes are persisted to RocksDB by default. -* Service endpoint: `/graphql` -* Schema (available after building): `fuel-client/assets/schema.sdl` +- Service endpoint: `/graphql` +- Schema (available after building): `fuel-client/assets/schema.sdl` The service expects a mutation defined as `submit` that receives a [Transaction](https://github.com/FuelLabs/fuel-tx) in hex encoded binary format, as [specified here](https://github.com/FuelLabs/fuel-specs/blob/master/specs/protocol/tx_format.md).