Skip to main content

sui_network/discovery/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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
58/// Per-source P2P addresses for trusted peers, keyed by peer then by address source.
59/// Returned by [`Sender::trusted_peer_p2p_addresses`].
60pub type TrustedPeerP2pAddresses =
61    BTreeMap<PeerId, BTreeMap<AddressSource, Vec<anemo::types::Address>>>;
62
63/// Message types for the discovery system mailbox.
64#[derive(Debug)]
65pub enum DiscoveryMessage {
66    /// An external source (e.g. admin API, node config) updated a peer's address.
67    PeerAddressChange {
68        peer_id: PeerId,
69        source: AddressSource,
70        addresses: Vec<anemo::types::Address>,
71    },
72    /// Node info was received from inbound RPC.
73    ReceivedNodeInfo {
74        peer_info: Box<SignedVersionedNodeInfo>,
75    },
76    /// A spawned peer-query task discovered updated info for a trusted peer,
77    /// signaling the event loop to update the stored peer addresses.
78    TrustedPeersUpdated,
79    /// A peer has been reported as continuously failing, triggering disconnect and cooldown.
80    PeerFailureReport { peer_id: PeerId },
81    /// Request a snapshot of per-source P2P addresses, for trusted peers with at least one
82    /// recorded address.
83    GetTrustedPeerP2pAddresses {
84        reply: oneshot::Sender<TrustedPeerP2pAddresses>,
85    },
86}
87
88/// A Handle to the Discovery subsystem. The Discovery system will be shut down once all Handles
89/// have been dropped.
90#[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/// A lightweight handle for sending messages to the discovery event loop
103/// without holding a shutdown reference.
104#[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    /// Snapshot of per-source P2P addresses for trusted peers only (configured/seed/allowlisted
132    /// peers and on-chain validators). Used by the address prober to enumerate probe candidates;
133    /// returns an empty map if the discovery event loop has shut down.
134    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
150/// The internal discovery state shared between the main event loop and the request handler
151struct 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/// The information necessary to dial another peer.
161///
162/// `NodeInfo` contains all the information that is shared with other nodes via the discovery
163/// service to advertise how a node can be reached.
164#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
165pub struct NodeInfo {
166    pub peer_id: PeerId,
167    pub addresses: Vec<Multiaddr>,
168
169    /// Creation time.
170    ///
171    /// This is used to determine which of two NodeInfo's from the same PeerId should be retained.
172    pub timestamp_ms: u64,
173
174    /// See docstring for `AccessType`.
175    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/// NodeInfoV2 supports multiple address types keyed by EndpointId.
209// TODO: Remove support for V1 once V2 is available in all production networks.
210#[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    /// Derive the P2P PeerId from the addresses map.
219    /// Returns None if no P2P endpoint is present.
220    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/// Versioned wrapper for NodeInfo types.
239#[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
296/// Verifies the signature and endpoint identity consistency of a
297/// `SignedVersionedNodeInfo`, returning a `VerifiedSignedVersionedNodeInfo`
298/// on success.
299///
300/// Checks:
301/// 1. The signature is valid for the P2P peer_id embedded in the node info.
302/// 2. Any non-P2P endpoint identities (e.g. `Consensus`) match the P2P
303///    peer_id, preventing a peer from advertising endpoints under another
304///    key.
305fn 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            // Each endpoint variant (P2p, Consensus) may appear at most once.
333            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                                // avoid crashing on ungraceful shutdown
425                            } else if e.is_panic() {
426                                // propagate panics.
427                                std::panic::resume_unwind(e.into_panic());
428                            } else {
429                                panic!("task failed: {e}");
430                            }
431                        },
432                    };
433                },
434                // Once the shutdown notification resolves we can terminate the event loop
435                _ = &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            // Populates `Consensus` EndpointId from our P2P (anemo) `PeerId`.
541            // This is safe because both the P2P and consensus networks use the same
542            // ed25519 network keypair. Both originate from `NodeConfig::network_key_pair()`.
543            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                    // Allowed peers accept inbound but aren't actively dialed,
576                    // so insert them directly without address priority tracking.
577                    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        // Update set of trusted peers from on-chain validator configs.
616        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        // Update stored addresses.
625        {
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            // When a Chain address arrives, check if there are cached Discovery
639            // P2P addresses in known_peers_v2 that should take priority. This
640            // handles the startup race where cached NodeInfo is loaded before
641            // chain_peers is populated (so update_known_peers_versioned couldn't
642            // forward the Discovery addresses at load time).
643            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        // Check if we should use the updated addresses.
671        self.reconfigure_peer_addresses(peer_id);
672    }
673
674    /// Reads the highest-priority addresses from `peer_addresses` for the given peer
675    /// and updates `network.known_peers()` if they differ from the current addresses.
676    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        // Record the active P2P address source for trusted peers (one gauge per
690        // peer, value-encoded). Untrusted peers are skipped to bound metric
691        // cardinality. A `None` priority (no installed source) records 0.
692        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            // Override existing connection if there might be one.
716            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                        // If this fails, ConnectionManager will retry.
723                        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        // update_known_peers_versioned forwards addresses through the mailbox.
753        // Drain it now so addresses are applied before the event loop starts.
754        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                    // Query the new node for any peers
785                    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        // Cull old peers older than a day
825        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        // Clear Discovery-sourced addresses for trusted peers whose
841        // signed info expired, allowing fallback to lower-priority sources.
842        for peer_id in culled_trusted_peers {
843            self.endpoint_manager
844                .clear_source(peer_id, AddressSource::Discovery);
845        }
846
847        // Clean out the pending_dials
848        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        // Spawn some dials
860        let state = self.state.read().unwrap();
861        let our_peer_id = self.network.peer_id();
862
863        // Recompute the number of distinct peers with an external address from the
864        // deduped union of both known-peer maps. With v3 enabled a peer can appear
865        // in both maps, so a HashSet avoids double-counting; recomputing each tick
866        // (rather than inc/dec) also avoids drift when expired entries are culled.
867        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        // Collect eligible peers from both known_peers (V2) and known_peers_v2 (V3),
882        // preferring fresher timestamps when a peer appears in both maps.
883        // Partition into preferred (not on cooldown) and cooldown peers.
884        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        // If we aren't connected to anything and we aren't presently trying to connect to anyone
959        // we need to try the configured peers with High affinity (seed peers)
960        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            // Ignore the result and just log the error if there is one
996            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    // Attempt connection to all high-affinity and seed peers.
1033    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                // Ignore the result and just log the error if there is one
1045                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    // Query V3 concurrently with V2 when enabled
1098    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    // V3 not enabled or our_info_v2 not available - just run V2
1152    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    // Query V2 to keep known_peers populated for V2 clients
1175    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    // Additionally query V3 when enabled, concurrently with V2
1211    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    // V3 not enabled or our_info_v2 not available - just run V2
1275    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    // only take the first MAX_PEERS_TO_SEND peers
1291    for peer_info in found_peers.into_iter().take(MAX_PEERS_TO_SEND + 1) {
1292        // +1 to account for the "own_info" of the serving peer
1293        // Skip peers whose timestamp is too far in the future from our clock
1294        // or that are too old
1295        if peer_info.timestamp_ms > now_unix.saturating_add(30 * 1_000) // 30 seconds
1296            || 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        // Skip entries that have too many addresses as a means to cap the size of a node's info
1314        if peer_info.addresses.len() > MAX_ADDRESSES_PER_PEER {
1315            continue;
1316        }
1317
1318        // verify that all addresses provided are valid anemo addresses
1319        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                // This should never happen.
1329                "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            // TODO: consider denylisting the source of bad NodeInfo from future requests.
1341            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
1358/// Returns true if any trusted peer was inserted or updated.
1359fn 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    // Filter and verify before taking the state write lock: signature
1377    // verification is CPU-heavy (up to MAX_PEERS_TO_SEND peers per call).
1378    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        // Collect discovered addresses of trusted peers to forward to the
1416        // EndpointManager after the merge below.
1417        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    // Merge under a brief write lock; the guard must drop before the
1429    // forwarding below.
1430    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    // Forward after the merge: handle_peer_address_change's Chain-source
1454    // fallback reads known_peers_v2, so this batch must be visible there
1455    // before these updates are processed. update_endpoint also acquires locks
1456    // in other subsystems, so it must not run under the state lock.
1457    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
1464/// A trusted peer is one that appears in the static configured_peers (seed/allowlisted)
1465/// or was dynamically added via on-chain validator info (chain_peers).
1466pub(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}