1use 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
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 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 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 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 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 pub fn network_sender(&self) -> Sender<NetworkMessage> {
397 self.network_sender_in.clone()
398 }
399
400 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 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 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 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 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}