1use crate::{
4 builder::ETH_REQUEST_CHANNEL_CAPACITY,
5 config::rng_secret_key,
6 error::NetworkError,
7 eth_requests::EthRequestHandler,
8 protocol::IntoRlpxSubProtocol,
9 transactions::{
10 config::{StrictEthAnnouncementFilter, TransactionPropagationKind},
11 constants::tx_manager::DEFAULT_TX_MANAGER_CHANNEL_MEMORY_LIMIT_BYTES,
12 policy::NetworkPolicies,
13 TransactionsHandle, TransactionsManager, TransactionsManagerConfig,
14 },
15 NetworkConfigBuilder, NetworkHandle, NetworkManager, PeersConfig,
16};
17use futures::{FutureExt, StreamExt};
18use pin_project::pin_project;
19use reth_chainspec::{ChainSpecProvider, EthereumHardforks, Hardforks};
20use reth_eth_wire::{
21 protocol::Protocol, DisconnectReason, EthNetworkPrimitives, HelloMessageWithProtocols,
22};
23use reth_ethereum_primitives::{PooledTransactionVariant, TransactionSigned};
24use reth_evm_ethereum::EthEvmConfig;
25use reth_metrics::common::mpsc::memory_bounded_channel;
26use reth_network_api::{
27 events::{PeerEvent, SessionInfo},
28 test_utils::{PeersHandle, PeersHandleProvider},
29 NetworkEvent, NetworkEventListenerProvider, NetworkInfo, PeerRequest, Peers,
30};
31use reth_network_p2p::error::{RequestError, RequestResult};
32use reth_network_peers::PeerId;
33use reth_storage_api::{
34 noop::NoopProvider, BalProvider, BlockNumReader, BlockReader, BlockReaderIdExt, HeaderProvider,
35 StateProviderFactory, StateRangeProviderFactory,
36};
37use reth_tasks::Runtime;
38use reth_tokio_util::EventStream;
39use reth_transaction_pool::{
40 blobstore::InMemoryBlobStore, test_utils::TestPool, EthTransactionPool, PoolTransaction,
41 TransactionPool, TransactionValidationTaskExecutor,
42};
43use secp256k1::SecretKey;
44use std::{
45 fmt,
46 future::Future,
47 net::{Ipv4Addr, SocketAddr, SocketAddrV4},
48 pin::Pin,
49 task::{Context, Poll},
50};
51use tokio::{
52 sync::{mpsc::channel, oneshot},
53 task::JoinHandle,
54};
55
56pub struct Testnet<C, Pool> {
58 peers: Vec<Peer<C, Pool>>,
60}
61
62impl<C> Testnet<C, TestPool>
65where
66 C: BlockNumReader + ChainSpecProvider<ChainSpec: Hardforks> + Clone + 'static,
67{
68 pub async fn from_configs(configs: impl IntoIterator<Item = PeerConfig<C>>) -> Self {
76 let peers = futures::future::try_join_all(configs.into_iter().map(PeerConfig::launch))
77 .await
78 .expect("failed to launch testnet peers");
79 Self { peers }
80 }
81
82 pub async fn create_with(num_peers: usize, provider: C) -> Self {
88 Self::from_configs((0..num_peers).map(|_| PeerConfig::new(provider.clone()))).await
89 }
90
91 pub async fn add_peer_with_config(
93 &mut self,
94 config: PeerConfig<C>,
95 ) -> Result<(), NetworkError> {
96 self.peers.push(config.launch().await?);
97 Ok(())
98 }
99}
100
101impl Testnet<NoopProvider, TestPool> {
102 pub async fn create(num_peers: usize) -> Self {
108 Self::create_with(num_peers, NoopProvider::default()).await
109 }
110}
111
112impl<C, Pool> Testnet<C, Pool>
113where
114 C: BlockReader + HeaderProvider + Clone + 'static,
115 Pool: TransactionPool,
116{
117 pub fn peers_mut(&mut self) -> &mut [Peer<C, Pool>] {
119 &mut self.peers
120 }
121
122 pub fn peers(&self) -> &[Peer<C, Pool>] {
124 &self.peers
125 }
126
127 pub fn map_pool<F, P>(self, f: F) -> Testnet<C, P>
129 where
130 F: Fn(Peer<C, Pool>) -> Peer<C, P>,
131 P: TransactionPool,
132 {
133 Testnet { peers: self.peers.into_iter().map(f).collect() }
134 }
135
136 pub fn for_each_mut<F>(&mut self, f: F)
138 where
139 F: FnMut(&mut Peer<C, Pool>),
140 {
141 self.peers.iter_mut().for_each(f)
142 }
143
144 pub fn with_request_handlers(mut self) -> Self
146 where
147 C: BalProvider,
148 {
149 self.for_each_mut(Peer::install_request_handler);
150 self
151 }
152}
153
154impl<C, Pool> Testnet<C, Pool>
155where
156 C: ChainSpecProvider<ChainSpec: EthereumHardforks>
157 + StateProviderFactory
158 + BlockReaderIdExt
159 + HeaderProvider<Header = alloy_consensus::Header>
160 + Clone
161 + 'static,
162 Pool: TransactionPool,
163{
164 pub fn with_eth_pool(
166 self,
167 ) -> Testnet<C, EthTransactionPool<C, InMemoryBlobStore, EthEvmConfig>> {
168 self.with_eth_pool_config(Default::default())
169 }
170
171 pub fn with_eth_pool_config(
173 self,
174 tx_manager_config: TransactionsManagerConfig,
175 ) -> Testnet<C, EthTransactionPool<C, InMemoryBlobStore, EthEvmConfig>> {
176 self.with_eth_pool_config_and_policy(tx_manager_config, Default::default())
177 }
178
179 pub fn with_eth_pool_config_and_policy(
181 self,
182 tx_manager_config: TransactionsManagerConfig,
183 policy: TransactionPropagationKind,
184 ) -> Testnet<C, EthTransactionPool<C, InMemoryBlobStore, EthEvmConfig>> {
185 self.map_pool(|peer| {
186 let blob_store = InMemoryBlobStore::default();
187 let validator = TransactionValidationTaskExecutor::eth(
188 peer.client.clone(),
189 EthEvmConfig::mainnet(),
190 blob_store.clone(),
191 Runtime::test(),
192 );
193 let pool = EthTransactionPool::eth_pool(validator, blob_store, Default::default());
194 peer.map_transactions_manager(pool, tx_manager_config.clone(), policy)
195 })
196 }
197}
198
199impl<C, Pool> Testnet<C, Pool>
200where
201 C: TestnetProvider + Clone,
202 Pool: TestnetPool,
203{
204 pub fn spawn(self) -> TestnetHandle<C, Pool> {
206 let (terminate, rx) = oneshot::channel::<oneshot::Sender<Self>>();
207 let peers = self.peers.iter().map(Peer::peer_handle).collect();
208 let mut net = self;
209 let handle = tokio::task::spawn(async move {
210 let tx = tokio::select! {
211 _ = &mut net => None,
212 tx = rx => tx.ok(),
213 };
214 if let Some(tx) = tx {
215 let _ = tx.send(net);
216 }
217 });
218
219 TestnetHandle { _handle: handle, peers, terminate }
220 }
221}
222
223impl<C, Pool> fmt::Debug for Testnet<C, Pool> {
224 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
225 f.debug_struct("Testnet").finish_non_exhaustive()
226 }
227}
228
229impl<C, Pool> Future for Testnet<C, Pool>
230where
231 C: TestnetProvider,
232 Pool: TestnetPool,
233{
234 type Output = ();
235
236 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
237 let this = self.get_mut();
238 for peer in &mut this.peers {
239 let _ = peer.poll_unpin(cx);
240 }
241 Poll::Pending
242 }
243}
244
245#[derive(Debug)]
247pub struct TestnetHandle<C, Pool> {
248 _handle: JoinHandle<()>,
249 peers: Vec<PeerHandle<Pool>>,
250 terminate: oneshot::Sender<oneshot::Sender<Testnet<C, Pool>>>,
251}
252
253impl<C, Pool> TestnetHandle<C, Pool> {
256 pub async fn terminate(self) -> Testnet<C, Pool> {
258 let (tx, rx) = oneshot::channel();
259 self.terminate.send(tx).unwrap();
260 rx.await.unwrap()
261 }
262
263 pub fn peers(&self) -> &[PeerHandle<Pool>] {
265 &self.peers
266 }
267
268 pub fn peers_array<const N: usize>(&self) -> &[PeerHandle<Pool>; N] {
275 self.peers.as_slice().try_into().unwrap_or_else(|_| {
276 panic!("expected {N} peers, but the testnet has {}", self.peers.len())
277 })
278 }
279
280 pub async fn connect_peers(&self) {
286 if self.peers.len() < 2 {
287 return
288 }
289
290 let streams = self.peers.iter().map(PeerHandle::event_stream).collect::<Vec<_>>();
292
293 for (idx, handle) in self.peers.iter().enumerate().take(self.peers.len() - 1) {
295 for neighbour in &self.peers[idx + 1..] {
296 handle.add_peer(neighbour);
297 }
298 }
299
300 let num_sessions_per_peer = self.peers.len() - 1;
302 let fut = streams.into_iter().map(|mut stream| async move {
303 stream.take_session_established(num_sessions_per_peer).await
304 });
305
306 futures::future::join_all(fut).await;
307 }
308}
309
310#[pin_project]
312#[derive(Debug)]
313pub struct Peer<C, Pool = TestPool> {
314 #[pin]
315 network: NetworkManager<EthNetworkPrimitives>,
316 #[pin]
317 request_handler: Option<EthRequestHandler<C, EthNetworkPrimitives>>,
318 #[pin]
319 transactions_manager: Option<TransactionsManager<Pool, EthNetworkPrimitives>>,
320 pool: Option<Pool>,
321 client: C,
322}
323
324impl<C, Pool> Peer<C, Pool>
327where
328 C: BlockReader + HeaderProvider + Clone + 'static,
329 Pool: TransactionPool,
330{
331 pub fn num_peers(&self) -> usize {
333 self.network.num_connected_peers()
334 }
335
336 pub fn add_rlpx_sub_protocol(&mut self, protocol: impl IntoRlpxSubProtocol) {
338 self.network.add_rlpx_sub_protocol(protocol);
339 }
340
341 pub fn peer_handle(&self) -> PeerHandle<Pool> {
343 PeerHandle {
344 network: self.network.handle().clone(),
345 pool: self.pool.clone(),
346 transactions: self.transactions_manager.as_ref().map(|mgr| mgr.handle()),
347 }
348 }
349
350 pub const fn local_addr(&self) -> SocketAddr {
352 self.network.local_addr()
353 }
354
355 pub fn peer_id(&self) -> PeerId {
357 *self.network.peer_id()
358 }
359
360 pub const fn network_mut(&mut self) -> &mut NetworkManager<EthNetworkPrimitives> {
362 &mut self.network
363 }
364
365 pub fn handle(&self) -> NetworkHandle<EthNetworkPrimitives> {
367 self.network.handle().clone()
368 }
369
370 pub const fn pool(&self) -> Option<&Pool> {
372 self.pool.as_ref()
373 }
374
375 pub fn install_request_handler(&mut self)
377 where
378 C: BalProvider,
379 {
380 let (tx, rx) = channel(ETH_REQUEST_CHANNEL_CAPACITY);
381 self.network.set_eth_request_handler(tx);
382 let peers = self.network.peers_handle();
383 let request_handler = EthRequestHandler::new(self.client.clone(), peers, rx);
384 self.request_handler = Some(request_handler);
385 }
386
387 pub fn install_transactions_manager(&mut self, pool: Pool) {
390 self.transactions_manager = Some(new_transactions_manager(
391 &mut self.network,
392 pool.clone(),
393 Default::default(),
394 Default::default(),
395 ));
396 self.pool = Some(pool);
397 }
398
399 pub fn map_transactions_manager<P>(
402 self,
403 pool: P,
404 config: TransactionsManagerConfig,
405 policy: TransactionPropagationKind,
406 ) -> Peer<C, P>
407 where
408 P: TransactionPool,
409 {
410 let Self { mut network, request_handler, client, .. } = self;
411 let transactions_manager =
412 new_transactions_manager(&mut network, pool.clone(), config, policy);
413 Peer {
414 network,
415 request_handler,
416 transactions_manager: Some(transactions_manager),
417 pool: Some(pool),
418 client,
419 }
420 }
421}
422
423impl<C, Pool> Future for Peer<C, Pool>
424where
425 C: TestnetProvider,
426 Pool: TestnetPool,
427{
428 type Output = ();
429
430 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
431 let this = self.project();
432
433 if let Some(request) = this.request_handler.as_pin_mut() {
434 let _ = request.poll(cx);
435 }
436
437 if let Some(tx_manager) = this.transactions_manager.as_pin_mut() {
438 let _ = tx_manager.poll(cx);
439 }
440
441 this.network.poll(cx)
442 }
443}
444
445#[derive(Debug)]
449pub struct PeerConfig<C = NoopProvider> {
450 client: C,
451 secret_key: SecretKey,
452 protocols: Option<Vec<Protocol>>,
453 peers_config: PeersConfig,
454}
455
456#[derive(Debug)]
458pub struct PeerHandle<Pool> {
459 network: NetworkHandle<EthNetworkPrimitives>,
460 transactions: Option<TransactionsHandle<EthNetworkPrimitives>>,
461 pool: Option<Pool>,
462}
463
464impl<Pool> PeerHandle<Pool> {
467 pub fn peer_id(&self) -> &PeerId {
469 self.network.peer_id()
470 }
471
472 pub fn peer_handle(&self) -> &PeersHandle {
474 self.network.peers_handle()
475 }
476
477 pub fn local_addr(&self) -> SocketAddr {
479 self.network.local_addr()
480 }
481
482 pub fn event_listener(&self) -> EventStream<NetworkEvent> {
484 self.network.event_listener()
485 }
486
487 pub const fn transactions(&self) -> Option<&TransactionsHandle> {
489 self.transactions.as_ref()
490 }
491
492 pub const fn pool(&self) -> Option<&Pool> {
494 self.pool.as_ref()
495 }
496
497 pub const fn network(&self) -> &NetworkHandle<EthNetworkPrimitives> {
499 &self.network
500 }
501
502 pub fn add_peer<P>(&self, other: &PeerHandle<P>) {
504 self.network.add_peer(*other.peer_id(), other.local_addr());
505 }
506
507 pub fn add_trusted_peer<P>(&self, other: &PeerHandle<P>) {
509 self.network.add_trusted_peer(*other.peer_id(), other.local_addr());
510 }
511
512 pub fn event_stream(&self) -> NetworkEventStream {
514 NetworkEventStream::new(self.event_listener())
515 }
516
517 pub async fn request<R>(
521 &self,
522 peer_id: PeerId,
523 request: impl FnOnce(oneshot::Sender<RequestResult<R>>) -> PeerRequest,
524 ) -> RequestResult<R> {
525 let (tx, rx) = oneshot::channel();
526 self.network.send_request(peer_id, request(tx));
527 rx.await.unwrap_or(Err(RequestError::ChannelClosed))
528 }
529}
530
531impl<C> PeerConfig<C> {
534 pub fn new(client: C) -> Self {
538 Self {
539 client,
540 secret_key: rng_secret_key(),
541 protocols: None,
542 peers_config: PeersConfig::test(),
543 }
544 }
545
546 pub const fn with_secret_key(mut self, secret_key: SecretKey) -> Self {
548 self.secret_key = secret_key;
549 self
550 }
551
552 pub fn with_protocols(mut self, protocols: impl IntoIterator<Item: Into<Protocol>>) -> Self {
554 self.protocols = Some(protocols.into_iter().map(Into::into).collect());
555 self
556 }
557
558 pub fn with_peers_config(mut self, peers_config: PeersConfig) -> Self {
560 self.peers_config = peers_config;
561 self
562 }
563
564 pub async fn launch(self) -> Result<Peer<C>, NetworkError>
566 where
567 C: BlockNumReader + ChainSpecProvider<ChainSpec: Hardforks> + Clone + 'static,
568 {
569 let Self { client, secret_key, protocols, peers_config } = self;
570 let mut builder = NetworkConfigBuilder::new(secret_key, Runtime::test())
571 .listener_addr(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, 0)))
572 .discovery_addr(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, 0)))
573 .disable_dns_discovery()
574 .disable_discv4_discovery()
575 .peer_config(peers_config);
576 if let Some(protocols) = protocols {
577 let snap_enabled = protocols.iter().any(|p| p.cap.name == Protocol::snap_2().cap.name);
580 let hello_message = HelloMessageWithProtocols::builder(builder.get_peer_id())
581 .protocols(protocols)
582 .build();
583 builder = builder.with_snap(snap_enabled).hello_message(hello_message);
584 }
585
586 let network = NetworkManager::new(builder.build(client.clone())).await?;
587 Ok(Peer { network, client, request_handler: None, transactions_manager: None, pool: None })
588 }
589}
590
591impl Default for PeerConfig {
592 fn default() -> Self {
593 Self::new(NoopProvider::default())
594 }
595}
596
597#[derive(Debug)]
601pub struct NetworkEventStream {
602 inner: EventStream<NetworkEvent>,
603}
604
605impl NetworkEventStream {
608 pub const fn new(inner: EventStream<NetworkEvent>) -> Self {
610 Self { inner }
611 }
612
613 pub async fn next_session_closed(&mut self) -> Option<(PeerId, Option<DisconnectReason>)> {
615 while let Some(ev) = self.inner.next().await {
616 if let NetworkEvent::Peer(PeerEvent::SessionClosed { peer_id, reason }) = ev {
617 return Some((peer_id, reason))
618 }
619 }
620 None
621 }
622
623 pub async fn next_session_established(&mut self) -> Option<PeerId> {
625 while let Some(ev) = self.inner.next().await {
626 match ev {
627 NetworkEvent::ActivePeerSession { info, .. } |
628 NetworkEvent::Peer(PeerEvent::SessionEstablished(info)) => {
629 return Some(info.peer_id)
630 }
631 _ => {}
632 }
633 }
634 None
635 }
636
637 pub async fn take_session_established(&mut self, mut num: usize) -> Vec<PeerId> {
639 if num == 0 {
640 return Vec::new();
641 }
642 let mut peers = Vec::with_capacity(num);
643 while let Some(ev) = self.inner.next().await {
644 if let NetworkEvent::ActivePeerSession { info: SessionInfo { peer_id, .. }, .. } = ev {
645 peers.push(peer_id);
646 num -= 1;
647 if num == 0 {
648 return peers;
649 }
650 }
651 }
652 peers
653 }
654
655 pub async fn peer_added(&mut self) -> Option<PeerId> {
657 match self.inner.next().await {
658 Some(NetworkEvent::Peer(PeerEvent::PeerAdded(peer_id))) => Some(peer_id),
659 _ => None,
660 }
661 }
662
663 pub async fn peer_removed(&mut self) -> Option<PeerId> {
665 match self.inner.next().await {
666 Some(NetworkEvent::Peer(PeerEvent::PeerRemoved(peer_id))) => Some(peer_id),
667 _ => None,
668 }
669 }
670}
671
672pub trait TestnetProvider:
674 BlockReader<
675 Block = reth_ethereum_primitives::Block,
676 Receipt = reth_ethereum_primitives::Receipt,
677 Header = alloy_consensus::Header,
678 > + HeaderProvider
679 + BalProvider
680 + StateProviderFactory
681 + StateRangeProviderFactory
682 + Unpin
683 + 'static
684{
685}
686
687impl<T> TestnetProvider for T where
688 T: BlockReader<
689 Block = reth_ethereum_primitives::Block,
690 Receipt = reth_ethereum_primitives::Receipt,
691 Header = alloy_consensus::Header,
692 > + HeaderProvider
693 + BalProvider
694 + StateProviderFactory
695 + StateRangeProviderFactory
696 + Unpin
697 + 'static
698{
699}
700
701pub trait TestnetPool:
704 TransactionPool<
705 Transaction: PoolTransaction<
706 Consensus = TransactionSigned,
707 Pooled = PooledTransactionVariant,
708 >,
709 > + Unpin
710 + 'static
711{
712}
713
714impl<T> TestnetPool for T where
715 T: TransactionPool<
716 Transaction: PoolTransaction<
717 Consensus = TransactionSigned,
718 Pooled = PooledTransactionVariant,
719 >,
720 > + Unpin
721 + 'static
722{
723}
724
725fn new_transactions_manager<Pool: TransactionPool>(
727 network: &mut NetworkManager<EthNetworkPrimitives>,
728 pool: Pool,
729 config: TransactionsManagerConfig,
730 policy: TransactionPropagationKind,
731) -> TransactionsManager<Pool, EthNetworkPrimitives> {
732 let (tx, rx) =
733 memory_bounded_channel(DEFAULT_TX_MANAGER_CHANNEL_MEMORY_LIMIT_BYTES, "test_tx_channel");
734 network.set_transactions(tx);
735 let policies = NetworkPolicies::new(policy, StrictEthAnnouncementFilter::default());
736 TransactionsManager::with_policy(network.handle().clone(), pool, rx, config, policies)
737}