Skip to main content

forest/libp2p/
service.rs

1// Copyright 2019-2026 ChainSafe Systems
2// SPDX-License-Identifier: Apache-2.0, MIT
3
4use std::time::{Duration, SystemTime, 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
77/// `Gossipsub` Filecoin blocks topic identifier.
78pub const PUBSUB_BLOCK_STR: &str = "/fil/blocks";
79/// `Gossipsub` Filecoin messages topic identifier.
80pub const PUBSUB_MSG_STR: &str = "/fil/msgs";
81
82/// Gossipsub topics Forest uses. Subscription, the subscription-filter
83/// whitelist, and peer-score params all iterate the variants, so adding one is
84/// handled everywhere.
85#[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    /// Full topic on `network_name`, e.g. `/fil/blocks/<net>`.
95    pub fn ident(self, network_name: impl std::fmt::Display) -> IdentTopic {
96        IdentTopic::new(format!("{self}/{network_name}"))
97    }
98}
99
100/// All gossipsub topics on `network_name`.
101pub 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/// Events emitted by this Service.
111#[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/// Message types that can come over `GossipSub`
133#[allow(clippy::large_enum_variant)]
134#[derive(Debug, Clone)]
135pub enum PubsubMessage {
136    /// Messages that come over the block topic
137    Block(GossipBlock),
138    /// Messages that come over the message topic
139    Message(SignedMessage),
140}
141
142/// Messages into the service to handle.
143#[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/// Network RPC API methods used to gather data from libp2p node.
169#[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
184/// The `Libp2pService` listens to events from the libp2p swarm.
185pub 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        // Subscribe to gossipsub topics with the network name suffix
230        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        // Hint at the multihash which has to go in the `/p2p/<multihash>` part of the
241        // peer's multiaddress. Useful if others want to use this node to bootstrap
242        // from.
243        info!("p2p network peer id: {}", swarm.local_peer_id());
244
245        // Listen on network endpoints before being detached and connecting to any peers.
246        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    /// Starts the libp2p service networking stack. This Future resolves when
291    /// shutdown occurs.
292    pub async fn run(mut self) -> anyhow::Result<()> {
293        info!("Running libp2p service");
294
295        // Bootstrap with Kademlia
296        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 peer_ops_rx_stream = self.peer_manager.peer_ops_rx().stream().fuse();
314        let metrics = if libp2p_metrics_enabled() {
315            Some(Metrics::new(&mut crate::metrics::collector_registry()))
316        } else {
317            None
318        };
319
320        const BOOTSTRAP_PEER_DIALER_INTERVAL: tokio::time::Duration =
321            tokio::time::Duration::from_secs(60);
322        let mut bootstrap_peer_dialer_interval_stream =
323            IntervalStream::new(tokio::time::interval_at(
324                tokio::time::Instant::now() + BOOTSTRAP_PEER_DIALER_INTERVAL,
325                BOOTSTRAP_PEER_DIALER_INTERVAL,
326            ))
327            .fuse();
328        loop {
329            select! {
330                swarm_event = swarm_stream.next() => match swarm_event {
331                    // outbound events
332                    Some(SwarmEvent::Behaviour(event)) => {
333                        if let Some(m) = &metrics {
334                            m.record(&event);
335                        }
336                        handle_forest_behaviour_event(
337                            swarm_stream.get_mut(),
338                            &bitswap_request_manager,
339                            &self.peer_manager,
340                            event,
341                            &self.cs,
342                            &self.genesis_cid,
343                            &self.network_sender_out,
344                            cx_response_tx.clone(),
345                            &pubsub_block_str,
346                            &pubsub_msg_str,).await;
347                    },
348                    None => { break; },
349                    _ => { },
350                },
351                rpc_message = network_stream.next() => match rpc_message {
352                    // Inbound messages
353                    Some(message) => {
354                        handle_network_message(
355                            swarm_stream.get_mut(),
356                            self.cs.db_owned(),
357                            bitswap_request_manager.shallow_clone(),
358                            message,
359                            &self.network_sender_out,
360                            &self.peer_manager).await;
361                    }
362                    None => { break; }
363                },
364                interval_event = interval.next() => if interval_event.is_some() {
365                    // Print peer count on an interval.
366                    trace!("Peers connected: {}", swarm_stream.get_mut().behaviour_mut().peers().len());
367                },
368                cs_pair_opt = cx_response_rx_stream.next() => {
369                    if let Some((_request_id, channel, cx_response)) = cs_pair_opt {
370                        let behaviour = swarm_stream.get_mut().behaviour_mut();
371                        if let Err(e) = behaviour.chain_exchange.send_response(channel, cx_response) {
372                            debug!("Error sending chain exchange response: {e:?}");
373                        }
374                    }
375                },
376                bitswap_outbound_request_opt = bitswap_outbound_request_stream.next() => {
377                    if let Some((peer, request)) = bitswap_outbound_request_opt {
378                        let bitswap = &mut swarm_stream.get_mut().behaviour_mut().bitswap;
379                        bitswap.send_request(&peer, request);
380                    }
381                }
382                peer_ops_opt = peer_ops_rx_stream.next() => {
383                    if let Some(peer_ops) = peer_ops_opt {
384                        handle_peer_ops(swarm_stream.get_mut(), peer_ops, &self.bootstrap_peers);
385                    }
386                },
387                _ = bootstrap_peer_dialer_interval_stream.next() => {
388                    dial_to_bootstrap_peers_if_needed(swarm_stream.get_mut(), &self.bootstrap_peers);
389                }
390            };
391        }
392        Ok(())
393    }
394
395    /// Returns a sender which allows sending messages to the libp2p service.
396    pub fn network_sender(&self) -> Sender<NetworkMessage> {
397        self.network_sender_in.clone()
398    }
399
400    /// Returns a receiver to listen to network events emitted from the service.
401    pub fn network_receiver(&self) -> flume::Receiver<NetworkEvent> {
402        self.network_receiver_out.clone()
403    }
404
405    pub fn peer_manager(&self) -> &Arc<PeerManager> {
406        &self.peer_manager
407    }
408}
409
410fn dial_to_bootstrap_peers_if_needed(
411    swarm: &mut Swarm<ForestBehaviour>,
412    bootstrap_peers: &HashMap<PeerId, Multiaddr>,
413) {
414    for (peer, ma) in bootstrap_peers {
415        if !swarm.behaviour().peers().contains(peer) {
416            info!("Re-dialing to bootstrap peer at {ma}");
417            if let Err(e) = swarm.dial(ma.clone()) {
418                warn!("{e}");
419            }
420        }
421    }
422}
423
424fn handle_peer_ops(
425    swarm: &mut Swarm<ForestBehaviour>,
426    peer_ops: PeerOperation,
427    bootstrap_peers: &HashMap<PeerId, Multiaddr>,
428) {
429    use PeerOperation::*;
430    match peer_ops {
431        Ban {
432            peer,
433            user_agent,
434            reason,
435        } => {
436            // Do not ban bootstrap nodes
437            if !bootstrap_peers.contains_key(&peer) {
438                let user_agent = user_agent.unwrap_or_default();
439                debug!(%peer, %user_agent, %reason, "Banning peer");
440                swarm.behaviour_mut().blocked_peers.block_peer(peer);
441            }
442        }
443        Unban(peer) => {
444            debug!(%peer, "Unbanning peer");
445            swarm.behaviour_mut().blocked_peers.unblock_peer(peer);
446        }
447    }
448}
449
450async fn handle_network_message(
451    swarm: &mut Swarm<ForestBehaviour>,
452    store: impl BitswapStoreReadWrite + ShallowClone,
453    bitswap_request_manager: Arc<BitswapRequestManager>,
454    message: NetworkMessage,
455    network_sender_out: &Sender<NetworkEvent>,
456    peer_manager: &PeerManager,
457) {
458    match message {
459        NetworkMessage::PubsubMessage { topic, message } => {
460            match swarm.behaviour_mut().publish(topic, message) {
461                Ok(_) => (),
462                Err(gossipsub::PublishError::Duplicate) => {
463                    // Safe to ignore, the message has already been published and deduped by gossipsub, so no need to log this as a warning.
464                    // This matches `go-libp2p-pubsub` behavior https://github.com/libp2p/go-libp2p-pubsub/blob/v0.15.0/topic.go#L240-L242
465                    debug!("Failed to send gossipsub message: duplicate");
466                }
467                Err(e) => {
468                    warn!("Failed to send gossipsub message: {e:#}");
469                }
470            }
471        }
472        NetworkMessage::HelloRequest {
473            peer_id,
474            request,
475            response_channel,
476        } => {
477            let _request_id =
478                swarm
479                    .behaviour_mut()
480                    .hello
481                    .send_request(&peer_id, request, response_channel);
482            emit_event(network_sender_out, NetworkEvent::HelloRequestOutbound).await;
483        }
484        NetworkMessage::ChainExchangeRequest {
485            peer_id,
486            request,
487            response_channel,
488        } => {
489            let _request_id = swarm.behaviour_mut().chain_exchange.send_request(
490                &peer_id,
491                request,
492                response_channel,
493            );
494            emit_event(
495                network_sender_out,
496                NetworkEvent::ChainExchangeRequestOutbound,
497            )
498            .await;
499        }
500        NetworkMessage::BitswapRequest {
501            cid,
502            response_channel,
503        } => {
504            bitswap_request_manager.get_block(
505                store,
506                cid,
507                BITSWAP_TIMEOUT,
508                Some(response_channel),
509                None,
510            );
511        }
512        NetworkMessage::JSONRPCRequest { method } => {
513            match method {
514                NetRPCMethods::AddrsListen(response_channel) => {
515                    let listeners = Swarm::listeners(swarm).cloned().collect();
516                    let peer_id = Swarm::local_peer_id(swarm);
517                    response_channel.send_or_warn((*peer_id, listeners));
518                }
519                NetRPCMethods::Peer(response_channel, peer) => {
520                    let addresses = swarm.behaviour().peer_addresses().get(&peer).cloned();
521                    response_channel.send_or_warn(addresses);
522                }
523                NetRPCMethods::Peers(response_channel) => {
524                    let peer_addresses = swarm.behaviour().peer_addresses();
525                    response_channel.send_or_warn(peer_addresses);
526                }
527                NetRPCMethods::ProtectPeer(tx, peer_ids) => {
528                    peer_ids.into_iter().for_each(|peer_id| {
529                        peer_manager.protect_peer(peer_id);
530                    });
531                    tx.send_or_warn(());
532                }
533                NetRPCMethods::ListProtectedPeers(tx) => {
534                    tx.send_or_warn(peer_manager.list_protected_peers());
535                }
536                NetRPCMethods::UnprotectPeer(tx, peer_ids) => {
537                    peer_ids.iter().for_each(|peer_id| {
538                        peer_manager.unprotect_peer(peer_id);
539                    });
540                    tx.send_or_warn(());
541                }
542                NetRPCMethods::Info(response_channel) => {
543                    response_channel.send_or_warn(swarm.network_info().into());
544                }
545                NetRPCMethods::Connect(response_channel, peer_id, addresses) => {
546                    let mut success = false;
547                    for mut multiaddr in addresses {
548                        multiaddr.push(Protocol::P2p(peer_id));
549
550                        match Swarm::dial(swarm, multiaddr.clone()) {
551                            Ok(_) => {
552                                info!("Dialed {multiaddr}");
553                                success = true;
554                                break;
555                            }
556                            Err(e) => {
557                                match e {
558                                    DialError::Denied { cause } => {
559                                        // try to get a more specific error cause
560                                        if let Some(cause) = cause.downcast_ref::<Exceeded>() {
561                                            error!(
562                                                "Denied dialing (limits exceeded) {multiaddr}: {cause}"
563                                            );
564                                        } else {
565                                            error!("Denied dialing {multiaddr}: {cause}")
566                                        }
567                                    }
568                                    e => {
569                                        error!("Failed to dial {multiaddr}: {e}");
570                                    }
571                                };
572                            }
573                        };
574                    }
575
576                    response_channel.send_or_warn(success);
577                }
578                NetRPCMethods::Disconnect(response_channel, peer_id) => {
579                    let _ = Swarm::disconnect_peer_id(swarm, peer_id);
580                    response_channel.send_or_warn(());
581                }
582                NetRPCMethods::AgentVersion(response_channel, peer_id) => {
583                    let agent_version = swarm.behaviour().peer_info(&peer_id).and_then(|info| {
584                        info.identify_info
585                            .as_ref()
586                            .map(|id| id.agent_version.clone())
587                    });
588                    response_channel.send_or_warn(agent_version);
589                }
590                NetRPCMethods::AutoNATStatus(response_channel) => {
591                    let nat_status = swarm.behaviour().discovery.nat_status();
592                    response_channel.send_or_warn(nat_status);
593                }
594            }
595        }
596    }
597}
598
599async fn handle_discovery_event(
600    peer_info_map: &HashMap<PeerId, PeerInfo>,
601    discovery_out: DiscoveryEvent,
602    network_sender_out: &Sender<NetworkEvent>,
603    peer_manager: &PeerManager,
604) {
605    match discovery_out {
606        DiscoveryEvent::PeerConnected(peer_id) => {
607            trace!("Peer connected, {peer_id}");
608            emit_event(network_sender_out, NetworkEvent::PeerConnected(peer_id)).await;
609        }
610        DiscoveryEvent::PeerDisconnected(peer_id) => {
611            trace!("Peer disconnected, {peer_id}");
612            emit_event(network_sender_out, NetworkEvent::PeerDisconnected(peer_id)).await;
613        }
614        DiscoveryEvent::Discovery(discovery_event) => match &*discovery_event {
615            DerivedDiscoveryBehaviourEvent::Identify(identify::Event::Received {
616                peer_id,
617                info,
618                ..
619            }) => {
620                let protocols = HashSet::from_iter(info.protocols.iter().map(|p| p.to_string()));
621                if !protocols.contains(super::hello::HELLO_PROTOCOL_NAME) {
622                    peer_manager
623                        .ban_peer_with_default_duration(
624                            *peer_id,
625                            "hello protocol unsupported",
626                            |p| get_user_agent(peer_info_map, p),
627                        )
628                        .await;
629                } else if !protocols.contains(super::chain_exchange::CHAIN_EXCHANGE_PROTOCOL_NAME) {
630                    peer_manager
631                        .ban_peer_with_default_duration(
632                            *peer_id,
633                            "chain exchange protocol unsupported",
634                            |p| get_user_agent(peer_info_map, p),
635                        )
636                        .await;
637                }
638            }
639            DerivedDiscoveryBehaviourEvent::Identify(_) => {}
640            _ => {}
641        },
642    }
643}
644
645async fn handle_gossip_event(
646    e: gossipsub::Event,
647    network_sender_out: &Sender<NetworkEvent>,
648    pubsub_block_str: &str,
649    pubsub_msg_str: &str,
650) {
651    if let gossipsub::Event::Message {
652        propagation_source: source,
653        message,
654        ..
655    } = e
656    {
657        let topic = message.topic.as_str();
658        let message = message.data;
659        trace!("Got a Gossip Message from {:?}", source);
660        if topic == pubsub_block_str {
661            match from_slice_with_fallback::<GossipBlock>(&message) {
662                Ok(b) => {
663                    emit_event(
664                        network_sender_out,
665                        NetworkEvent::PubsubMessage {
666                            message: PubsubMessage::Block(b),
667                        },
668                    )
669                    .await;
670                }
671                Err(e) => {
672                    warn!("Gossip Block from peer {source:?} could not be deserialized: {e:#}",);
673                }
674            }
675        } else if topic == pubsub_msg_str {
676            match from_slice_with_fallback::<SignedMessage>(&message) {
677                Ok(m) => {
678                    emit_event(
679                        network_sender_out,
680                        NetworkEvent::PubsubMessage {
681                            message: PubsubMessage::Message(m),
682                        },
683                    )
684                    .await;
685                }
686                Err(e) => {
687                    warn!("Gossip Message from peer {source:?} could not be deserialized: {e:#}");
688                }
689            }
690        } else {
691            warn!("Getting gossip messages from unknown topic: {topic}");
692        }
693    }
694}
695
696async fn handle_hello_event(
697    peer_info_map: &HashMap<PeerId, PeerInfo>,
698    hello: &mut HelloBehaviour,
699    event: request_response::Event<HelloRequest, HelloResponse, HelloResponse>,
700    peer_manager: &PeerManager,
701    genesis_cid: &Cid,
702    network_sender_out: &Sender<NetworkEvent>,
703) {
704    match event {
705        request_response::Event::Message { peer, message, .. } => match message {
706            request_response::Message::Request {
707                request, channel, ..
708            } => {
709                emit_event(network_sender_out, NetworkEvent::HelloRequestInbound).await;
710
711                let arrival = SystemTime::now()
712                    .duration_since(UNIX_EPOCH)
713                    .expect("System time before unix epoch")
714                    .as_nanos()
715                    .try_into()
716                    .expect("System time since unix epoch should not exceed u64");
717
718                trace!("Received hello request: {:?}", request);
719                if &request.genesis_cid != genesis_cid {
720                    peer_manager
721                        .ban_peer_with_default_duration(
722                            peer,
723                            format!(
724                                "Genesis hash mismatch: {} received, {genesis_cid} expected",
725                                request.genesis_cid
726                            ),
727                            |p| get_user_agent(peer_info_map, p),
728                        )
729                        .await;
730                } else {
731                    let sent = SystemTime::now()
732                        .duration_since(UNIX_EPOCH)
733                        .expect("System time before unix epoch")
734                        .as_nanos()
735                        .try_into()
736                        .expect("System time since unix epoch should not exceed u64");
737
738                    // Send hello response immediately, no need to have the overhead of emitting
739                    // channel and polling future here.
740                    if let Err(e) = hello.send_response(channel, HelloResponse { arrival, sent }) {
741                        warn!("Failed to send HelloResponse: {e:?}");
742                    } else {
743                        emit_event(
744                            network_sender_out,
745                            NetworkEvent::HelloResponseOutbound {
746                                source: peer,
747                                request,
748                            },
749                        )
750                        .await;
751                    }
752                }
753            }
754            request_response::Message::Response {
755                request_id,
756                response,
757            } => {
758                emit_event(network_sender_out, NetworkEvent::HelloResponseInbound).await;
759                hello.handle_response(&request_id, response).await;
760            }
761        },
762        request_response::Event::OutboundFailure {
763            request_id,
764            peer,
765            error,
766            ..
767        } => {
768            hello.on_outbound_failure(&request_id);
769            match error {
770                request_response::OutboundFailure::UnsupportedProtocols => {
771                    peer_manager
772                        .ban_peer_with_default_duration(peer, "Hello protocol unsupported", |p| {
773                            get_user_agent(peer_info_map, p)
774                        })
775                        .await;
776                }
777                _ => {
778                    peer_manager.mark_peer_bad(peer, format!("Hello outbound failure {error}"));
779                }
780            }
781        }
782        request_response::Event::InboundFailure { .. } => {}
783        request_response::Event::ResponseSent { .. } => (),
784    }
785}
786
787async fn handle_ping_event(ping_event: ping::Event) {
788    match ping_event.result {
789        Ok(rtt) => {
790            trace!(
791                "PingSuccess::Ping rtt to {} is {} ms",
792                ping_event.peer,
793                rtt.as_millis()
794            );
795        }
796        Err(ping::Failure::Unsupported) => {
797            debug!(peer=%ping_event.peer, "Ping protocol unsupported");
798        }
799        Err(ping::Failure::Timeout) => {
800            debug!("Ping timeout: {}", ping_event.peer);
801        }
802        Err(ping::Failure::Other { error }) => {
803            debug!("Ping failure: {error}");
804        }
805    }
806}
807
808async fn handle_chain_exchange_event(
809    chain_exchange: &mut ChainExchangeBehaviour,
810    ce_event: request_response::Event<ChainExchangeRequest, ChainExchangeResponse>,
811    db: &ChainStore,
812    network_sender_out: &Sender<NetworkEvent>,
813    cx_response_tx: Sender<(
814        request_response::InboundRequestId,
815        request_response::ResponseChannel<ChainExchangeResponse>,
816        ChainExchangeResponse,
817    )>,
818) {
819    const CHAIN_EXCHANGE_RESPONSE_TIMEOUT: Duration = Duration::from_mins(5);
820    match ce_event {
821        request_response::Event::Message { peer, message, .. } => match message {
822            request_response::Message::Request {
823                request,
824                channel,
825                request_id,
826            } => {
827                let Some(per_peer_permit) = chain_exchange.try_acquire_peer_permit(peer) else {
828                    debug!("Rejecting chain_exchange request from {peer}: per-peer cap reached");
829                    let _ = chain_exchange.send_response(
830                        channel,
831                        ChainExchangeResponse::go_away("per-peer concurrent request cap reached"),
832                    );
833                    return;
834                };
835                let Some(global_permit) = chain_exchange.try_acquire_request_permit() else {
836                    debug!("Rejecting chain_exchange request from {peer}: global cap reached");
837                    let _ = chain_exchange.send_response(
838                        channel,
839                        ChainExchangeResponse::go_away("global concurrent request cap reached"),
840                    );
841                    return;
842                };
843
844                trace!(
845                    "Received chain_exchange request (request_id:{request_id}, peer_id: {peer:?})",
846                );
847                emit_event(
848                    network_sender_out,
849                    NetworkEvent::ChainExchangeRequestInbound,
850                )
851                .await;
852
853                let db = db.shallow_clone();
854                tokio::task::spawn(tokio::time::timeout(
855                    CHAIN_EXCHANGE_RESPONSE_TIMEOUT,
856                    async move {
857                        let _per_peer_permit = per_peer_permit;
858                        let _global_permit = global_permit;
859                        if let Err(e) = cx_response_tx
860                            .send_async((
861                                request_id,
862                                channel,
863                                make_chain_exchange_response(&db, &request),
864                            ))
865                            .await
866                        {
867                            debug!("Failed to send ChainExchangeResponse: {e:?}");
868                        }
869                    },
870                ));
871            }
872            request_response::Message::Response {
873                request_id,
874                response,
875            } => {
876                emit_event(
877                    network_sender_out,
878                    NetworkEvent::ChainExchangeResponseInbound,
879                )
880                .await;
881                chain_exchange
882                    .handle_inbound_response(&request_id, response)
883                    .await;
884            }
885        },
886        request_response::Event::OutboundFailure {
887            request_id, error, ..
888        } => {
889            chain_exchange.on_outbound_error(&request_id, error);
890        }
891        request_response::Event::InboundFailure { peer, error, .. } => {
892            debug!(
893                "ChainExchange inbound error (peer: {:?}): {:?}",
894                peer, error
895            );
896        }
897        request_response::Event::ResponseSent { .. } => {
898            emit_event(
899                network_sender_out,
900                NetworkEvent::ChainExchangeResponseOutbound,
901            )
902            .await;
903        }
904    }
905}
906
907#[allow(clippy::too_many_arguments)]
908async fn handle_forest_behaviour_event(
909    swarm: &mut Swarm<ForestBehaviour>,
910    bitswap_request_manager: &Arc<BitswapRequestManager>,
911    peer_manager: &PeerManager,
912    event: ForestBehaviourEvent,
913    db: &ChainStore,
914    genesis_cid: &Cid,
915    network_sender_out: &Sender<NetworkEvent>,
916    cx_response_tx: Sender<(
917        request_response::InboundRequestId,
918        request_response::ResponseChannel<ChainExchangeResponse>,
919        ChainExchangeResponse,
920    )>,
921    pubsub_block_str: &str,
922    pubsub_msg_str: &str,
923) {
924    match event {
925        ForestBehaviourEvent::Discovery(discovery_out) => {
926            handle_discovery_event(
927                &swarm.behaviour().discovery.peer_info,
928                discovery_out,
929                network_sender_out,
930                peer_manager,
931            )
932            .await
933        }
934        ForestBehaviourEvent::Gossipsub(e) => {
935            handle_gossip_event(e, network_sender_out, pubsub_block_str, pubsub_msg_str).await
936        }
937        ForestBehaviourEvent::Hello(rr_event) => {
938            let behaviour_mut = swarm.behaviour_mut();
939            handle_hello_event(
940                &behaviour_mut.discovery.peer_info,
941                &mut behaviour_mut.hello,
942                rr_event,
943                peer_manager,
944                genesis_cid,
945                network_sender_out,
946            )
947            .await
948        }
949        ForestBehaviourEvent::Bitswap(event) => {
950            if let Err(e) = bitswap_request_manager.handle_event(
951                &mut swarm.behaviour_mut().bitswap,
952                db.db(),
953                event,
954            ) {
955                warn!("bitswap: {e:#}");
956            }
957        }
958        ForestBehaviourEvent::Ping(ping_event) => handle_ping_event(ping_event).await,
959        ForestBehaviourEvent::ConnectionLimits(_) => {}
960        ForestBehaviourEvent::BlockedPeers(_) => {}
961        ForestBehaviourEvent::ChainExchange(ce_event) => {
962            handle_chain_exchange_event(
963                &mut swarm.behaviour_mut().chain_exchange,
964                ce_event,
965                db,
966                network_sender_out,
967                cx_response_tx,
968            )
969            .await
970        }
971    }
972}
973
974async fn emit_event(sender: &Sender<NetworkEvent>, event: NetworkEvent) {
975    if sender.send_async(event).await.is_err() {
976        error!("Failed to emit event: Network channel receiver has been dropped");
977    }
978}
979
980fn get_user_agent(peer_info_map: &HashMap<PeerId, PeerInfo>, peer: &PeerId) -> Option<String> {
981    peer_info_map
982        .get(peer)
983        .and_then(|i| i.identify_info.as_ref())
984        .map(|i| i.agent_version.clone())
985}