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, 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 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                    // outbound events
335                    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                    // Inbound messages
356                    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                    // Print peer count on an interval.
369                    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    /// Returns a sender which allows sending messages to the libp2p service.
405    pub fn network_sender(&self) -> Sender<NetworkMessage> {
406        self.network_sender_in.clone()
407    }
408
409    /// Returns a receiver to listen to network events emitted from the service.
410    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            // Do not ban bootstrap nodes
446            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                    // Safe to ignore, the message has already been published and deduped by gossipsub, so no need to log this as a warning.
473                    // This matches `go-libp2p-pubsub` behavior https://github.com/libp2p/go-libp2p-pubsub/blob/v0.15.0/topic.go#L240-L242
474                    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                                        // try to get a more specific error cause
569                                        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
705/// Saturating: these timestamps only feed peer latency estimates, so a skewed clock should degrade
706/// those rather than take the node down.
707fn 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                    // Send hello response immediately, no need to have the overhead of emitting
744                    // channel and polling future here.
745                    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}