Skip to main content

reth_network/
discovery.rs

1//! Discovery support for the network.
2
3use crate::{
4    cache::LruMap,
5    error::{NetworkError, ServiceKind},
6};
7use enr::Enr;
8use futures::StreamExt;
9use reth_discv4::{DiscoveryUpdate, Discv4, Discv4Config};
10use reth_discv5::{DiscoveredPeer, Discv5};
11use reth_dns_discovery::{
12    DnsDiscoveryConfig, DnsDiscoveryHandle, DnsDiscoveryService, DnsNodeRecordUpdate, DnsResolver,
13};
14use reth_ethereum_forks::{EnrForkIdEntry, ForkId};
15use reth_net_nat::{NatResolver, ResolveNatInterval};
16use reth_network_api::{DiscoveredEvent, DiscoveryEvent};
17use reth_network_peers::{NodeRecord, PeerId};
18use reth_network_types::PeerAddr;
19use secp256k1::SecretKey;
20use std::{
21    collections::VecDeque,
22    net::{IpAddr, SocketAddr},
23    pin::Pin,
24    sync::Arc,
25    task::{ready, Context, Poll},
26    time::Duration,
27};
28use tokio::{net::UdpSocket, sync::mpsc, task::JoinHandle};
29use tokio_stream::{wrappers::ReceiverStream, Stream};
30use tracing::{debug, trace};
31
32/// Default max capacity for cache of discovered peers.
33///
34/// Default is 10 000 peers.
35pub const DEFAULT_MAX_CAPACITY_DISCOVERED_PEERS_CACHE: u32 = 10_000;
36
37/// How often discv5-only nodes refresh their external IP.
38///
39/// Matches the default interval of discv4.
40const RESOLVE_EXTERNAL_IP_INTERVAL: Duration = Duration::from_secs(60 * 5);
41
42/// An abstraction over the configured discovery protocol.
43///
44/// Listens for new discovered nodes and emits events for discovered nodes and their
45/// address.
46#[derive(Debug)]
47pub struct Discovery {
48    /// All nodes discovered via discovery protocol.
49    ///
50    /// These nodes can be ephemeral and are updated via the discovery protocol.
51    discovered_nodes: LruMap<PeerId, PeerAddr>,
52    /// Local ENR of the discovery v4 service (discv5 ENR has same [`PeerId`]).
53    local_enr: NodeRecord,
54    /// Handler to interact with the Discovery v4 service
55    discv4: Option<Discv4>,
56    /// All KAD table updates from the discv4 service.
57    discv4_updates: Option<ReceiverStream<DiscoveryUpdate>>,
58    /// The handle to the spawned discv4 service
59    _discv4_service: Option<JoinHandle<()>>,
60    /// Handler to interact with the Discovery v5 service
61    discv5: Option<Discv5>,
62    /// All KAD table updates from the discv5 service.
63    discv5_updates: Option<ReceiverStream<discv5::Event>>,
64    /// Background task that, in shared-port mode, drains `UnrecognizedFrame`s from discv5 and
65    /// feeds them into the discv4 ingress so packets advance without polling `Discovery`.
66    _discv5_forwarder: Option<JoinHandle<()>>,
67    /// Handler to interact with the DNS discovery service
68    _dns_discovery: Option<DnsDiscoveryHandle>,
69    /// Updates from the DNS discovery service.
70    dns_discovery_updates: Option<ReceiverStream<DnsNodeRecordUpdate>>,
71    /// The handle to the spawned DNS discovery service
72    _dns_disc_service: Option<JoinHandle<()>>,
73    /// Events buffered until polled.
74    queued_events: VecDeque<DiscoveryEvent>,
75    /// List of listeners subscribed to discovery events.
76    discovery_listeners: Vec<mpsc::UnboundedSender<DiscoveryEvent>>,
77    /// Resolves the external IP for discv5 when discv4 is disabled.
78    nat_resolver: Option<ResolveNatInterval>,
79}
80
81impl Discovery {
82    /// Spawns the discovery service.
83    ///
84    /// This will spawn the [`reth_discv4::Discv4Service`] onto a new task and establish a listener
85    /// channel to receive all discovered nodes. When only discv5 is enabled, `nat` periodically
86    /// resolves its advertised external IP without waiting for peer votes. When ENR updates are
87    /// disabled or discv5 checks inbound connectivity, only explicit `extip` or `extaddr` settings
88    /// are applied.
89    pub async fn new(
90        tcp_addr: SocketAddr,
91        discovery_v4_addr: SocketAddr,
92        sk: SecretKey,
93        discv4_config: Option<Discv4Config>,
94        mut discv5_config: Option<reth_discv5::Config>, // contains discv5 listen address
95        dns_discovery_config: Option<DnsDiscoveryConfig>,
96        nat: Option<NatResolver>,
97    ) -> Result<Self, NetworkError> {
98        // setup discv4 with the discovery address and tcp port
99        let local_enr =
100            NodeRecord::from_secret_key(discovery_v4_addr, &sk).with_tcp_port(tcp_addr.port());
101
102        // For IPv6 we set IPV6_V6ONLY=true so an IPv4 sibling socket on the same port doesn't
103        // clash with the IPv6 one (Linux's default of V6ONLY=0 has IPv6 also claim the IPv4
104        // port via mapped addresses), matching how discv5 binds its `DualStack` sockets.
105        let bind_socket = async |addr: SocketAddr| {
106            let result = match addr {
107                SocketAddr::V4(_) => UdpSocket::bind(addr).await,
108                SocketAddr::V6(_) => {
109                    use socket2::{Domain, Protocol, Socket, Type};
110                    (|| {
111                        let socket = Socket::new(Domain::IPV6, Type::DGRAM, Some(Protocol::UDP))?;
112                        socket.set_only_v6(true)?;
113                        socket.set_nonblocking(true)?;
114                        socket.bind(&addr.into())?;
115                        UdpSocket::from_std(socket.into())
116                    })()
117                }
118            };
119            result
120                .map(Arc::new)
121                .map_err(|err| NetworkError::from_io_error(err, ServiceKind::Discovery(addr)))
122        };
123
124        // In shared-port mode, bind the shared socket and start discv4 without its own receive
125        // loop. Unrecognized frames from discv5 will be forwarded to the ingress handler.
126        let (discv4, discv4_updates, _discv4_service, discv4_ingress, shared_socket) =
127            if let Some(config) = discv4_config {
128                if let Some(discv5_config) = &mut discv5_config &&
129                    discv5_config.has_matching_socket(discovery_v4_addr)
130                {
131                    let socket = bind_socket(discovery_v4_addr).await?;
132
133                    let (discv4, mut discv4_service, ingress) = Discv4::bind_shared(
134                        socket.clone(),
135                        local_enr,
136                        sk,
137                        config,
138                    )
139                    .map_err(|err| {
140                        NetworkError::from_io_error(err, ServiceKind::Discovery(discovery_v4_addr))
141                    })?;
142
143                    let discv4_updates = discv4_service.update_stream();
144                    let discv4_service = discv4_service.spawn();
145                    debug!(target:"net", ?discovery_v4_addr, "started discovery v4 (shared port)");
146                    (
147                        Some(discv4),
148                        Some(discv4_updates),
149                        Some(discv4_service),
150                        Some(ingress),
151                        Some(socket),
152                    )
153                } else {
154                    let (discv4, mut discv4_service) =
155                        Discv4::bind(discovery_v4_addr, local_enr, sk, config).await.map_err(
156                            |err| {
157                                NetworkError::from_io_error(
158                                    err,
159                                    ServiceKind::Discovery(discovery_v4_addr),
160                                )
161                            },
162                        )?;
163                    let discv4_updates = discv4_service.update_stream();
164                    // spawn the service
165                    let discv4_service = discv4_service.spawn();
166
167                    debug!(target:"net", ?discovery_v4_addr, "started discovery v4");
168
169                    (Some(discv4), Some(discv4_updates), Some(discv4_service), None, None)
170                }
171            } else {
172                (None, None, None, None, None)
173            };
174
175        // Discv4 resolves its own external IP when enabled.
176        let nat_resolver = if discv4.is_none() &&
177            let Some(config) = &mut discv5_config
178        {
179            let config = config.discv5_config_mut();
180            nat.filter(|nat| match nat {
181                NatResolver::None => false,
182                // Use the address or hostname explicitly configured by the operator.
183                NatResolver::ExternalIp(_) | NatResolver::ExternalAddr(_) => true,
184                // Preserve manual addresses and let discv5 handle reachability checks. Its check
185                // for incoming connections starts when peer votes change the advertised address.
186                _ => config.enr_update && config.auto_nat_listen_duration.is_none(),
187            })
188            .map(|nat| ResolveNatInterval::interval(nat, RESOLVE_EXTERNAL_IP_INTERVAL))
189        } else {
190            None
191        };
192
193        // Start discv5, wiring in the shared socket if in shared-port mode.
194        let (discv5, discv5_updates) = if let Some(mut config) = discv5_config {
195            // Set OS-assigned advertised RLPx ports to the bound listener port.
196            set_bound_rlpx_port_if_unset(&mut config, tcp_addr.port());
197
198            if let Some(socket) = shared_socket {
199                let discv5_cfg = config.discv5_config_mut();
200
201                // The shared socket covers discv4's address family; bind the opposite family
202                // only if discv5 was configured for dual-stack.
203                let (mut ipv4, mut ipv6) = (None, None);
204                if discovery_v4_addr.is_ipv4() {
205                    ipv4 = Some(socket);
206                    if let Some(addr) = reth_discv5::config::ipv6(&discv5_cfg.listen_config) {
207                        ipv6 = Some(bind_socket(SocketAddr::V6(addr)).await?);
208                    }
209                } else {
210                    ipv6 = Some(socket);
211                    if let Some(addr) = reth_discv5::config::ipv4(&discv5_cfg.listen_config) {
212                        ipv4 = Some(bind_socket(SocketAddr::V4(addr)).await?);
213                    }
214                }
215
216                discv5_cfg.listen_config = discv5::ListenConfig::FromSockets { ipv4, ipv6 };
217            }
218
219            let (discv5, discv5_updates) = Discv5::start(&sk, config).await?;
220            debug!(target:"net", discovery_v5_enr=?discv5.local_enr(), "started discovery v5");
221            (Some(discv5), Some(discv5_updates))
222        } else {
223            (None, None)
224        };
225
226        // In shared-port mode, spawn a task that peels `UnrecognizedFrame` events off the discv5
227        // update stream and feeds them into discv4's ingress. Other events are forwarded through
228        // a new channel that `Discovery::poll` reads. This keeps both protocols moving without
229        // requiring the main `Discovery::poll` loop to be driven for packets to be routed.
230        let (discv5_updates, _discv5_forwarder) = match (discv4_ingress, discv5_updates) {
231            (Some(mut ingress), Some(mut updates)) => {
232                let (tx, rx) = mpsc::channel(updates.max_capacity());
233                let handle = tokio::spawn(async move {
234                    while let Some(event) = updates.recv().await {
235                        if let discv5::Event::UnrecognizedFrame(frame) = &event {
236                            ingress.handle_packet(&frame.packet, frame.src_address).await;
237                            continue;
238                        }
239                        if tx.send(event).await.is_err() {
240                            break;
241                        }
242                    }
243                });
244                (Some(ReceiverStream::new(rx)), Some(handle))
245            }
246            (_, updates) => (updates.map(ReceiverStream::new), None),
247        };
248
249        // setup DNS discovery
250        let (_dns_discovery, dns_discovery_updates, _dns_disc_service) =
251            if let Some(dns_config) = dns_discovery_config {
252                let (mut service, dns_disc) = DnsDiscoveryService::new_pair(
253                    Arc::new(DnsResolver::from_system_conf()?),
254                    dns_config,
255                );
256                let dns_discovery_updates = service.node_record_stream();
257                let dns_disc_service = service.spawn();
258                (Some(dns_disc), Some(dns_discovery_updates), Some(dns_disc_service))
259            } else {
260                (None, None, None)
261            };
262
263        Ok(Self {
264            discovery_listeners: Default::default(),
265            local_enr,
266            discv4,
267            discv4_updates,
268            _discv4_service,
269            discv5,
270            discv5_updates,
271            _discv5_forwarder,
272            discovered_nodes: LruMap::new(DEFAULT_MAX_CAPACITY_DISCOVERED_PEERS_CACHE),
273            queued_events: Default::default(),
274            _dns_disc_service,
275            _dns_discovery,
276            dns_discovery_updates,
277            nat_resolver,
278        })
279    }
280
281    /// Registers a listener for receiving [`DiscoveryEvent`] updates.
282    pub(crate) fn add_listener(&mut self, tx: mpsc::UnboundedSender<DiscoveryEvent>) {
283        self.discovery_listeners.push(tx);
284    }
285
286    /// Notifies all registered listeners with the provided `event`.
287    #[inline]
288    fn notify_listeners(&mut self, event: &DiscoveryEvent) {
289        self.discovery_listeners.retain_mut(|listener| listener.send(event.clone()).is_ok());
290    }
291
292    /// Updates the `eth:ForkId` field in discv4/discv5.
293    pub(crate) fn update_fork_id(&self, fork_id: ForkId) {
294        if let Some(discv4) = &self.discv4 {
295            // use forward-compatible forkid entry
296            discv4.set_eip868_rlp(b"eth".to_vec(), EnrForkIdEntry::from(fork_id))
297        }
298        if let Some(discv5) = &self.discv5 {
299            discv5
300                .encode_and_set_eip868_in_local_enr(b"eth".to_vec(), EnrForkIdEntry::from(fork_id))
301        }
302    }
303
304    /// Bans the [`IpAddr`] in the discovery service.
305    pub(crate) fn ban_ip(&self, ip: IpAddr) {
306        if let Some(discv4) = &self.discv4 {
307            discv4.ban_ip(ip)
308        }
309        if let Some(discv5) = &self.discv5 {
310            discv5.ban_ip(ip)
311        }
312    }
313
314    /// Bans the [`PeerId`] and [`IpAddr`] in the discovery service.
315    pub(crate) fn ban(&self, peer_id: PeerId, ip: IpAddr) {
316        if let Some(discv4) = &self.discv4 {
317            discv4.ban(peer_id, ip)
318        }
319        if let Some(discv5) = &self.discv5 {
320            discv5.ban(peer_id, ip)
321        }
322    }
323
324    /// Returns a shared reference to the discv4.
325    pub fn discv4(&self) -> Option<Discv4> {
326        self.discv4.clone()
327    }
328
329    /// Returns the id with which the local node identifies itself in the network
330    pub(crate) const fn local_id(&self) -> PeerId {
331        self.local_enr.id // local discv4 and discv5 have same id, since signed with same secret key
332    }
333
334    /// Add a node to the discv4 table.
335    pub(crate) fn add_discv4_node(&self, node: NodeRecord) {
336        if let Some(discv4) = &self.discv4 {
337            discv4.add_node(node);
338        }
339    }
340
341    /// Returns discv5 handle.
342    pub fn discv5(&self) -> Option<Discv5> {
343        self.discv5.clone()
344    }
345
346    /// Add a node to the discv4 table.
347    pub(crate) fn add_discv5_node(&self, enr: Enr<SecretKey>) -> Result<(), NetworkError> {
348        if let Some(discv5) = &self.discv5 {
349            discv5.add_node(enr).map_err(NetworkError::Discv5Error)?;
350        }
351
352        Ok(())
353    }
354
355    /// Processes an incoming [`NodeRecord`] update from a discovery service
356    fn on_node_record_update(&mut self, record: NodeRecord, fork_id: Option<ForkId>) {
357        let peer_id = record.id;
358        let tcp_addr = record.tcp_addr();
359        if tcp_addr.port() == 0 {
360            // useless peer for p2p
361            return
362        }
363        let udp_addr = record.udp_addr();
364        let addr = PeerAddr::new(tcp_addr, Some(udp_addr));
365        _ =
366            self.discovered_nodes.get_or_insert(peer_id, || {
367                self.queued_events.push_back(DiscoveryEvent::NewNode(
368                    DiscoveredEvent::EventQueued { peer_id, addr, fork_id },
369                ));
370
371                addr
372            })
373    }
374
375    fn on_discv4_update(&mut self, update: DiscoveryUpdate) {
376        match update {
377            DiscoveryUpdate::Added(record) | DiscoveryUpdate::DiscoveredAtCapacity(record) => {
378                self.on_node_record_update(record, None);
379            }
380            DiscoveryUpdate::EnrForkId(node, fork_id) => {
381                self.queued_events.push_back(DiscoveryEvent::EnrForkId(node, fork_id))
382            }
383            DiscoveryUpdate::Removed(peer_id) => {
384                self.discovered_nodes.remove(&peer_id);
385            }
386            DiscoveryUpdate::Batch(updates) => {
387                for update in updates {
388                    self.on_discv4_update(update);
389                }
390            }
391        }
392    }
393
394    pub(crate) fn poll(&mut self, cx: &mut Context<'_>) -> Poll<DiscoveryEvent> {
395        loop {
396            // Drain all buffered events first
397            if let Some(event) = self.queued_events.pop_front() {
398                self.notify_listeners(&event);
399                return Poll::Ready(event)
400            }
401
402            // drain the discv4 update stream
403            while let Some(Poll::Ready(Some(update))) =
404                self.discv4_updates.as_mut().map(|updates| updates.poll_next_unpin(cx))
405            {
406                self.on_discv4_update(update)
407            }
408
409            // drain the discv5 update stream
410            while let Some(Poll::Ready(Some(update))) =
411                self.discv5_updates.as_mut().map(|updates| updates.poll_next_unpin(cx))
412            {
413                if let Some(discv5) = self.discv5.as_mut() &&
414                    let Some(DiscoveredPeer { node_record, fork_id }) =
415                        discv5.on_discv5_update(update)
416                {
417                    self.on_node_record_update(node_record, fork_id);
418                }
419            }
420
421            // drain the dns update stream
422            while let Some(Poll::Ready(Some(update))) =
423                self.dns_discovery_updates.as_mut().map(|updates| updates.poll_next_unpin(cx))
424            {
425                self.add_discv4_node(update.node_record);
426                if let Err(err) = self.add_discv5_node(update.enr) {
427                    trace!(target: "net::discovery",
428                        %err,
429                        "failed adding node discovered by dns to discv5"
430                    );
431                }
432                self.on_node_record_update(update.node_record, update.fork_id);
433            }
434
435            while let Some(Poll::Ready(ip)) =
436                self.nat_resolver.as_mut().map(|resolver| resolver.poll_tick(cx))
437            {
438                if let Some(ip) = ip {
439                    self.on_external_ip(ip);
440                }
441            }
442
443            if self.queued_events.is_empty() {
444                return Poll::Pending
445            }
446        }
447    }
448
449    /// Updates discv5's advertised IP without changing ports or address families.
450    ///
451    /// Discv4 owns its NAT resolution and updates its own node record and ENR, so only
452    /// discv5 needs an update here.
453    fn on_external_ip(&self, external_ip: IpAddr) {
454        let Some(discv5) = &self.discv5 else { return };
455        let enr = discv5.local_enr();
456        // TCP and UDP share the ENR IP fields. Preserve their ports, including ports learned
457        // from peer votes, and only update address families the service is listening on.
458        let result = match external_ip {
459            IpAddr::V4(ip) if enr.udp4().is_some() && enr.ip4() != Some(ip) => {
460                discv5.with_discv5(|discv5| discv5.enr_insert("ip", &ip))
461            }
462            IpAddr::V6(ip) if enr.udp6().is_some() && enr.ip6() != Some(ip) => {
463                discv5.with_discv5(|discv5| discv5.enr_insert("ip6", &ip))
464            }
465            _ => return,
466        };
467        match result {
468            Ok(_) => debug!(target: "net::discovery", ?external_ip, "Updated external IP"),
469            Err(err) => {
470                debug!(target: "net::discovery", ?external_ip, %err, "Failed to update external IP");
471            }
472        }
473    }
474}
475
476const fn set_bound_rlpx_port_if_unset(config: &mut reth_discv5::Config, port: u16) {
477    if config.rlpx_socket().port() == 0 {
478        config.set_rlpx_port(port);
479    }
480}
481
482impl Drop for Discovery {
483    fn drop(&mut self) {
484        if let Some(discv4) = &self.discv4 {
485            discv4.terminate();
486        }
487        if let Some(handle) = self._discv4_service.take() {
488            handle.abort();
489        }
490        if let Some(handle) = self._discv5_forwarder.take() {
491            handle.abort();
492        }
493        if let Some(handle) = self._dns_disc_service.take() {
494            handle.abort();
495        }
496    }
497}
498
499impl Stream for Discovery {
500    type Item = DiscoveryEvent;
501
502    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
503        Poll::Ready(Some(ready!(self.get_mut().poll(cx))))
504    }
505}
506
507#[cfg(test)]
508impl Discovery {
509    /// Returns a Discovery instance that does nothing and is intended for testing purposes.
510    ///
511    /// NOTE: This instance does nothing
512    pub(crate) fn noop() -> Self {
513        let (_discovery_listeners, _): (mpsc::UnboundedSender<DiscoveryEvent>, _) =
514            mpsc::unbounded_channel();
515
516        Self {
517            discovered_nodes: LruMap::new(0),
518            local_enr: NodeRecord {
519                address: IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED),
520                tcp_port: 0,
521                udp_port: 0,
522                id: PeerId::random(),
523            },
524            discv4: Default::default(),
525            discv4_updates: Default::default(),
526            _discv4_service: Default::default(),
527            _discv5_forwarder: None,
528            discv5: None,
529            discv5_updates: None,
530            queued_events: Default::default(),
531            _dns_discovery: None,
532            dns_discovery_updates: None,
533            _dns_disc_service: None,
534            discovery_listeners: Default::default(),
535            nat_resolver: None,
536        }
537    }
538}
539
540#[cfg(test)]
541mod tests {
542    use super::*;
543    use secp256k1::SECP256K1;
544    use std::net::{Ipv4Addr, SocketAddrV4};
545
546    #[tokio::test(flavor = "multi_thread")]
547    async fn test_discovery_setup() {
548        let (secret_key, _) = SECP256K1.generate_keypair(&mut rand_08::thread_rng());
549        let discovery_addr = SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, 0));
550        let _discovery = Discovery::new(
551            discovery_addr,
552            discovery_addr,
553            secret_key,
554            Default::default(),
555            None,
556            Default::default(),
557            None,
558        )
559        .await
560        .unwrap();
561    }
562
563    use reth_discv4::Discv4ConfigBuilder;
564    use reth_discv5::{enr::EnrCombinedKeyWrapper, enr_to_discv4_id};
565    use tracing::trace;
566
567    async fn start_discovery_node(udp_port_discv4: u16, udp_port_discv5: u16) -> Discovery {
568        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
569
570        let discv4_addr = format!("127.0.0.1:{udp_port_discv4}").parse().unwrap();
571        let discv5_addr: SocketAddr = format!("127.0.0.1:{udp_port_discv5}").parse().unwrap();
572
573        // disable `NatResolver`
574        let discv4_config = Discv4ConfigBuilder::default().external_ip_resolver(None).build();
575
576        let discv5_listen_config = discv5::ListenConfig::from(discv5_addr);
577        let discv5_config = reth_discv5::Config::builder(discv5_addr)
578            .discv5_config(discv5::ConfigBuilder::new(discv5_listen_config).build())
579            .build();
580
581        Discovery::new(
582            discv4_addr,
583            discv4_addr,
584            secret_key,
585            Some(discv4_config),
586            Some(discv5_config),
587            None,
588            None,
589        )
590        .await
591        .expect("should build discv5 with discv4 downgrade")
592    }
593
594    async fn start_discv5_only(config: reth_discv5::Config, nat: Option<NatResolver>) -> Discovery {
595        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
596        Discovery::new(
597            "127.0.0.1:30303".parse().unwrap(),
598            "127.0.0.1:0".parse().unwrap(),
599            secret_key,
600            None,
601            Some(config),
602            None,
603            nat,
604        )
605        .await
606        .expect("should start discv5")
607    }
608
609    #[tokio::test]
610    async fn discv5_only_advertises_resolved_external_ip() {
611        let socket = Arc::new(UdpSocket::bind("0.0.0.0:0").await.unwrap());
612        let udp_port = socket.local_addr().unwrap().port();
613        let external_ip = "203.0.113.7".parse().unwrap();
614        // The advertised TCP port may differ from the listener port behind NAT.
615        let mut config = reth_discv5::Config::builder("127.0.0.1:30309".parse().unwrap()).build();
616        // Install pre-bound sockets after the builder normalizes its listen addresses.
617        let discv5_config = config.discv5_config_mut();
618        discv5_config.listen_config =
619            discv5::ListenConfig::FromSockets { ipv4: Some(socket), ipv6: None };
620        // Explicit NAT addresses still apply when peer-driven ENR updates are disabled.
621        discv5_config.enr_update = false;
622        let mut discovery =
623            start_discv5_only(config, Some(NatResolver::ExternalIp(external_ip))).await;
624        let discv5 = discovery.discv5().unwrap();
625        assert!(discv5.local_enr().ip4().is_none());
626
627        tokio::time::timeout(
628            Duration::from_secs(5),
629            futures::future::poll_fn(|cx| {
630                let _ = discovery.poll(cx);
631                if discv5.node_record().is_some_and(|record| record.address == external_ip) {
632                    Poll::Ready(())
633                } else {
634                    Poll::Pending
635                }
636            }),
637        )
638        .await
639        .expect("external IP should be advertised");
640
641        let record = discv5.node_record().unwrap();
642        assert_eq!(record.address, external_ip);
643        assert_eq!(record.udp_port, udp_port);
644        assert_eq!(record.tcp_port, 30309);
645
646        let seq = discv5.local_enr().seq();
647        discovery.on_external_ip(external_ip);
648        assert_eq!(discv5.local_enr().seq(), seq);
649        discovery.on_external_ip("2001:db8::7".parse().unwrap());
650        assert_eq!(discv5.local_enr().seq(), seq);
651
652        // Peer votes can change the address and mapped UDP port between NAT resolutions.
653        discv5.with_discv5(|discv5| {
654            discv5.update_local_enr_socket((Ipv4Addr::LOCALHOST, 30310).into(), false);
655        });
656        discovery.on_external_ip(external_ip);
657        assert_eq!(discv5.node_record().unwrap(), NodeRecord { udp_port: 30310, ..record });
658    }
659
660    #[tokio::test]
661    async fn discv5_only_respects_enr_update_policy() {
662        for (enr_update, auto_nat_listen_duration) in
663            [(false, None), (true, Some(Duration::from_secs(300)))]
664        {
665            let mut config = reth_discv5::Config::builder((Ipv4Addr::LOCALHOST, 0).into()).build();
666            let discv5_config = config.discv5_config_mut();
667            discv5_config.listen_config =
668                discv5::ListenConfig::Ipv4 { ip: Ipv4Addr::LOCALHOST, port: 0 };
669            discv5_config.enr_update = enr_update;
670            // Set this after building: the upstream builder clears the reachability timer.
671            discv5_config.auto_nat_listen_duration = auto_nat_listen_duration;
672            let discovery = start_discv5_only(config, Some(NatResolver::Any)).await;
673
674            // No automatic lookup may overwrite a pinned address or bypass reachability checks.
675            assert!(discovery.nat_resolver.is_none());
676        }
677    }
678
679    #[test]
680    fn discv5_enr_advertises_bound_rlpx_port() {
681        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
682
683        let bound_rlpx_port = 30307;
684        let discv5_addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
685        let mut discv5_config = reth_discv5::Config::builder((Ipv4Addr::LOCALHOST, 0).into())
686            .discv5_config(discv5::ConfigBuilder::new(discv5_addr.into()).build())
687            .build();
688
689        set_bound_rlpx_port_if_unset(&mut discv5_config, bound_rlpx_port);
690
691        let (enr, _, _, _) = reth_discv5::build_local_enr(&secret_key, &discv5_config);
692        assert_eq!(enr.tcp4(), Some(bound_rlpx_port));
693    }
694
695    #[test]
696    fn discv5_enr_preserves_configured_rlpx_port() {
697        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
698
699        let advertised_addr: SocketAddr = "127.0.0.1:30308".parse().unwrap();
700        let discv5_addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
701        let mut discv5_config = reth_discv5::Config::builder(advertised_addr)
702            .discv5_config(discv5::ConfigBuilder::new(discv5_addr.into()).build())
703            .build();
704
705        set_bound_rlpx_port_if_unset(&mut discv5_config, 30307);
706
707        let (enr, _, _, _) = reth_discv5::build_local_enr(&secret_key, &discv5_config);
708        assert_eq!(enr.tcp4(), Some(advertised_addr.port()));
709    }
710
711    #[tokio::test(flavor = "multi_thread")]
712    async fn discv5_and_discv4_same_pk() {
713        reth_tracing::init_test_tracing();
714
715        // set up test
716        let mut node_1 = start_discovery_node(40014, 40015).await;
717        let discv4_enr_1 = node_1.discv4.as_ref().unwrap().node_record();
718        let discv5_enr_node_1 =
719            node_1.discv5.as_ref().unwrap().with_discv5(|discv5| discv5.local_enr());
720        let discv4_id_1 = discv4_enr_1.id;
721        let discv5_id_1 = discv5_enr_node_1.node_id();
722
723        let mut node_2 = start_discovery_node(40024, 40025).await;
724        let discv4_enr_2 = node_2.discv4.as_ref().unwrap().node_record();
725        let discv5_enr_node_2 =
726            node_2.discv5.as_ref().unwrap().with_discv5(|discv5| discv5.local_enr());
727        let discv4_id_2 = discv4_enr_2.id;
728        let discv5_id_2 = discv5_enr_node_2.node_id();
729
730        trace!(target: "net::discovery::tests",
731            node_1_node_id=format!("{:#}", discv5_id_1),
732            node_2_node_id=format!("{:#}", discv5_id_2),
733            "started nodes"
734        );
735
736        // test
737
738        // assert discovery version 4 and version 5 nodes have same id
739        assert_eq!(discv4_id_1, enr_to_discv4_id(&discv5_enr_node_1).unwrap());
740        assert_eq!(discv4_id_2, enr_to_discv4_id(&discv5_enr_node_2).unwrap());
741
742        // add node_2:discv4 manually to node_1:discv4
743        node_1.add_discv4_node(discv4_enr_2);
744
745        // verify node_2:discv4 discovered node_1:discv4 and vv
746        let event_node_1 = node_1.next().await.unwrap();
747        let event_node_2 = node_2.next().await.unwrap();
748
749        assert_eq!(
750            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued {
751                peer_id: discv4_id_2,
752                addr: PeerAddr::new(discv4_enr_2.tcp_addr(), Some(discv4_enr_2.udp_addr())),
753                fork_id: None
754            }),
755            event_node_1
756        );
757        assert_eq!(
758            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued {
759                peer_id: discv4_id_1,
760                addr: PeerAddr::new(discv4_enr_1.tcp_addr(), Some(discv4_enr_1.udp_addr())),
761                fork_id: None
762            }),
763            event_node_2
764        );
765
766        assert_eq!(1, node_1.discovered_nodes.len());
767        assert_eq!(1, node_2.discovered_nodes.len());
768
769        // add node_2:discv5 to node_1:discv5, manual insertion won't emit an event
770        node_1.add_discv5_node(EnrCombinedKeyWrapper(discv5_enr_node_2.clone()).into()).unwrap();
771        // verify node_2 is in KBuckets of node_1:discv5
772        assert!(node_1
773            .discv5
774            .as_ref()
775            .unwrap()
776            .with_discv5(|discv5| discv5.table_entries_id().contains(&discv5_id_2)));
777
778        // manually trigger connection from node_1:discv5 to node_2:discv5
779        node_1
780            .discv5
781            .as_ref()
782            .unwrap()
783            .with_discv5(|discv5| discv5.send_ping(discv5_enr_node_2.clone()))
784            .await
785            .unwrap();
786
787        // this won't emit an event, since the nodes already discovered each other on discv4, the
788        // number of nodes stored for each node on this level remains 1.
789        assert_eq!(1, node_1.discovered_nodes.len());
790        assert_eq!(1, node_2.discovered_nodes.len());
791    }
792
793    /// Starts a discovery node with discv4 and discv5 sharing the same UDP port.
794    async fn start_shared_port_node(port: u16) -> Discovery {
795        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
796        let disc_addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap();
797        // Use a non-zero TCP port so the node record isn't filtered out by
798        // `on_node_record_update` (which drops peers with tcp port == 0).
799        let tcp_addr: SocketAddr = "127.0.0.1:30303".parse().unwrap();
800
801        let discv4_config = Discv4ConfigBuilder::default().external_ip_resolver(None).build();
802
803        let discv5_listen_config = discv5::ListenConfig::from(disc_addr);
804        let discv5_config = reth_discv5::Config::builder(tcp_addr)
805            .discv5_config(discv5::ConfigBuilder::new(discv5_listen_config).build())
806            .build();
807
808        // Both protocols use the same address, triggering shared-port mode
809        Discovery::new(
810            tcp_addr,
811            disc_addr,
812            secret_key,
813            Some(discv4_config),
814            Some(discv5_config),
815            None,
816            None,
817        )
818        .await
819        .expect("should start with shared port")
820    }
821
822    #[tokio::test(flavor = "multi_thread")]
823    async fn test_shared_port_setup() {
824        reth_tracing::init_test_tracing();
825
826        // Use port 0 so the OS picks a free port
827        let node = start_shared_port_node(0).await;
828
829        // Both protocols should be active
830        assert!(node.discv4.is_some(), "discv4 should be running");
831        assert!(node.discv5.is_some(), "discv5 should be running");
832    }
833
834    #[tokio::test(flavor = "multi_thread")]
835    async fn test_shared_port_discv5_discovery() {
836        reth_tracing::init_test_tracing();
837
838        let mut node_1 = start_shared_port_node(0).await;
839        let mut node_2 = start_shared_port_node(0).await;
840
841        let discv5_enr_1 = node_1.discv5.as_ref().unwrap().with_discv5(|discv5| discv5.local_enr());
842        let discv5_enr_2 = node_2.discv5.as_ref().unwrap().with_discv5(|discv5| discv5.local_enr());
843
844        let peer_id_1 = enr_to_discv4_id(&discv5_enr_1).unwrap();
845        let peer_id_2 = enr_to_discv4_id(&discv5_enr_2).unwrap();
846
847        // Add node_2's ENR to node_1's discv5 kbuckets and trigger a ping to establish a session.
848        // send_ping awaits the PONG, so the handshake completes before we poll the Discovery
849        // stream. The discv5 service runs its own background task.
850        node_1.add_discv5_node(EnrCombinedKeyWrapper(discv5_enr_2.clone()).into()).unwrap();
851        node_1
852            .discv5
853            .as_ref()
854            .unwrap()
855            .with_discv5(|discv5| discv5.send_ping(discv5_enr_2))
856            .await
857            .unwrap();
858
859        // Both SessionEstablished events should now be buffered in the update channels.
860        // Drive both nodes concurrently to collect them.
861        let mut event_1 = None;
862        let mut event_2 = None;
863        let timeout = tokio::time::sleep(std::time::Duration::from_secs(5));
864        tokio::pin!(timeout);
865        loop {
866            tokio::select! {
867                ev = node_1.next(), if event_1.is_none() => {
868                    event_1 = ev;
869                }
870                ev = node_2.next(), if event_2.is_none() => {
871                    event_2 = ev;
872                }
873                _ = &mut timeout => {
874                    panic!("timed out waiting for discv5 discovery events");
875                }
876            }
877            if event_1.is_some() && event_2.is_some() {
878                break;
879            }
880        }
881
882        assert!(matches!(
883            event_1.unwrap(),
884            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued { peer_id, .. })
885                if peer_id == peer_id_2
886        ));
887        assert!(matches!(
888            event_2.unwrap(),
889            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued { peer_id, .. })
890                if peer_id == peer_id_1
891        ));
892    }
893
894    #[tokio::test(flavor = "multi_thread")]
895    async fn test_shared_port_discv4_discovery() {
896        reth_tracing::init_test_tracing();
897
898        let mut node_1 = start_shared_port_node(0).await;
899        let mut node_2 = start_shared_port_node(0).await;
900
901        let enr_1 = node_1.discv4.as_ref().unwrap().node_record();
902        let enr_2 = node_2.discv4.as_ref().unwrap().node_record();
903
904        // Introduce node_2 to node_1 via discv4
905        node_1.add_discv4_node(enr_2);
906
907        // Both nodes should discover each other via discv4 ping/pong
908        let event_1 = node_1.next().await.unwrap();
909        let event_2 = node_2.next().await.unwrap();
910
911        assert_eq!(
912            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued {
913                peer_id: enr_2.id,
914                addr: PeerAddr::new(enr_2.tcp_addr(), Some(enr_2.udp_addr())),
915                fork_id: None
916            }),
917            event_1
918        );
919        assert_eq!(
920            DiscoveryEvent::NewNode(DiscoveredEvent::EventQueued {
921                peer_id: enr_1.id,
922                addr: PeerAddr::new(enr_1.tcp_addr(), Some(enr_1.udp_addr())),
923                fork_id: None
924            }),
925            event_2
926        );
927    }
928
929    /// Verifies that shared-port mode binds correctly when discv5 is configured for dual-stack.
930    /// On Linux this exercises the IPv6 V6ONLY path: without it, the IPv4 sibling would clash
931    /// with the IPv6 socket bound to the same port.
932    #[tokio::test(flavor = "multi_thread")]
933    async fn test_shared_port_dual_stack() {
934        reth_tracing::init_test_tracing();
935
936        // Find a port that's free on the v4 wildcard so we can use it for both v4 and v6.
937        let probe = UdpSocket::bind("0.0.0.0:0").await.expect("probe bind");
938        let port = probe.local_addr().unwrap().port();
939        drop(probe);
940
941        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
942        let v4_addr: SocketAddr = format!("0.0.0.0:{port}").parse().unwrap();
943        let tcp_addr: SocketAddr = "0.0.0.0:30303".parse().unwrap();
944
945        let discv4_config = Discv4ConfigBuilder::default().external_ip_resolver(None).build();
946
947        let discv5_listen_config = discv5::ListenConfig::DualStack {
948            ipv4: std::net::Ipv4Addr::UNSPECIFIED,
949            ipv4_port: port,
950            ipv6: std::net::Ipv6Addr::UNSPECIFIED,
951            ipv6_port: port,
952        };
953        let discv5_config = reth_discv5::Config::builder(tcp_addr)
954            .discv5_config(discv5::ConfigBuilder::new(discv5_listen_config).build())
955            .build();
956
957        Discovery::new(
958            tcp_addr,
959            v4_addr,
960            secret_key,
961            Some(discv4_config),
962            Some(discv5_config),
963            None,
964            None,
965        )
966        .await
967        .expect("discovery should start with shared port + dual-stack");
968    }
969}