1use anemo::types::PeerAffinity;
5use anemo::types::PeerInfo;
6use anemo::{Network, Peer, PeerId, Request, Response, types::PeerEvent};
7use fastcrypto::ed25519::{Ed25519PublicKey, Ed25519Signature};
8use futures::StreamExt;
9use mysten_common::debug_fatal;
10use serde::{Deserialize, Serialize};
11use shared_crypto::intent::IntentScope;
12use std::{
13 collections::{BTreeMap, HashMap, HashSet},
14 path::PathBuf,
15 sync::{Arc, RwLock},
16 time::{Duration, Instant},
17};
18
19use crate::endpoint_manager::{AddressSource, EndpointId, EndpointManager};
20use store::{load_stored_peers, save_stored_peers};
21use sui_config::p2p::{AccessType, DiscoveryConfig, P2pConfig};
22use sui_types::crypto::{NetworkKeyPair, NetworkPublicKey, Signer, ToFromBytes, VerifyingKey};
23use sui_types::digests::Digest;
24use sui_types::message_envelope::{Envelope, Message, VerifiedEnvelope};
25use sui_types::multiaddr::Multiaddr;
26use tap::{Pipe, TapFallible};
27use tokio::sync::broadcast::error::RecvError;
28use tokio::sync::mpsc;
29use tokio::{
30 sync::oneshot,
31 task::{AbortHandle, JoinSet},
32};
33use tracing::{debug, info, trace};
34
35const TIMEOUT: Duration = Duration::from_secs(1);
36const ONE_DAY_MILLISECONDS: u64 = 24 * 60 * 60 * 1_000;
37const MAX_ADDRESS_LENGTH: usize = 300;
38const MAX_PEERS_TO_SEND: usize = 200;
39const MAX_ADDRESSES_PER_PEER: usize = 2;
40
41mod generated {
42 include!(concat!(env!("OUT_DIR"), "/sui.Discovery.rs"));
43}
44mod builder;
45mod metrics;
46mod server;
47mod store;
48#[cfg(test)]
49mod tests;
50
51pub use builder::{Builder, UnstartedDiscovery};
52pub use generated::{
53 discovery_client::DiscoveryClient,
54 discovery_server::{Discovery, DiscoveryServer},
55};
56pub use server::{GetKnownPeersRequestV3, GetKnownPeersResponseV2, GetKnownPeersResponseV3};
57
58pub type TrustedPeerP2pAddresses =
61 BTreeMap<PeerId, BTreeMap<AddressSource, Vec<anemo::types::Address>>>;
62
63#[derive(Debug)]
65pub enum DiscoveryMessage {
66 PeerAddressChange {
68 peer_id: PeerId,
69 source: AddressSource,
70 addresses: Vec<anemo::types::Address>,
71 },
72 ReceivedNodeInfo {
74 peer_info: Box<SignedVersionedNodeInfo>,
75 },
76 TrustedPeersUpdated,
79 PeerFailureReport { peer_id: PeerId },
81 GetTrustedPeerP2pAddresses {
84 reply: oneshot::Sender<TrustedPeerP2pAddresses>,
85 },
86}
87
88#[derive(Clone, Debug)]
91pub struct Handle {
92 pub(super) _shutdown_handle: Arc<oneshot::Sender<()>>,
93 pub(super) sender: Sender,
94}
95
96impl Handle {
97 pub fn sender(&self) -> Sender {
98 self.sender.clone()
99 }
100}
101
102#[derive(Clone, Debug)]
105pub struct Sender {
106 pub(super) sender: mpsc::Sender<DiscoveryMessage>,
107}
108
109impl Sender {
110 pub fn peer_address_change(
111 &self,
112 peer_id: PeerId,
113 source: AddressSource,
114 addresses: Vec<anemo::types::Address>,
115 ) {
116 self.sender
117 .try_send(DiscoveryMessage::PeerAddressChange {
118 peer_id,
119 source,
120 addresses,
121 })
122 .expect("Discovery mailbox should not overflow or be closed")
123 }
124
125 pub fn report_peer_failure(&self, peer_id: PeerId) {
126 let _ = self
127 .sender
128 .try_send(DiscoveryMessage::PeerFailureReport { peer_id });
129 }
130
131 pub async fn trusted_peer_p2p_addresses(&self) -> TrustedPeerP2pAddresses {
135 let (s, r) = oneshot::channel();
136 if self
137 .sender
138 .send(DiscoveryMessage::GetTrustedPeerP2pAddresses { reply: s })
139 .await
140 .is_err()
141 {
142 return TrustedPeerP2pAddresses::default();
143 }
144 r.await.unwrap_or_default()
145 }
146}
147
148use self::metrics::Metrics;
149
150struct State {
152 our_info: Option<SignedNodeInfo>,
153 our_info_v2: Option<SignedVersionedNodeInfo>,
154 connected_peers: HashMap<PeerId, ()>,
155 known_peers: HashMap<PeerId, VerifiedSignedNodeInfo>,
156 known_peers_v2: HashMap<PeerId, VerifiedSignedVersionedNodeInfo>,
157 peer_addresses: HashMap<PeerId, BTreeMap<AddressSource, Vec<anemo::types::Address>>>,
158}
159
160#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
165pub struct NodeInfo {
166 pub peer_id: PeerId,
167 pub addresses: Vec<Multiaddr>,
168
169 pub timestamp_ms: u64,
173
174 pub access_type: AccessType,
176}
177
178impl NodeInfo {
179 fn sign(self, keypair: &NetworkKeyPair) -> SignedNodeInfo {
180 let msg = bcs::to_bytes(&self).expect("BCS serialization should not fail");
181 let sig = keypair.sign(&msg);
182 SignedNodeInfo::new_from_data_and_sig(self, sig)
183 }
184}
185
186pub type SignedNodeInfo = Envelope<NodeInfo, Ed25519Signature>;
187
188pub type VerifiedSignedNodeInfo = VerifiedEnvelope<NodeInfo, Ed25519Signature>;
189
190#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
191pub struct NodeInfoDigest(Digest);
192
193impl NodeInfoDigest {
194 pub const fn new(digest: [u8; 32]) -> Self {
195 Self(Digest::new(digest))
196 }
197}
198
199impl Message for NodeInfo {
200 type DigestType = NodeInfoDigest;
201 const SCOPE: IntentScope = IntentScope::DiscoveryPeers;
202
203 fn digest(&self) -> Self::DigestType {
204 unreachable!("NodeInfoDigest is not used today")
205 }
206}
207
208#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
211pub struct NodeInfoV2 {
212 pub addresses: BTreeMap<EndpointId, Vec<Multiaddr>>,
213 pub timestamp_ms: u64,
214 pub access_type: AccessType,
215}
216
217impl NodeInfoV2 {
218 pub fn peer_id(&self) -> Option<PeerId> {
221 self.addresses.keys().find_map(|k| match k {
222 EndpointId::P2p(peer_id) => Some(*peer_id),
223 EndpointId::Consensus(_) => None,
224 })
225 }
226
227 pub fn p2p_addresses(&self) -> &[Multiaddr] {
228 self.addresses
229 .iter()
230 .find_map(|(k, v)| match k {
231 EndpointId::P2p(_) => Some(v.as_slice()),
232 EndpointId::Consensus(_) => None,
233 })
234 .unwrap_or(&[])
235 }
236}
237
238#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
240pub enum VersionedNodeInfo {
241 V1(NodeInfo),
242 V2(NodeInfoV2),
243}
244
245impl VersionedNodeInfo {
246 pub fn peer_id(&self) -> Option<PeerId> {
247 match self {
248 VersionedNodeInfo::V1(info) => Some(info.peer_id),
249 VersionedNodeInfo::V2(info) => info.peer_id(),
250 }
251 }
252
253 pub fn timestamp_ms(&self) -> u64 {
254 match self {
255 VersionedNodeInfo::V1(info) => info.timestamp_ms,
256 VersionedNodeInfo::V2(info) => info.timestamp_ms,
257 }
258 }
259
260 pub fn access_type(&self) -> AccessType {
261 match self {
262 VersionedNodeInfo::V1(info) => info.access_type,
263 VersionedNodeInfo::V2(info) => info.access_type,
264 }
265 }
266
267 pub fn p2p_addresses(&self) -> &[Multiaddr] {
268 match self {
269 VersionedNodeInfo::V1(info) => &info.addresses,
270 VersionedNodeInfo::V2(info) => info.p2p_addresses(),
271 }
272 }
273
274 pub fn sign(self, keypair: &NetworkKeyPair) -> SignedVersionedNodeInfo {
275 let msg = bcs::to_bytes(&self).expect("BCS serialization should not fail");
276 let sig = keypair.sign(&msg);
277 SignedVersionedNodeInfo::new_from_data_and_sig(self, sig)
278 }
279}
280
281#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
282pub struct VersionedNodeInfoDigest(Digest);
283
284impl Message for VersionedNodeInfo {
285 type DigestType = VersionedNodeInfoDigest;
286 const SCOPE: IntentScope = IntentScope::DiscoveryPeers;
287
288 fn digest(&self) -> Self::DigestType {
289 unreachable!("VersionedNodeInfoDigest is not used today")
290 }
291}
292
293pub type SignedVersionedNodeInfo = Envelope<VersionedNodeInfo, Ed25519Signature>;
294pub type VerifiedSignedVersionedNodeInfo = VerifiedEnvelope<VersionedNodeInfo, Ed25519Signature>;
295
296fn verify_versioned_node_info(
306 peer_info: &SignedVersionedNodeInfo,
307) -> Result<VerifiedSignedVersionedNodeInfo, &'static str> {
308 let peer_id = peer_info.peer_id().ok_or("missing P2P peer_id")?;
309
310 let public_key =
311 Ed25519PublicKey::from_bytes(&peer_id.0).map_err(|_| "invalid peer_id public key")?;
312
313 let msg = bcs::to_bytes(peer_info.data()).expect("BCS serialization should not fail");
314 public_key
315 .verify(&msg, peer_info.auth_sig())
316 .map_err(|_| "signature verification failed")?;
317
318 match peer_info.data() {
319 VersionedNodeInfo::V1(info) => {
320 if info.addresses.len() > MAX_ADDRESSES_PER_PEER {
321 return Err("too many addresses");
322 }
323 if !info
324 .addresses
325 .iter()
326 .all(|addr| addr.len() < MAX_ADDRESS_LENGTH && addr.to_anemo_address().is_ok())
327 {
328 return Err("invalid address");
329 }
330 }
331 VersionedNodeInfo::V2(info_v2) => {
332 let mut seen_variants = Vec::new();
334 for endpoint_id in info_v2.addresses.keys() {
335 let variant = std::mem::discriminant(endpoint_id);
336 if seen_variants.contains(&variant) {
337 return Err("duplicate endpoint variant");
338 }
339 seen_variants.push(variant);
340 }
341
342 for (endpoint_id, addrs) in &info_v2.addresses {
343 if addrs.len() > MAX_ADDRESSES_PER_PEER {
344 return Err("too many addresses for endpoint");
345 }
346 if !addrs.iter().all(|addr| addr.len() < MAX_ADDRESS_LENGTH) {
347 return Err("address too long");
348 }
349 if matches!(endpoint_id, EndpointId::P2p(_))
350 && !addrs.iter().all(|addr| addr.to_anemo_address().is_ok())
351 {
352 return Err("invalid P2P address");
353 }
354 }
355
356 let identities_valid = info_v2.addresses.keys().all(|eid| match eid {
357 EndpointId::P2p(_) => true,
358 EndpointId::Consensus(pubkey) => pubkey.as_bytes() == peer_id.0,
359 });
360 if !identities_valid {
361 return Err("non-P2P endpoint identity mismatch");
362 }
363 }
364 }
365
366 Ok(VerifiedSignedVersionedNodeInfo::new_from_verified(
367 peer_info.clone(),
368 ))
369}
370
371struct DiscoveryEventLoop {
372 config: P2pConfig,
373 discovery_config: Arc<DiscoveryConfig>,
374 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
375 chain_peers: Arc<RwLock<HashSet<PeerId>>>,
376 unidentified_seed_peers: Vec<anemo::types::Address>,
377 network: Network,
378 keypair: NetworkKeyPair,
379 tasks: JoinSet<()>,
380 pending_dials: HashMap<PeerId, AbortHandle>,
381 dial_seed_peers_task: Option<AbortHandle>,
382 shutdown_handle: oneshot::Receiver<()>,
383 state: Arc<RwLock<State>>,
384 mailbox: mpsc::Receiver<DiscoveryMessage>,
385 mailbox_tx: mpsc::Sender<DiscoveryMessage>,
386 metrics: Metrics,
387 consensus_external_address: Option<Multiaddr>,
388 endpoint_manager: EndpointManager,
389 store_path: Option<PathBuf>,
390 peer_cooldowns: HashMap<PeerId, Instant>,
391}
392
393impl DiscoveryEventLoop {
394 pub async fn start(mut self) {
395 info!("Discovery started");
396
397 self.construct_our_info();
398 self.configure_preferred_peers();
399 self.load_stored_peers_on_startup();
400
401 let mut interval = tokio::time::interval(self.discovery_config.interval_period());
402 let mut peer_events = {
403 let (subscriber, _peers) = self.network.subscribe().unwrap();
404 subscriber
405 };
406
407 loop {
408 tokio::select! {
409 now = interval.tick() => {
410 let now_unix = now_unix();
411 self.handle_tick(now.into_std(), now_unix);
412 }
413 peer_event = peer_events.recv() => {
414 self.handle_peer_event(peer_event);
415 },
416 Some(message) = self.mailbox.recv() => {
417 self.handle_message(message);
418 }
419 Some(task_result) = self.tasks.join_next() => {
420 match task_result {
421 Ok(()) => {},
422 Err(e) => {
423 if e.is_cancelled() {
424 } else if e.is_panic() {
426 std::panic::resume_unwind(e.into_panic());
428 } else {
429 panic!("task failed: {e}");
430 }
431 },
432 };
433 },
434 _ = &mut self.shutdown_handle => {
436 break;
437 }
438 }
439 }
440
441 self.save_stored_peers();
442 info!("Discovery ended");
443 }
444
445 fn handle_message(&mut self, message: DiscoveryMessage) {
446 match message {
447 DiscoveryMessage::PeerAddressChange {
448 peer_id,
449 source,
450 addresses,
451 } => {
452 self.handle_peer_address_change(peer_id, source, addresses);
453 }
454 DiscoveryMessage::ReceivedNodeInfo { peer_info } => {
455 let changed = update_known_peers_versioned(
456 self.state.clone(),
457 vec![*peer_info],
458 self.configured_peers.clone(),
459 &self.chain_peers,
460 &self.endpoint_manager,
461 );
462 if changed {
463 self.save_stored_peers();
464 }
465 }
466 DiscoveryMessage::TrustedPeersUpdated => {
467 self.save_stored_peers();
468 }
469 DiscoveryMessage::PeerFailureReport { peer_id } => {
470 self.handle_peer_failure_report(peer_id);
471 }
472 DiscoveryMessage::GetTrustedPeerP2pAddresses { reply } => {
473 let _ = reply.send(self.trusted_peer_p2p_addresses());
474 }
475 }
476 }
477
478 fn trusted_peer_p2p_addresses(&self) -> TrustedPeerP2pAddresses {
479 let state = self.state.read().unwrap();
480 state
481 .peer_addresses
482 .iter()
483 .filter(|(peer_id, _)| {
484 is_trusted_peer(peer_id, &self.configured_peers, &self.chain_peers)
485 })
486 .map(|(peer_id, sources)| (*peer_id, sources.clone()))
487 .collect()
488 }
489
490 fn handle_peer_failure_report(&mut self, peer_id: PeerId) {
491 if self.is_trusted_peer(&peer_id) {
492 info!(?peer_id, "ignoring failure report for trusted peer");
493 return;
494 }
495 let min_peers = self.discovery_config.min_peers_for_disconnect();
496 let connected_count = self.state.read().unwrap().connected_peers.len();
497 if connected_count < min_peers {
498 info!(
499 ?peer_id,
500 connected_count, min_peers, "skipping disconnect, too few connected peers"
501 );
502 return;
503 }
504 info!(
505 ?peer_id,
506 "peer failure reported, disconnecting and adding cooldown"
507 );
508 let _ = self.network.disconnect(peer_id);
509 self.peer_cooldowns.insert(peer_id, Instant::now());
510 }
511
512 fn construct_our_info(&mut self) {
513 if self.state.read().unwrap().our_info.is_some() {
514 return;
515 }
516
517 let peer_id = self.network.peer_id();
518 let timestamp_ms = now_unix();
519 let access_type = self.discovery_config.access_type();
520
521 let addresses: Vec<Multiaddr> = self
522 .config
523 .external_address
524 .clone()
525 .and_then(|addr| addr.to_anemo_address().ok().map(|_| addr))
526 .into_iter()
527 .collect();
528
529 let our_info = NodeInfo {
530 peer_id,
531 addresses: addresses.clone(),
532 timestamp_ms,
533 access_type,
534 }
535 .sign(&self.keypair);
536
537 let mut addresses_map = BTreeMap::new();
538 addresses_map.insert(EndpointId::P2p(peer_id), addresses);
539 if let Some(consensus_addr) = &self.consensus_external_address {
540 let network_pubkey =
544 NetworkPublicKey::from_bytes(&peer_id.0).expect("PeerId is a valid public key");
545 addresses_map.insert(
546 EndpointId::Consensus(network_pubkey),
547 vec![consensus_addr.clone()],
548 );
549 }
550 let our_info_v2 = VersionedNodeInfo::V2(NodeInfoV2 {
551 addresses: addresses_map,
552 timestamp_ms,
553 access_type,
554 })
555 .sign(&self.keypair);
556
557 let mut state = self.state.write().unwrap();
558 state.our_info = Some(our_info);
559 state.our_info_v2 = Some(our_info_v2);
560 }
561
562 fn configure_preferred_peers(&mut self) {
563 let peers: Vec<_> = self.configured_peers.values().cloned().collect();
564 for peer_info in peers {
565 debug!(?peer_info, "Add configured preferred peer");
566 match peer_info.affinity {
567 PeerAffinity::High => {
568 self.handle_peer_address_change(
569 peer_info.peer_id,
570 AddressSource::Seed,
571 peer_info.address,
572 );
573 }
574 _ => {
575 self.network.known_peers().insert(peer_info);
578 }
579 }
580 }
581 }
582
583 fn update_our_info_timestamp(&mut self, now_unix: u64) {
584 let state = &mut self.state.write().unwrap();
585
586 if let Some(our_info) = &state.our_info {
587 let mut data = our_info.data().clone();
588 data.timestamp_ms = now_unix;
589 state.our_info = Some(data.sign(&self.keypair));
590 }
591
592 if let Some(our_info_v2) = &state.our_info_v2 {
593 let mut data = our_info_v2.data().clone();
594 match &mut data {
595 VersionedNodeInfo::V1(info) => info.timestamp_ms = now_unix,
596 VersionedNodeInfo::V2(info) => info.timestamp_ms = now_unix,
597 }
598 state.our_info_v2 = Some(data.sign(&self.keypair));
599 }
600 }
601
602 fn handle_peer_address_change(
603 &mut self,
604 peer_id: PeerId,
605 source: AddressSource,
606 addresses: Vec<anemo::types::Address>,
607 ) {
608 debug!(
609 ?peer_id,
610 ?source,
611 ?addresses,
612 "Received peer address change"
613 );
614
615 if source == AddressSource::Chain {
617 if addresses.is_empty() {
618 self.chain_peers.write().unwrap().remove(&peer_id);
619 } else {
620 self.chain_peers.write().unwrap().insert(peer_id);
621 }
622 }
623
624 {
626 let mut state = self.state.write().unwrap();
627 let source_map = state.peer_addresses.entry(peer_id).or_default();
628
629 if addresses.is_empty() {
630 source_map.remove(&source);
631 if source_map.is_empty() {
632 state.peer_addresses.remove(&peer_id);
633 }
634 } else {
635 source_map.insert(source, addresses);
636 }
637
638 if source == AddressSource::Chain
644 && !state
645 .peer_addresses
646 .get(&peer_id)
647 .is_some_and(|s| s.contains_key(&AddressSource::Discovery))
648 && let Some(addrs) = state
649 .known_peers_v2
650 .get(&peer_id)
651 .and_then(|info| match info.data() {
652 VersionedNodeInfo::V2(v2) => Some(v2.p2p_addresses()),
653 _ => None,
654 })
655 {
656 let anemo_addrs: Vec<_> = addrs
657 .iter()
658 .filter_map(|a| a.to_anemo_address().ok())
659 .collect();
660 if !anemo_addrs.is_empty() {
661 state
662 .peer_addresses
663 .entry(peer_id)
664 .or_default()
665 .insert(AddressSource::Discovery, anemo_addrs);
666 }
667 }
668 }
669
670 self.reconfigure_peer_addresses(peer_id);
672 }
673
674 fn reconfigure_peer_addresses(&mut self, peer_id: PeerId) {
677 let priority = self
678 .state
679 .read()
680 .unwrap()
681 .peer_addresses
682 .get(&peer_id)
683 .and_then(|sources| {
684 sources
685 .first_key_value()
686 .map(|(source, addrs)| (*source, addrs.clone()))
687 });
688
689 if self.is_trusted_peer(&peer_id) {
693 self.metrics.set_active_p2p_address_source(
694 &peer_id.to_string(),
695 priority.as_ref().map(|(source, _)| *source),
696 );
697 }
698
699 let priority_addresses = priority.map(|(_, addrs)| addrs).unwrap_or_default();
700 let current_addresses = self
701 .network
702 .known_peers()
703 .get(&peer_id)
704 .map(|info| info.address.clone())
705 .unwrap_or_default();
706 if priority_addresses != current_addresses {
707 let new_peer_info = PeerInfo {
708 peer_id,
709 affinity: PeerAffinity::High,
710 address: priority_addresses.clone(),
711 };
712
713 self.network.known_peers().insert(new_peer_info);
714
715 if !current_addresses.is_empty() {
717 let _ = self.network.disconnect(peer_id);
718
719 if let Some(address) = priority_addresses.first().cloned() {
720 let network = self.network.clone();
721 self.tasks.spawn(async move {
722 let _ = network.connect_with_peer_id(address, peer_id).await;
724 });
725 }
726 }
727 }
728 }
729
730 fn load_stored_peers_on_startup(&mut self) {
731 let Some(path) = &self.store_path else {
732 return;
733 };
734 let entries = load_stored_peers(path);
735 if entries.is_empty() {
736 return;
737 }
738 info!(
739 count = entries.len(),
740 "Loaded stored peer addresses from {}",
741 path.display()
742 );
743
744 update_known_peers_versioned(
745 self.state.clone(),
746 entries,
747 self.configured_peers.clone(),
748 &self.chain_peers,
749 &self.endpoint_manager,
750 );
751
752 while let Ok(msg) = self.mailbox.try_recv() {
755 self.handle_message(msg);
756 }
757 }
758
759 fn save_stored_peers(&self) {
760 let Some(path) = &self.store_path else {
761 return;
762 };
763 let state = self.state.read().unwrap();
764 let peers_to_save: Vec<SignedVersionedNodeInfo> = state
765 .known_peers_v2
766 .iter()
767 .filter(|(pid, _)| self.is_trusted_peer(pid))
768 .map(|(_, verified)| verified.inner().clone())
769 .collect();
770 drop(state);
771 save_stored_peers(path, &peers_to_save);
772 }
773
774 fn handle_peer_event(&mut self, peer_event: Result<PeerEvent, RecvError>) {
775 match peer_event {
776 Ok(PeerEvent::NewPeer(peer_id)) => {
777 if let Some(peer) = self.network.peer(peer_id) {
778 self.state
779 .write()
780 .unwrap()
781 .connected_peers
782 .insert(peer_id, ());
783
784 self.tasks.spawn(query_peer_for_their_known_peers(
786 peer,
787 self.discovery_config.clone(),
788 self.state.clone(),
789 self.configured_peers.clone(),
790 self.chain_peers.clone(),
791 self.endpoint_manager.clone(),
792 self.mailbox_sender(),
793 ));
794 }
795 }
796 Ok(PeerEvent::LostPeer(peer_id, _)) => {
797 self.state.write().unwrap().connected_peers.remove(&peer_id);
798 }
799
800 Err(RecvError::Closed) => {
801 panic!("PeerEvent channel shouldn't be able to be closed");
802 }
803
804 Err(RecvError::Lagged(_)) => {
805 trace!("State-Sync fell behind processing PeerEvents");
806 }
807 }
808 }
809
810 fn handle_tick(&mut self, _now: std::time::Instant, now_unix: u64) {
811 self.update_our_info_timestamp(now_unix);
812
813 self.tasks
814 .spawn(query_connected_peers_for_their_known_peers(
815 self.network.clone(),
816 self.discovery_config.clone(),
817 self.state.clone(),
818 self.configured_peers.clone(),
819 self.chain_peers.clone(),
820 self.endpoint_manager.clone(),
821 self.mailbox_sender(),
822 ));
823
824 let mut culled_trusted_peers = Vec::new();
826 {
827 let mut state = self.state.write().unwrap();
828 state
829 .known_peers
830 .retain(|_k, v| now_unix.saturating_sub(v.timestamp_ms) < ONE_DAY_MILLISECONDS);
831 state.known_peers_v2.retain(|k, v| {
832 let keep = now_unix.saturating_sub(v.timestamp_ms()) < ONE_DAY_MILLISECONDS;
833 if !keep && self.is_trusted_peer(k) {
834 culled_trusted_peers.push(*k);
835 }
836 keep
837 });
838 }
839
840 for peer_id in culled_trusted_peers {
843 self.endpoint_manager
844 .clear_source(peer_id, AddressSource::Discovery);
845 }
846
847 self.pending_dials.retain(|_k, v| !v.is_finished());
849 if let Some(abort_handle) = &self.dial_seed_peers_task
850 && abort_handle.is_finished()
851 {
852 self.dial_seed_peers_task = None;
853 }
854
855 let cooldown = self.discovery_config.peer_failure_cooldown();
856 self.peer_cooldowns
857 .retain(|_, since| since.elapsed() < cooldown);
858
859 let state = self.state.read().unwrap();
861 let our_peer_id = self.network.peer_id();
862
863 let mut peers_with_external_address: HashSet<PeerId> = HashSet::new();
868 for (peer_id, info) in state.known_peers.iter() {
869 if !info.addresses.is_empty() {
870 peers_with_external_address.insert(*peer_id);
871 }
872 }
873 for (peer_id, info) in state.known_peers_v2.iter() {
874 if !info.p2p_addresses().is_empty() {
875 peers_with_external_address.insert(*peer_id);
876 }
877 }
878 self.metrics
879 .set_num_peers_with_external_address(peers_with_external_address.len() as i64);
880
881 let mut preferred: HashMap<PeerId, NodeInfo> = HashMap::new();
885 let mut cooldown_peers: HashMap<PeerId, NodeInfo> = HashMap::new();
886
887 for (peer_id, info) in state.known_peers.iter() {
888 if *peer_id != our_peer_id
889 && !info.addresses.is_empty()
890 && !state.connected_peers.contains_key(peer_id)
891 && !self.pending_dials.contains_key(peer_id)
892 && !state.peer_addresses.contains_key(peer_id)
893 {
894 if self.peer_cooldowns.contains_key(peer_id) {
895 cooldown_peers.insert(*peer_id, info.data().clone());
896 } else {
897 preferred.insert(*peer_id, info.data().clone());
898 }
899 }
900 }
901 for (peer_id, info) in state.known_peers_v2.iter() {
902 let p2p_addresses = info.p2p_addresses();
903 if *peer_id != our_peer_id
904 && !p2p_addresses.is_empty()
905 && !state.connected_peers.contains_key(peer_id)
906 && !self.pending_dials.contains_key(peer_id)
907 && !state.peer_addresses.contains_key(peer_id)
908 {
909 let node_info = NodeInfo {
910 peer_id: *peer_id,
911 addresses: p2p_addresses.to_vec(),
912 timestamp_ms: info.timestamp_ms(),
913 access_type: info.access_type(),
914 };
915 if self.peer_cooldowns.contains_key(peer_id) {
916 if cooldown_peers
917 .get(peer_id)
918 .is_none_or(|existing| info.timestamp_ms() > existing.timestamp_ms)
919 {
920 cooldown_peers.insert(*peer_id, node_info);
921 }
922 } else if preferred
923 .get(peer_id)
924 .is_none_or(|existing| info.timestamp_ms() > existing.timestamp_ms)
925 {
926 preferred.insert(*peer_id, node_info);
927 }
928 }
929 }
930
931 let number_of_connections = state.connected_peers.len();
932 let number_to_dial = self
933 .discovery_config
934 .target_concurrent_connections()
935 .saturating_sub(number_of_connections);
936
937 use rand::seq::IteratorRandom;
938 let mut rng = rand::thread_rng();
939
940 let mut to_dial: Vec<_> = preferred
941 .into_iter()
942 .choose_multiple(&mut rng, number_to_dial);
943
944 let remaining = number_to_dial.saturating_sub(to_dial.len());
945 to_dial.extend(
946 cooldown_peers
947 .into_iter()
948 .choose_multiple(&mut rng, remaining),
949 );
950
951 for (peer_id, info) in to_dial {
952 let abort_handle = self
953 .tasks
954 .spawn(try_to_connect_to_peer(self.network.clone(), info));
955 self.pending_dials.insert(peer_id, abort_handle);
956 }
957
958 let has_peers_to_dial = || {
961 self.configured_peers
962 .values()
963 .any(|p| p.affinity == PeerAffinity::High)
964 || !self.unidentified_seed_peers.is_empty()
965 };
966 if self.dial_seed_peers_task.is_none()
967 && state.connected_peers.is_empty()
968 && self.pending_dials.is_empty()
969 && has_peers_to_dial()
970 {
971 let abort_handle = self.tasks.spawn(try_to_connect_to_seed_peers(
972 self.network.clone(),
973 self.discovery_config.clone(),
974 self.configured_peers.clone(),
975 self.unidentified_seed_peers.clone(),
976 ));
977
978 self.dial_seed_peers_task = Some(abort_handle);
979 }
980 }
981
982 fn is_trusted_peer(&self, peer_id: &PeerId) -> bool {
983 is_trusted_peer(peer_id, &self.configured_peers, &self.chain_peers)
984 }
985
986 fn mailbox_sender(&self) -> mpsc::Sender<DiscoveryMessage> {
987 self.mailbox_tx.clone()
988 }
989}
990
991async fn try_to_connect_to_peer(network: Network, info: NodeInfo) {
992 debug!("Connecting to peer {info:?}");
993 for multiaddr in &info.addresses {
994 if let Ok(address) = multiaddr.to_anemo_address() {
995 if network
997 .connect_with_peer_id(address, info.peer_id)
998 .await
999 .tap_err(|e| {
1000 debug!(
1001 "error dialing {} at address '{}': {e}",
1002 info.peer_id.short_display(4),
1003 multiaddr
1004 )
1005 })
1006 .is_ok()
1007 {
1008 return;
1009 }
1010 }
1011 }
1012}
1013
1014async fn try_to_connect_to_seed_peers(
1015 network: Network,
1016 config: Arc<DiscoveryConfig>,
1017 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
1018 unidentified_seed_peers: Vec<anemo::types::Address>,
1019) {
1020 let high_affinity_peers: Vec<_> = configured_peers
1021 .values()
1022 .filter(|p| p.affinity == PeerAffinity::High)
1023 .cloned()
1024 .collect();
1025 debug!(
1026 ?high_affinity_peers,
1027 ?unidentified_seed_peers,
1028 "Connecting to seed peers"
1029 );
1030 let network = &network;
1031
1032 let with_peer_id = high_affinity_peers.into_iter().flat_map(|peer_info| {
1034 peer_info
1035 .address
1036 .into_iter()
1037 .map(move |addr| (Some(peer_info.peer_id), addr))
1038 });
1039 let without_peer_id = unidentified_seed_peers.into_iter().map(|addr| (None, addr));
1040 futures::stream::iter(with_peer_id.chain(without_peer_id))
1041 .for_each_concurrent(
1042 config.target_concurrent_connections(),
1043 |(peer_id, address)| async move {
1044 let _ = if let Some(peer_id) = peer_id {
1046 network
1047 .connect_with_peer_id(address.clone(), peer_id)
1048 .await
1049 .tap_err(|e| {
1050 debug!(
1051 "error dialing peer {} at '{}': {e}",
1052 peer_id.short_display(4),
1053 address
1054 )
1055 })
1056 } else {
1057 network
1058 .connect(address.clone())
1059 .await
1060 .tap_err(|e| debug!("error dialing address '{}': {e}", address))
1061 };
1062 },
1063 )
1064 .await;
1065}
1066
1067async fn query_peer_for_known_peers_v2(peer: Peer) -> Option<Vec<SignedNodeInfo>> {
1068 let mut client = DiscoveryClient::new(peer);
1069 let request = Request::new(()).with_timeout(TIMEOUT);
1070 client
1071 .get_known_peers_v2(request)
1072 .await
1073 .ok()
1074 .map(Response::into_inner)
1075 .map(
1076 |GetKnownPeersResponseV2 {
1077 own_info,
1078 mut known_peers,
1079 }| {
1080 if !own_info.addresses.is_empty() {
1081 known_peers.push(own_info)
1082 }
1083 known_peers
1084 },
1085 )
1086}
1087
1088async fn query_peer_for_their_known_peers(
1089 peer: Peer,
1090 discovery_config: Arc<DiscoveryConfig>,
1091 state: Arc<RwLock<State>>,
1092 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
1093 chain_peers: Arc<RwLock<HashSet<PeerId>>>,
1094 endpoint_manager: EndpointManager,
1095 mailbox_tx: mpsc::Sender<DiscoveryMessage>,
1096) {
1097 if discovery_config.use_get_known_peers_v3() {
1099 let our_info_v2 = state.read().unwrap().our_info_v2.clone();
1100 if let Some(own_info) = our_info_v2 {
1101 let peer_for_v3 = peer.clone();
1102 let v3_query = async move {
1103 let mut client = DiscoveryClient::new(peer_for_v3);
1104 let request =
1105 Request::new(GetKnownPeersRequestV3 { own_info }).with_timeout(TIMEOUT);
1106 client
1107 .get_known_peers_v3(request)
1108 .await
1109 .ok()
1110 .map(Response::into_inner)
1111 .map(
1112 |GetKnownPeersResponseV3 {
1113 own_info,
1114 mut known_peers,
1115 }| {
1116 if !own_info.p2p_addresses().is_empty() {
1117 known_peers.push(own_info)
1118 }
1119 known_peers
1120 },
1121 )
1122 };
1123
1124 let (found_peers_v2, found_peers_v3) =
1125 tokio::join!(query_peer_for_known_peers_v2(peer), v3_query);
1126
1127 if let Some(found_peers) = found_peers_v2 {
1128 update_known_peers(
1129 state.clone(),
1130 found_peers,
1131 configured_peers.clone(),
1132 &chain_peers,
1133 );
1134 }
1135 if let Some(found_peers) = found_peers_v3 {
1136 let changed = update_known_peers_versioned(
1137 state,
1138 found_peers,
1139 configured_peers,
1140 &chain_peers,
1141 &endpoint_manager,
1142 );
1143 if changed {
1144 let _ = mailbox_tx.try_send(DiscoveryMessage::TrustedPeersUpdated);
1145 }
1146 }
1147 return;
1148 }
1149 }
1150
1151 if let Some(found_peers) = query_peer_for_known_peers_v2(peer).await {
1153 update_known_peers(state, found_peers, configured_peers, &chain_peers);
1154 }
1155}
1156
1157async fn query_connected_peers_for_their_known_peers(
1158 network: Network,
1159 config: Arc<DiscoveryConfig>,
1160 state: Arc<RwLock<State>>,
1161 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
1162 chain_peers: Arc<RwLock<HashSet<PeerId>>>,
1163 endpoint_manager: EndpointManager,
1164 mailbox_tx: mpsc::Sender<DiscoveryMessage>,
1165) {
1166 use rand::seq::IteratorRandom;
1167
1168 let peers_to_query: Vec<_> = network
1169 .peers()
1170 .into_iter()
1171 .flat_map(|id| network.peer(id))
1172 .choose_multiple(&mut rand::thread_rng(), config.peers_to_query());
1173
1174 let v2_query = {
1176 let peers = peers_to_query.clone();
1177 let peers_to_query_count = config.peers_to_query();
1178 async move {
1179 peers
1180 .into_iter()
1181 .map(DiscoveryClient::new)
1182 .map(|mut client| async move {
1183 let request = Request::new(()).with_timeout(TIMEOUT);
1184 client
1185 .get_known_peers_v2(request)
1186 .await
1187 .ok()
1188 .map(Response::into_inner)
1189 .map(
1190 |GetKnownPeersResponseV2 {
1191 own_info,
1192 mut known_peers,
1193 }| {
1194 if !own_info.addresses.is_empty() {
1195 known_peers.push(own_info)
1196 }
1197 known_peers
1198 },
1199 )
1200 })
1201 .pipe(futures::stream::iter)
1202 .buffer_unordered(peers_to_query_count)
1203 .filter_map(std::future::ready)
1204 .flat_map(futures::stream::iter)
1205 .collect::<Vec<_>>()
1206 .await
1207 }
1208 };
1209
1210 if config.use_get_known_peers_v3() {
1212 let our_info_v2 = state.read().unwrap().our_info_v2.clone();
1213 if let Some(own_info) = our_info_v2 {
1214 let v3_query = {
1215 let peers_to_query_count = config.peers_to_query();
1216 async move {
1217 peers_to_query
1218 .into_iter()
1219 .map(DiscoveryClient::new)
1220 .map(|mut client| {
1221 let own_info = own_info.clone();
1222 async move {
1223 let request = Request::new(GetKnownPeersRequestV3 { own_info })
1224 .with_timeout(TIMEOUT);
1225 client
1226 .get_known_peers_v3(request)
1227 .await
1228 .ok()
1229 .map(Response::into_inner)
1230 .map(
1231 |GetKnownPeersResponseV3 {
1232 own_info,
1233 mut known_peers,
1234 }| {
1235 if !own_info.p2p_addresses().is_empty() {
1236 known_peers.push(own_info)
1237 }
1238 known_peers
1239 },
1240 )
1241 }
1242 })
1243 .pipe(futures::stream::iter)
1244 .buffer_unordered(peers_to_query_count)
1245 .filter_map(std::future::ready)
1246 .flat_map(futures::stream::iter)
1247 .collect::<Vec<_>>()
1248 .await
1249 }
1250 };
1251
1252 let (found_peers_v2, found_peers_v3) = tokio::join!(v2_query, v3_query);
1253
1254 update_known_peers(
1255 state.clone(),
1256 found_peers_v2,
1257 configured_peers.clone(),
1258 &chain_peers,
1259 );
1260 let changed = update_known_peers_versioned(
1261 state,
1262 found_peers_v3,
1263 configured_peers,
1264 &chain_peers,
1265 &endpoint_manager,
1266 );
1267 if changed {
1268 let _ = mailbox_tx.try_send(DiscoveryMessage::TrustedPeersUpdated);
1269 }
1270 return;
1271 }
1272 }
1273
1274 let found_peers_v2 = v2_query.await;
1276 update_known_peers(state, found_peers_v2, configured_peers, &chain_peers);
1277}
1278
1279fn update_known_peers(
1280 state: Arc<RwLock<State>>,
1281 found_peers: Vec<SignedNodeInfo>,
1282 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
1283 chain_peers: &Arc<RwLock<HashSet<PeerId>>>,
1284) {
1285 use std::collections::hash_map::Entry;
1286
1287 let now_unix = now_unix();
1288 let our_peer_id = state.read().unwrap().our_info.clone().unwrap().peer_id;
1289 let known_peers = &mut state.write().unwrap().known_peers;
1290 for peer_info in found_peers.into_iter().take(MAX_PEERS_TO_SEND + 1) {
1292 if peer_info.timestamp_ms > now_unix.saturating_add(30 * 1_000) || now_unix.saturating_sub(peer_info.timestamp_ms) > ONE_DAY_MILLISECONDS
1297 {
1298 continue;
1299 }
1300
1301 if peer_info.peer_id == our_peer_id {
1302 continue;
1303 }
1304
1305 let is_restricted = match peer_info.access_type {
1306 AccessType::Public => false,
1307 AccessType::Private | AccessType::Trusted => true,
1308 };
1309 if is_restricted && !is_trusted_peer(&peer_info.peer_id, &configured_peers, chain_peers) {
1310 continue;
1311 }
1312
1313 if peer_info.addresses.len() > MAX_ADDRESSES_PER_PEER {
1315 continue;
1316 }
1317
1318 if !peer_info
1320 .addresses
1321 .iter()
1322 .all(|addr| addr.len() < MAX_ADDRESS_LENGTH && addr.to_anemo_address().is_ok())
1323 {
1324 continue;
1325 }
1326 let Ok(public_key) = Ed25519PublicKey::from_bytes(&peer_info.peer_id.0) else {
1327 debug_fatal!(
1328 "Failed to convert anemo PeerId {:?} to Ed25519PublicKey",
1330 peer_info.peer_id
1331 );
1332 continue;
1333 };
1334 let msg = bcs::to_bytes(peer_info.data()).expect("BCS serialization should not fail");
1335 if let Err(e) = public_key.verify(&msg, peer_info.auth_sig()) {
1336 info!(
1337 "Discovery failed to verify signature for NodeInfo for peer {:?}: {e:?}",
1338 peer_info.peer_id
1339 );
1340 continue;
1342 }
1343 let peer = VerifiedSignedNodeInfo::new_from_verified(peer_info);
1344
1345 match known_peers.entry(peer.peer_id) {
1346 Entry::Occupied(mut o) => {
1347 if peer.timestamp_ms > o.get().timestamp_ms {
1348 o.insert(peer);
1349 }
1350 }
1351 Entry::Vacant(v) => {
1352 v.insert(peer);
1353 }
1354 }
1355 }
1356}
1357
1358fn update_known_peers_versioned(
1360 state: Arc<RwLock<State>>,
1361 found_peers: Vec<SignedVersionedNodeInfo>,
1362 configured_peers: Arc<HashMap<PeerId, PeerInfo>>,
1363 chain_peers: &Arc<RwLock<HashSet<PeerId>>>,
1364 endpoint_manager: &EndpointManager,
1365) -> bool {
1366 use std::collections::hash_map::Entry;
1367
1368 let now_unix = now_unix();
1369 let our_peer_id = state
1370 .read()
1371 .unwrap()
1372 .our_info_v2
1373 .as_ref()
1374 .and_then(|info| info.peer_id());
1375
1376 let mut verified_peers = Vec::new();
1379 let mut endpoint_updates = Vec::new();
1380
1381 for peer_info in found_peers.into_iter().take(MAX_PEERS_TO_SEND + 1) {
1382 let timestamp_ms = peer_info.timestamp_ms();
1383
1384 if timestamp_ms > now_unix.saturating_add(30 * 1_000)
1385 || now_unix.saturating_sub(timestamp_ms) > ONE_DAY_MILLISECONDS
1386 {
1387 continue;
1388 }
1389
1390 let Some(peer_id) = peer_info.peer_id() else {
1391 continue;
1392 };
1393
1394 if Some(peer_id) == our_peer_id {
1395 continue;
1396 }
1397
1398 let is_restricted = match peer_info.access_type() {
1399 AccessType::Public => false,
1400 AccessType::Private | AccessType::Trusted => true,
1401 };
1402 let is_trusted = is_trusted_peer(&peer_id, &configured_peers, chain_peers);
1403 if is_restricted && !is_trusted {
1404 continue;
1405 }
1406
1407 let peer = match verify_versioned_node_info(&peer_info) {
1408 Ok(verified) => verified,
1409 Err(reason) => {
1410 info!("Discovery rejecting VersionedNodeInfo for peer {peer_id:?}: {reason}");
1411 continue;
1412 }
1413 };
1414
1415 if is_trusted && let VersionedNodeInfo::V2(info_v2) = peer_info.data() {
1418 for (endpoint_id, addrs) in &info_v2.addresses {
1419 if !addrs.is_empty() {
1420 endpoint_updates.push((endpoint_id.clone(), addrs.clone()));
1421 }
1422 }
1423 }
1424
1425 verified_peers.push((peer_id, peer, is_trusted));
1426 }
1427
1428 let mut trusted_peer_changed = false;
1431 {
1432 let known_peers_v2 = &mut state.write().unwrap().known_peers_v2;
1433 for (peer_id, peer, is_trusted) in verified_peers {
1434 match known_peers_v2.entry(peer_id) {
1435 Entry::Occupied(mut o) => {
1436 if peer.timestamp_ms() > o.get().timestamp_ms() {
1437 o.insert(peer);
1438 if is_trusted {
1439 trusted_peer_changed = true;
1440 }
1441 }
1442 }
1443 Entry::Vacant(v) => {
1444 v.insert(peer);
1445 if is_trusted {
1446 trusted_peer_changed = true;
1447 }
1448 }
1449 }
1450 }
1451 }
1452
1453 for (endpoint_id, addrs) in endpoint_updates {
1458 let _ = endpoint_manager.update_endpoint(endpoint_id, AddressSource::Discovery, addrs);
1459 }
1460
1461 trusted_peer_changed
1462}
1463
1464pub(super) fn is_trusted_peer(
1467 peer_id: &PeerId,
1468 configured_peers: &HashMap<PeerId, PeerInfo>,
1469 chain_peers: &Arc<RwLock<HashSet<PeerId>>>,
1470) -> bool {
1471 configured_peers.contains_key(peer_id) || chain_peers.read().unwrap().contains(peer_id)
1472}
1473
1474pub(super) fn now_unix() -> u64 {
1475 use std::time::{SystemTime, UNIX_EPOCH};
1476
1477 SystemTime::now()
1478 .duration_since(UNIX_EPOCH)
1479 .unwrap()
1480 .as_millis() as u64
1481}