1use std::time::{Duration, UNIX_EPOCH};
5
6use crate::prelude::*;
7use crate::{blocks::GossipBlock, rpc::net::NetInfoResult};
8use crate::{chain::ChainStore, utils::encoding::from_slice_with_fallback};
9use crate::{
10 libp2p_bitswap::{BitswapStoreReadWrite, request_manager::BitswapRequestManager},
11 utils::flume::FlumeSenderExt as _,
12};
13use crate::{message::SignedMessage, networks::GenesisNetworkName};
14use ahash::{HashMap, HashSet};
15use anyhow::Context as _;
16use flume::Sender;
17use futures::{select, stream::StreamExt as _};
18pub use libp2p::gossipsub::{IdentTopic, Topic};
19use libp2p::{
20 PeerId, Swarm, SwarmBuilder,
21 autonat::NatStatus,
22 connection_limits::Exceeded,
23 core::Multiaddr,
24 gossipsub, identify,
25 identity::Keypair,
26 metrics::{Metrics, Recorder},
27 multiaddr::Protocol,
28 noise, ping, request_response,
29 swarm::{DialError, SwarmEvent},
30 tcp, yamux,
31};
32use nonzero_ext::nonzero;
33use tokio_stream::wrappers::IntervalStream;
34use tracing::{debug, error, info, trace, warn};
35
36use super::{
37 ForestBehaviour, ForestBehaviourEvent, Libp2pConfig,
38 chain_exchange::{ChainExchangeRequest, ChainExchangeResponse, make_chain_exchange_response},
39 discovery::{DerivedDiscoveryBehaviourEvent, PeerInfo},
40};
41use crate::libp2p::{
42 PeerManager, PeerOperation,
43 chain_exchange::ChainExchangeBehaviour,
44 discovery::DiscoveryEvent,
45 hello::{HelloBehaviour, HelloRequest, HelloResponse},
46 rpc::RequestResponseError,
47};
48
49pub(in crate::libp2p) mod metrics {
50 use prometheus_client::metrics::{family::Family, gauge::Gauge};
51 use std::sync::LazyLock;
52
53 use crate::metrics::KindLabel;
54
55 pub static NETWORK_CONTAINER_CAPACITIES: LazyLock<Family<KindLabel, Gauge>> = {
56 LazyLock::new(|| {
57 let metric = Family::default();
58 crate::metrics::default_registry().register(
59 "network_container_capacities",
60 "Capacity for each container",
61 metric.clone(),
62 );
63 metric
64 })
65 };
66
67 pub mod values {
68 use crate::metrics::KindLabel;
69
70 pub const HELLO_REQUEST_TABLE: KindLabel = KindLabel::new("hello_request_table");
71 pub const CHAIN_EXCHANGE_REQUEST_TABLE: KindLabel = KindLabel::new("cx_request_table");
72 }
73}
74
75crate::def_is_env_truthy!(libp2p_metrics_enabled, "FOREST_LIBP2P_METRICS_ENABLED");
76
77pub const PUBSUB_BLOCK_STR: &str = "/fil/blocks";
79pub const PUBSUB_MSG_STR: &str = "/fil/msgs";
81
82#[derive(Copy, Clone, Debug, strum::EnumIter, derive_more::Display)]
86pub enum PubsubTopic {
87 #[display("{PUBSUB_BLOCK_STR}")]
88 Blocks,
89 #[display("{PUBSUB_MSG_STR}")]
90 Messages,
91}
92
93impl PubsubTopic {
94 pub fn ident(self, network_name: impl std::fmt::Display) -> IdentTopic {
96 IdentTopic::new(format!("{self}/{network_name}"))
97 }
98}
99
100pub fn pubsub_topics(
102 network_name: impl std::fmt::Display + Copy,
103) -> impl Iterator<Item = IdentTopic> {
104 use strum::IntoEnumIterator as _;
105 PubsubTopic::iter().map(move |t| t.ident(network_name))
106}
107
108pub const BITSWAP_TIMEOUT: Duration = Duration::from_secs(30);
109
110#[allow(clippy::large_enum_variant)]
112#[derive(Debug)]
113pub enum NetworkEvent {
114 PubsubMessage {
115 message: PubsubMessage,
116 },
117 HelloRequestInbound,
118 HelloResponseOutbound {
119 source: PeerId,
120 request: HelloRequest,
121 },
122 HelloRequestOutbound,
123 HelloResponseInbound,
124 ChainExchangeRequestOutbound,
125 ChainExchangeResponseInbound,
126 ChainExchangeRequestInbound,
127 ChainExchangeResponseOutbound,
128 PeerConnected(PeerId),
129 PeerDisconnected(PeerId),
130}
131
132#[allow(clippy::large_enum_variant)]
134#[derive(Debug, Clone)]
135pub enum PubsubMessage {
136 Block(GossipBlock),
138 Message(SignedMessage),
140}
141
142#[derive(Debug)]
144pub enum NetworkMessage {
145 PubsubMessage {
146 topic: IdentTopic,
147 message: Vec<u8>,
148 },
149 ChainExchangeRequest {
150 peer_id: PeerId,
151 request: ChainExchangeRequest,
152 response_channel: flume::Sender<Result<ChainExchangeResponse, RequestResponseError>>,
153 },
154 HelloRequest {
155 peer_id: PeerId,
156 request: HelloRequest,
157 response_channel: flume::Sender<HelloResponse>,
158 },
159 BitswapRequest {
160 cid: Cid,
161 response_channel: flume::Sender<bool>,
162 },
163 JSONRPCRequest {
164 method: NetRPCMethods,
165 },
166}
167
168#[derive(Debug)]
170pub enum NetRPCMethods {
171 AddrsListen(flume::Sender<(PeerId, HashSet<Multiaddr>)>),
172 Peer(flume::Sender<Option<HashSet<Multiaddr>>>, PeerId),
173 Peers(flume::Sender<HashMap<PeerId, HashSet<Multiaddr>>>),
174 ProtectPeer(flume::Sender<()>, HashSet<PeerId>),
175 UnprotectPeer(flume::Sender<()>, HashSet<PeerId>),
176 ListProtectedPeers(flume::Sender<HashSet<PeerId>>),
177 Info(flume::Sender<NetInfoResult>),
178 Connect(flume::Sender<bool>, PeerId, HashSet<Multiaddr>),
179 Disconnect(flume::Sender<()>, PeerId),
180 AgentVersion(flume::Sender<Option<String>>, PeerId),
181 AutoNATStatus(flume::Sender<NatStatus>),
182}
183
184pub struct Libp2pService {
186 swarm: Swarm<ForestBehaviour>,
187 bootstrap_peers: HashMap<PeerId, Multiaddr>,
188 cs: ChainStore,
189 peer_manager: Arc<PeerManager>,
190 network_receiver_in: flume::Receiver<NetworkMessage>,
191 network_sender_in: Sender<NetworkMessage>,
192 network_receiver_out: flume::Receiver<NetworkEvent>,
193 network_sender_out: Sender<NetworkEvent>,
194 network_name: String,
195 genesis_cid: Cid,
196}
197
198impl Libp2pService {
199 pub async fn new(
200 config: Libp2pConfig,
201 cs: ChainStore,
202 peer_manager: Arc<PeerManager>,
203 net_keypair: Keypair,
204 network_name: GenesisNetworkName,
205 genesis_cid: Cid,
206 ) -> anyhow::Result<Self> {
207 let behaviour =
208 ForestBehaviour::new(&net_keypair, &config, &network_name, peer_manager.clone())
209 .await?;
210 let mut swarm = SwarmBuilder::with_existing_identity(net_keypair)
211 .with_tokio()
212 .with_tcp(
213 tcp::Config::default().nodelay(true),
214 noise::Config::new,
215 yamux::Config::default,
216 )?
217 .with_quic()
218 .with_dns()?
219 .with_bandwidth_metrics(&mut crate::metrics::collector_registry())
220 .with_behaviour(|_| behaviour)?
221 .with_swarm_config(|config| {
222 config
223 .with_notify_handler_buffer_size(nonzero!(20usize))
224 .with_per_connection_event_buffer_size(64)
225 .with_idle_connection_timeout(Duration::from_secs(60 * 10))
226 })
227 .build();
228
229 for topic in pubsub_topics(&network_name) {
231 swarm
232 .behaviour_mut()
233 .subscribe(&topic)
234 .with_context(|| format!("Failed to subscribe gossipsub topic {topic}"))?;
235 }
236
237 let (network_sender_in, network_receiver_in) = flume::unbounded();
238 let (network_sender_out, network_receiver_out) = flume::unbounded();
239
240 info!("p2p network peer id: {}", swarm.local_peer_id());
244
245 for addr in &config.listening_multiaddrs {
247 match swarm.listen_on(addr.clone()) {
248 Ok(id) => loop {
249 if let SwarmEvent::NewListenAddr {
250 address,
251 listener_id,
252 } = swarm.select_next_some().await
253 && id == listener_id
254 {
255 info!("p2p peer is now listening on: {address}");
256 break;
257 }
258 },
259 Err(err) => error!("Fail to listen on {addr}: {err}"),
260 }
261 }
262
263 if swarm.listeners().count() == 0 {
264 anyhow::bail!("p2p peer failed to listen on any network endpoints");
265 }
266
267 let bootstrap_peers = config
268 .bootstrap_peers
269 .iter()
270 .filter_map(|ma| match ma.iter().last() {
271 Some(Protocol::P2p(peer)) => Some((peer, ma.clone())),
272 _ => None,
273 })
274 .collect();
275
276 Ok(Libp2pService {
277 swarm,
278 bootstrap_peers,
279 cs,
280 peer_manager,
281 network_receiver_in,
282 network_sender_in,
283 network_receiver_out,
284 network_sender_out,
285 network_name: network_name.into(),
286 genesis_cid,
287 })
288 }
289
290 pub async fn run(mut self) -> anyhow::Result<()> {
293 info!("Running libp2p service");
294
295 if let Err(e) = self.swarm.behaviour_mut().bootstrap() {
297 warn!("Failed to bootstrap with Kademlia: {e:#}");
298 }
299
300 let bitswap_request_manager = self.swarm.behaviour().bitswap.request_manager();
301 let mut swarm_stream = self.swarm.fuse();
302 let mut network_stream = self.network_receiver_in.stream().fuse();
303 let mut interval =
304 IntervalStream::new(tokio::time::interval(Duration::from_secs(15))).fuse();
305 let pubsub_block_str = PubsubTopic::Blocks.ident(&self.network_name).to_string();
306 let pubsub_msg_str = PubsubTopic::Messages.ident(&self.network_name).to_string();
307
308 let (cx_response_tx, cx_response_rx) = flume::unbounded();
309
310 let mut cx_response_rx_stream = cx_response_rx.stream().fuse();
311 let mut bitswap_outbound_request_stream =
312 bitswap_request_manager.outbound_request_stream().fuse();
313 let mut bitswap_serve_response_stream = bitswap_request_manager
314 .outbound_serve_response_stream()
315 .fuse();
316 let mut peer_ops_rx_stream = self.peer_manager.peer_ops_rx().stream().fuse();
317 let metrics = if libp2p_metrics_enabled() {
318 Some(Metrics::new(&mut crate::metrics::collector_registry()))
319 } else {
320 None
321 };
322
323 const BOOTSTRAP_PEER_DIALER_INTERVAL: tokio::time::Duration =
324 tokio::time::Duration::from_secs(60);
325 let mut bootstrap_peer_dialer_interval_stream =
326 IntervalStream::new(tokio::time::interval_at(
327 tokio::time::Instant::now() + BOOTSTRAP_PEER_DIALER_INTERVAL,
328 BOOTSTRAP_PEER_DIALER_INTERVAL,
329 ))
330 .fuse();
331 loop {
332 select! {
333 swarm_event = swarm_stream.next() => match swarm_event {
334 Some(SwarmEvent::Behaviour(event)) => {
336 if let Some(m) = &metrics {
337 m.record(&event);
338 }
339 handle_forest_behaviour_event(
340 swarm_stream.get_mut(),
341 &bitswap_request_manager,
342 &self.peer_manager,
343 event,
344 &self.cs,
345 &self.genesis_cid,
346 &self.network_sender_out,
347 cx_response_tx.clone(),
348 &pubsub_block_str,
349 &pubsub_msg_str,).await;
350 },
351 None => { break; },
352 _ => { },
353 },
354 rpc_message = network_stream.next() => match rpc_message {
355 Some(message) => {
357 handle_network_message(
358 swarm_stream.get_mut(),
359 self.cs.db_owned(),
360 bitswap_request_manager.shallow_clone(),
361 message,
362 &self.network_sender_out,
363 &self.peer_manager).await;
364 }
365 None => { break; }
366 },
367 interval_event = interval.next() => if interval_event.is_some() {
368 trace!("Peers connected: {}", swarm_stream.get_mut().behaviour_mut().peers().len());
370 },
371 cs_pair_opt = cx_response_rx_stream.next() => {
372 if let Some((_request_id, channel, cx_response)) = cs_pair_opt {
373 let behaviour = swarm_stream.get_mut().behaviour_mut();
374 if let Err(e) = behaviour.chain_exchange.send_response(channel, cx_response) {
375 debug!("Error sending chain exchange response: {e:?}");
376 }
377 }
378 },
379 bitswap_outbound_request_opt = bitswap_outbound_request_stream.next() => {
380 if let Some((peer, request)) = bitswap_outbound_request_opt {
381 let bitswap = &mut swarm_stream.get_mut().behaviour_mut().bitswap;
382 bitswap.send_request(&peer, request);
383 }
384 }
385 bitswap_serve_response_opt = bitswap_serve_response_stream.next() => {
386 if let Some((peer, cid, response)) = bitswap_serve_response_opt {
387 let bitswap = &mut swarm_stream.get_mut().behaviour_mut().bitswap;
388 bitswap.send_response(&peer, (cid, response));
389 }
390 }
391 peer_ops_opt = peer_ops_rx_stream.next() => {
392 if let Some(peer_ops) = peer_ops_opt {
393 handle_peer_ops(swarm_stream.get_mut(), peer_ops, &self.bootstrap_peers);
394 }
395 },
396 _ = bootstrap_peer_dialer_interval_stream.next() => {
397 dial_to_bootstrap_peers_if_needed(swarm_stream.get_mut(), &self.bootstrap_peers);
398 }
399 };
400 }
401 Ok(())
402 }
403
404 pub fn network_sender(&self) -> Sender<NetworkMessage> {
406 self.network_sender_in.clone()
407 }
408
409 pub fn network_receiver(&self) -> flume::Receiver<NetworkEvent> {
411 self.network_receiver_out.clone()
412 }
413
414 pub fn peer_manager(&self) -> &Arc<PeerManager> {
415 &self.peer_manager
416 }
417}
418
419fn dial_to_bootstrap_peers_if_needed(
420 swarm: &mut Swarm<ForestBehaviour>,
421 bootstrap_peers: &HashMap<PeerId, Multiaddr>,
422) {
423 for (peer, ma) in bootstrap_peers {
424 if !swarm.behaviour().peers().contains(peer) {
425 info!("Re-dialing to bootstrap peer at {ma}");
426 if let Err(e) = swarm.dial(ma.clone()) {
427 warn!("{e}");
428 }
429 }
430 }
431}
432
433fn handle_peer_ops(
434 swarm: &mut Swarm<ForestBehaviour>,
435 peer_ops: PeerOperation,
436 bootstrap_peers: &HashMap<PeerId, Multiaddr>,
437) {
438 use PeerOperation::*;
439 match peer_ops {
440 Ban {
441 peer,
442 user_agent,
443 reason,
444 } => {
445 if !bootstrap_peers.contains_key(&peer) {
447 let user_agent = user_agent.unwrap_or_default();
448 debug!(%peer, %user_agent, %reason, "Banning peer");
449 swarm.behaviour_mut().blocked_peers.block_peer(peer);
450 }
451 }
452 Unban(peer) => {
453 debug!(%peer, "Unbanning peer");
454 swarm.behaviour_mut().blocked_peers.unblock_peer(peer);
455 }
456 }
457}
458
459async fn handle_network_message(
460 swarm: &mut Swarm<ForestBehaviour>,
461 store: impl BitswapStoreReadWrite + ShallowClone,
462 bitswap_request_manager: Arc<BitswapRequestManager>,
463 message: NetworkMessage,
464 network_sender_out: &Sender<NetworkEvent>,
465 peer_manager: &PeerManager,
466) {
467 match message {
468 NetworkMessage::PubsubMessage { topic, message } => {
469 match swarm.behaviour_mut().publish(topic, message) {
470 Ok(_) => (),
471 Err(gossipsub::PublishError::Duplicate) => {
472 debug!("Failed to send gossipsub message: duplicate");
475 }
476 Err(e) => {
477 warn!("Failed to send gossipsub message: {e:#}");
478 }
479 }
480 }
481 NetworkMessage::HelloRequest {
482 peer_id,
483 request,
484 response_channel,
485 } => {
486 let _request_id =
487 swarm
488 .behaviour_mut()
489 .hello
490 .send_request(&peer_id, request, response_channel);
491 emit_event(network_sender_out, NetworkEvent::HelloRequestOutbound).await;
492 }
493 NetworkMessage::ChainExchangeRequest {
494 peer_id,
495 request,
496 response_channel,
497 } => {
498 let _request_id = swarm.behaviour_mut().chain_exchange.send_request(
499 &peer_id,
500 request,
501 response_channel,
502 );
503 emit_event(
504 network_sender_out,
505 NetworkEvent::ChainExchangeRequestOutbound,
506 )
507 .await;
508 }
509 NetworkMessage::BitswapRequest {
510 cid,
511 response_channel,
512 } => {
513 bitswap_request_manager.get_block(
514 store,
515 cid,
516 BITSWAP_TIMEOUT,
517 Some(response_channel),
518 None,
519 );
520 }
521 NetworkMessage::JSONRPCRequest { method } => {
522 match method {
523 NetRPCMethods::AddrsListen(response_channel) => {
524 let listeners = Swarm::listeners(swarm).cloned().collect();
525 let peer_id = Swarm::local_peer_id(swarm);
526 response_channel.send_or_warn((*peer_id, listeners));
527 }
528 NetRPCMethods::Peer(response_channel, peer) => {
529 let addresses = swarm.behaviour().peer_addresses().get(&peer).cloned();
530 response_channel.send_or_warn(addresses);
531 }
532 NetRPCMethods::Peers(response_channel) => {
533 let peer_addresses = swarm.behaviour().peer_addresses();
534 response_channel.send_or_warn(peer_addresses);
535 }
536 NetRPCMethods::ProtectPeer(tx, peer_ids) => {
537 peer_ids.into_iter().for_each(|peer_id| {
538 peer_manager.protect_peer(peer_id);
539 });
540 tx.send_or_warn(());
541 }
542 NetRPCMethods::ListProtectedPeers(tx) => {
543 tx.send_or_warn(peer_manager.list_protected_peers());
544 }
545 NetRPCMethods::UnprotectPeer(tx, peer_ids) => {
546 peer_ids.iter().for_each(|peer_id| {
547 peer_manager.unprotect_peer(peer_id);
548 });
549 tx.send_or_warn(());
550 }
551 NetRPCMethods::Info(response_channel) => {
552 response_channel.send_or_warn(swarm.network_info().into());
553 }
554 NetRPCMethods::Connect(response_channel, peer_id, addresses) => {
555 let mut success = false;
556 for mut multiaddr in addresses {
557 multiaddr.push(Protocol::P2p(peer_id));
558
559 match Swarm::dial(swarm, multiaddr.clone()) {
560 Ok(_) => {
561 info!("Dialed {multiaddr}");
562 success = true;
563 break;
564 }
565 Err(e) => {
566 match e {
567 DialError::Denied { cause } => {
568 if let Some(cause) = cause.downcast_ref::<Exceeded>() {
570 error!(
571 "Denied dialing (limits exceeded) {multiaddr}: {cause}"
572 );
573 } else {
574 error!("Denied dialing {multiaddr}: {cause}")
575 }
576 }
577 e => {
578 error!("Failed to dial {multiaddr}: {e}");
579 }
580 };
581 }
582 };
583 }
584
585 response_channel.send_or_warn(success);
586 }
587 NetRPCMethods::Disconnect(response_channel, peer_id) => {
588 let _ = Swarm::disconnect_peer_id(swarm, peer_id);
589 response_channel.send_or_warn(());
590 }
591 NetRPCMethods::AgentVersion(response_channel, peer_id) => {
592 let agent_version = swarm.behaviour().peer_info(&peer_id).and_then(|info| {
593 info.identify_info
594 .as_ref()
595 .map(|id| id.agent_version.clone())
596 });
597 response_channel.send_or_warn(agent_version);
598 }
599 NetRPCMethods::AutoNATStatus(response_channel) => {
600 let nat_status = swarm.behaviour().discovery.nat_status();
601 response_channel.send_or_warn(nat_status);
602 }
603 }
604 }
605 }
606}
607
608async fn handle_discovery_event(
609 peer_info_map: &HashMap<PeerId, PeerInfo>,
610 discovery_out: DiscoveryEvent,
611 network_sender_out: &Sender<NetworkEvent>,
612 peer_manager: &PeerManager,
613) {
614 match discovery_out {
615 DiscoveryEvent::PeerConnected(peer_id) => {
616 trace!("Peer connected, {peer_id}");
617 emit_event(network_sender_out, NetworkEvent::PeerConnected(peer_id)).await;
618 }
619 DiscoveryEvent::PeerDisconnected(peer_id) => {
620 trace!("Peer disconnected, {peer_id}");
621 emit_event(network_sender_out, NetworkEvent::PeerDisconnected(peer_id)).await;
622 }
623 DiscoveryEvent::Discovery(discovery_event) => match &*discovery_event {
624 DerivedDiscoveryBehaviourEvent::Identify(identify::Event::Received {
625 peer_id,
626 info,
627 ..
628 }) => {
629 let protocols = HashSet::from_iter(info.protocols.iter().map(|p| p.to_string()));
630 if !protocols.contains(super::hello::HELLO_PROTOCOL_NAME) {
631 peer_manager
632 .ban_peer_with_default_duration(
633 *peer_id,
634 "hello protocol unsupported",
635 |p| get_user_agent(peer_info_map, p),
636 )
637 .await;
638 } else if !protocols.contains(super::chain_exchange::CHAIN_EXCHANGE_PROTOCOL_NAME) {
639 peer_manager
640 .ban_peer_with_default_duration(
641 *peer_id,
642 "chain exchange protocol unsupported",
643 |p| get_user_agent(peer_info_map, p),
644 )
645 .await;
646 }
647 }
648 DerivedDiscoveryBehaviourEvent::Identify(_) => {}
649 _ => {}
650 },
651 }
652}
653
654async fn handle_gossip_event(
655 e: gossipsub::Event,
656 network_sender_out: &Sender<NetworkEvent>,
657 pubsub_block_str: &str,
658 pubsub_msg_str: &str,
659) {
660 if let gossipsub::Event::Message {
661 propagation_source: source,
662 message,
663 ..
664 } = e
665 {
666 let topic = message.topic.as_str();
667 let message = message.data;
668 trace!("Got a Gossip Message from {:?}", source);
669 if topic == pubsub_block_str {
670 match from_slice_with_fallback::<GossipBlock>(&message) {
671 Ok(b) => {
672 emit_event(
673 network_sender_out,
674 NetworkEvent::PubsubMessage {
675 message: PubsubMessage::Block(b),
676 },
677 )
678 .await;
679 }
680 Err(e) => {
681 warn!("Gossip Block from peer {source:?} could not be deserialized: {e:#}",);
682 }
683 }
684 } else if topic == pubsub_msg_str {
685 match from_slice_with_fallback::<SignedMessage>(&message) {
686 Ok(m) => {
687 emit_event(
688 network_sender_out,
689 NetworkEvent::PubsubMessage {
690 message: PubsubMessage::Message(m),
691 },
692 )
693 .await;
694 }
695 Err(e) => {
696 warn!("Gossip Message from peer {source:?} could not be deserialized: {e:#}");
697 }
698 }
699 } else {
700 warn!("Getting gossip messages from unknown topic: {topic}");
701 }
702 }
703}
704
705fn nanos_since_unix_epoch() -> u64 {
708 u64::try_from(UNIX_EPOCH.elapsed().unwrap_or_default().as_nanos()).unwrap_or(u64::MAX)
709}
710
711async fn handle_hello_event(
712 peer_info_map: &HashMap<PeerId, PeerInfo>,
713 hello: &mut HelloBehaviour,
714 event: request_response::Event<HelloRequest, HelloResponse, HelloResponse>,
715 peer_manager: &PeerManager,
716 genesis_cid: &Cid,
717 network_sender_out: &Sender<NetworkEvent>,
718) {
719 match event {
720 request_response::Event::Message { peer, message, .. } => match message {
721 request_response::Message::Request {
722 request, channel, ..
723 } => {
724 emit_event(network_sender_out, NetworkEvent::HelloRequestInbound).await;
725
726 let arrival = nanos_since_unix_epoch();
727
728 trace!("Received hello request: {:?}", request);
729 if &request.genesis_cid != genesis_cid {
730 peer_manager
731 .ban_peer_with_default_duration(
732 peer,
733 format!(
734 "Genesis hash mismatch: {} received, {genesis_cid} expected",
735 request.genesis_cid
736 ),
737 |p| get_user_agent(peer_info_map, p),
738 )
739 .await;
740 } else {
741 let sent = nanos_since_unix_epoch();
742
743 if let Err(e) = hello.send_response(channel, HelloResponse { arrival, sent }) {
746 warn!("Failed to send HelloResponse: {e:?}");
747 } else {
748 emit_event(
749 network_sender_out,
750 NetworkEvent::HelloResponseOutbound {
751 source: peer,
752 request,
753 },
754 )
755 .await;
756 }
757 }
758 }
759 request_response::Message::Response {
760 request_id,
761 response,
762 } => {
763 emit_event(network_sender_out, NetworkEvent::HelloResponseInbound).await;
764 hello.handle_response(&request_id, response).await;
765 }
766 },
767 request_response::Event::OutboundFailure {
768 request_id,
769 peer,
770 error,
771 ..
772 } => {
773 hello.on_outbound_failure(&request_id);
774 match error {
775 request_response::OutboundFailure::UnsupportedProtocols => {
776 peer_manager
777 .ban_peer_with_default_duration(peer, "Hello protocol unsupported", |p| {
778 get_user_agent(peer_info_map, p)
779 })
780 .await;
781 }
782 _ => {
783 peer_manager.mark_peer_bad(peer, format!("Hello outbound failure {error}"));
784 }
785 }
786 }
787 request_response::Event::InboundFailure { .. } => {}
788 request_response::Event::ResponseSent { .. } => (),
789 }
790}
791
792async fn handle_ping_event(ping_event: ping::Event) {
793 match ping_event.result {
794 Ok(rtt) => {
795 trace!(
796 "PingSuccess::Ping rtt to {} is {} ms",
797 ping_event.peer,
798 rtt.as_millis()
799 );
800 }
801 Err(ping::Failure::Unsupported) => {
802 debug!(peer=%ping_event.peer, "Ping protocol unsupported");
803 }
804 Err(ping::Failure::Timeout) => {
805 debug!("Ping timeout: {}", ping_event.peer);
806 }
807 Err(ping::Failure::Other { error }) => {
808 debug!("Ping failure: {error}");
809 }
810 }
811}
812
813async fn handle_chain_exchange_event(
814 chain_exchange: &mut ChainExchangeBehaviour,
815 ce_event: request_response::Event<ChainExchangeRequest, ChainExchangeResponse>,
816 db: &ChainStore,
817 network_sender_out: &Sender<NetworkEvent>,
818 cx_response_tx: Sender<(
819 request_response::InboundRequestId,
820 request_response::ResponseChannel<ChainExchangeResponse>,
821 ChainExchangeResponse,
822 )>,
823) {
824 const CHAIN_EXCHANGE_RESPONSE_TIMEOUT: Duration = Duration::from_mins(5);
825 match ce_event {
826 request_response::Event::Message { peer, message, .. } => match message {
827 request_response::Message::Request {
828 request,
829 channel,
830 request_id,
831 } => {
832 let Some(per_peer_permit) = chain_exchange.try_acquire_peer_permit(peer) else {
833 debug!("Rejecting chain_exchange request from {peer}: per-peer cap reached");
834 let _ = chain_exchange.send_response(
835 channel,
836 ChainExchangeResponse::go_away("per-peer concurrent request cap reached"),
837 );
838 return;
839 };
840 let Some(global_permit) = chain_exchange.try_acquire_request_permit() else {
841 debug!("Rejecting chain_exchange request from {peer}: global cap reached");
842 let _ = chain_exchange.send_response(
843 channel,
844 ChainExchangeResponse::go_away("global concurrent request cap reached"),
845 );
846 return;
847 };
848
849 trace!(
850 "Received chain_exchange request (request_id:{request_id}, peer_id: {peer:?})",
851 );
852 emit_event(
853 network_sender_out,
854 NetworkEvent::ChainExchangeRequestInbound,
855 )
856 .await;
857
858 let db = db.shallow_clone();
859 tokio::task::spawn(tokio::time::timeout(
860 CHAIN_EXCHANGE_RESPONSE_TIMEOUT,
861 async move {
862 let _per_peer_permit = per_peer_permit;
863 let _global_permit = global_permit;
864 if let Err(e) = cx_response_tx
865 .send_async((
866 request_id,
867 channel,
868 make_chain_exchange_response(&db, &request),
869 ))
870 .await
871 {
872 debug!("Failed to send ChainExchangeResponse: {e:?}");
873 }
874 },
875 ));
876 }
877 request_response::Message::Response {
878 request_id,
879 response,
880 } => {
881 emit_event(
882 network_sender_out,
883 NetworkEvent::ChainExchangeResponseInbound,
884 )
885 .await;
886 chain_exchange
887 .handle_inbound_response(&request_id, response)
888 .await;
889 }
890 },
891 request_response::Event::OutboundFailure {
892 request_id, error, ..
893 } => {
894 chain_exchange.on_outbound_error(&request_id, error);
895 }
896 request_response::Event::InboundFailure { peer, error, .. } => {
897 debug!(
898 "ChainExchange inbound error (peer: {:?}): {:?}",
899 peer, error
900 );
901 }
902 request_response::Event::ResponseSent { .. } => {
903 emit_event(
904 network_sender_out,
905 NetworkEvent::ChainExchangeResponseOutbound,
906 )
907 .await;
908 }
909 }
910}
911
912#[allow(clippy::too_many_arguments)]
913async fn handle_forest_behaviour_event(
914 swarm: &mut Swarm<ForestBehaviour>,
915 bitswap_request_manager: &Arc<BitswapRequestManager>,
916 peer_manager: &PeerManager,
917 event: ForestBehaviourEvent,
918 db: &ChainStore,
919 genesis_cid: &Cid,
920 network_sender_out: &Sender<NetworkEvent>,
921 cx_response_tx: Sender<(
922 request_response::InboundRequestId,
923 request_response::ResponseChannel<ChainExchangeResponse>,
924 ChainExchangeResponse,
925 )>,
926 pubsub_block_str: &str,
927 pubsub_msg_str: &str,
928) {
929 match event {
930 ForestBehaviourEvent::Discovery(discovery_out) => {
931 handle_discovery_event(
932 &swarm.behaviour().discovery.peer_info,
933 discovery_out,
934 network_sender_out,
935 peer_manager,
936 )
937 .await
938 }
939 ForestBehaviourEvent::Gossipsub(e) => {
940 handle_gossip_event(e, network_sender_out, pubsub_block_str, pubsub_msg_str).await
941 }
942 ForestBehaviourEvent::Hello(rr_event) => {
943 let behaviour_mut = swarm.behaviour_mut();
944 handle_hello_event(
945 &behaviour_mut.discovery.peer_info,
946 &mut behaviour_mut.hello,
947 rr_event,
948 peer_manager,
949 genesis_cid,
950 network_sender_out,
951 )
952 .await
953 }
954 ForestBehaviourEvent::Bitswap(event) => {
955 if let Err(e) = bitswap_request_manager.handle_event(
956 &mut swarm.behaviour_mut().bitswap,
957 db.db(),
958 event,
959 ) {
960 warn!("bitswap: {e:#}");
961 }
962 }
963 ForestBehaviourEvent::Ping(ping_event) => handle_ping_event(ping_event).await,
964 ForestBehaviourEvent::ConnectionLimits(_) => {}
965 ForestBehaviourEvent::BlockedPeers(_) => {}
966 ForestBehaviourEvent::ChainExchange(ce_event) => {
967 handle_chain_exchange_event(
968 &mut swarm.behaviour_mut().chain_exchange,
969 ce_event,
970 db,
971 network_sender_out,
972 cx_response_tx,
973 )
974 .await
975 }
976 }
977}
978
979async fn emit_event(sender: &Sender<NetworkEvent>, event: NetworkEvent) {
980 if sender.send_async(event).await.is_err() {
981 error!("Failed to emit event: Network channel receiver has been dropped");
982 }
983}
984
985fn get_user_agent(peer_info_map: &HashMap<PeerId, PeerInfo>, peer: &PeerId) -> Option<String> {
986 peer_info_map
987 .get(peer)
988 .and_then(|i| i.identify_info.as_ref())
989 .map(|i| i.agent_version.clone())
990}