1use 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
32pub const DEFAULT_MAX_CAPACITY_DISCOVERED_PEERS_CACHE: u32 = 10_000;
36
37const RESOLVE_EXTERNAL_IP_INTERVAL: Duration = Duration::from_secs(60 * 5);
41
42#[derive(Debug)]
47pub struct Discovery {
48 discovered_nodes: LruMap<PeerId, PeerAddr>,
52 local_enr: NodeRecord,
54 discv4: Option<Discv4>,
56 discv4_updates: Option<ReceiverStream<DiscoveryUpdate>>,
58 _discv4_service: Option<JoinHandle<()>>,
60 discv5: Option<Discv5>,
62 discv5_updates: Option<ReceiverStream<discv5::Event>>,
64 _discv5_forwarder: Option<JoinHandle<()>>,
67 _dns_discovery: Option<DnsDiscoveryHandle>,
69 dns_discovery_updates: Option<ReceiverStream<DnsNodeRecordUpdate>>,
71 _dns_disc_service: Option<JoinHandle<()>>,
73 queued_events: VecDeque<DiscoveryEvent>,
75 discovery_listeners: Vec<mpsc::UnboundedSender<DiscoveryEvent>>,
77 nat_resolver: Option<ResolveNatInterval>,
79}
80
81impl Discovery {
82 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>, dns_discovery_config: Option<DnsDiscoveryConfig>,
96 nat: Option<NatResolver>,
97 ) -> Result<Self, NetworkError> {
98 let local_enr =
100 NodeRecord::from_secret_key(discovery_v4_addr, &sk).with_tcp_port(tcp_addr.port());
101
102 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 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 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 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 NatResolver::ExternalIp(_) | NatResolver::ExternalAddr(_) => true,
184 _ => 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 let (discv5, discv5_updates) = if let Some(mut config) = discv5_config {
195 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 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 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 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 pub(crate) fn add_listener(&mut self, tx: mpsc::UnboundedSender<DiscoveryEvent>) {
283 self.discovery_listeners.push(tx);
284 }
285
286 #[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 pub(crate) fn update_fork_id(&self, fork_id: ForkId) {
294 if let Some(discv4) = &self.discv4 {
295 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 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 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 pub fn discv4(&self) -> Option<Discv4> {
326 self.discv4.clone()
327 }
328
329 pub(crate) const fn local_id(&self) -> PeerId {
331 self.local_enr.id }
333
334 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 pub fn discv5(&self) -> Option<Discv5> {
343 self.discv5.clone()
344 }
345
346 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 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 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 if let Some(event) = self.queued_events.pop_front() {
398 self.notify_listeners(&event);
399 return Poll::Ready(event)
400 }
401
402 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 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 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 fn on_external_ip(&self, external_ip: IpAddr) {
454 let Some(discv5) = &self.discv5 else { return };
455 let enr = discv5.local_enr();
456 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 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 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 let mut config = reth_discv5::Config::builder("127.0.0.1:30309".parse().unwrap()).build();
616 let discv5_config = config.discv5_config_mut();
618 discv5_config.listen_config =
619 discv5::ListenConfig::FromSockets { ipv4: Some(socket), ipv6: None };
620 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 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 discv5_config.auto_nat_listen_duration = auto_nat_listen_duration;
672 let discovery = start_discv5_only(config, Some(NatResolver::Any)).await;
673
674 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 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 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 node_1.add_discv4_node(discv4_enr_2);
744
745 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 node_1.add_discv5_node(EnrCombinedKeyWrapper(discv5_enr_node_2.clone()).into()).unwrap();
771 assert!(node_1
773 .discv5
774 .as_ref()
775 .unwrap()
776 .with_discv5(|discv5| discv5.table_entries_id().contains(&discv5_id_2)));
777
778 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 assert_eq!(1, node_1.discovered_nodes.len());
790 assert_eq!(1, node_2.discovered_nodes.len());
791 }
792
793 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 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 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 let node = start_shared_port_node(0).await;
828
829 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 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 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 node_1.add_discv4_node(enr_2);
906
907 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 #[tokio::test(flavor = "multi_thread")]
933 async fn test_shared_port_dual_stack() {
934 reth_tracing::init_test_tracing();
935
936 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}