1use crate::{
2 network::NetworkTestContext,
3 payload::PayloadTestContext,
4 rpc::RpcTestContext,
5 wait::{poll_until, POLL_INTERVAL, WAIT_TIMEOUT},
6};
7use alloy_consensus::{transaction::TxHashRef, BlockHeader};
8use alloy_eips::BlockId;
9use alloy_network::{Ethereum, IntoWallet};
10use alloy_primitives::{BlockHash, BlockNumber, Bytes, Sealable, B256};
11use alloy_provider::{
12 fillers::{FillProvider, RecommendedFillers, TxFiller},
13 Provider, ProviderBuilder, RootProvider,
14};
15use alloy_rpc_types_engine::{ExecutionPayloadEnvelopeV5, ForkchoiceState, ForkchoiceUpdated};
16use alloy_rpc_types_eth::BlockNumberOrTag;
17use eyre::{ensure, eyre, Ok};
18use futures_util::{
19 future::{select, Either},
20 Future,
21};
22use jsonrpsee::{core::client::ClientT, http_client::HttpClient};
23use reth_chainspec::EthereumHardforks;
24use reth_network_api::test_utils::PeersHandleProvider;
25use reth_node_api::{Block, BlockBody, BlockTy, FullNodeComponents, PayloadTypes, PrimitivesTy};
26use reth_node_builder::{rpc::RethRpcAddOns, FullNode, NodeTypes};
27use reth_payload_primitives::BuiltPayload;
28use reth_provider::{
29 BlockNumReader, BlockReader, BlockReaderIdExt, CanonStateNotificationStream,
30 CanonStateSubscriptions, DatabaseProviderFactory, HeaderProvider, StageCheckpointReader,
31};
32use reth_rpc_api::TestingBuildBlockRequestV1;
33use reth_rpc_builder::auth::AuthServerHandle;
34use reth_rpc_eth_api::{
35 helpers::{EthApiSpec, EthTransactions, LoadReceipt, TraceExt},
36 EthApiTypes, RpcReceipt,
37};
38use reth_stages_types::StageId;
39use reth_transaction_pool::TransactionPool;
40use std::{
41 pin::{pin, Pin},
42 sync::Arc,
43 time::Duration,
44};
45use tokio_stream::StreamExt;
46use url::Url;
47
48pub const RPC_PROVIDER_POLL_INTERVAL: Duration = Duration::from_millis(10);
54
55#[expect(missing_debug_implementations)]
57pub struct NodeTestContext<Node, AddOns>
58where
59 Node: FullNodeComponents,
60 AddOns: RethRpcAddOns<Node>,
61{
62 pub inner: FullNode<Node, AddOns>,
64 pub payload: PayloadTestContext<<Node::Types as NodeTypes>::Payload>,
66 pub network: NetworkTestContext<Node::Network>,
68 pub rpc: RpcTestContext<Node, AddOns::EthApi>,
70 pub canonical_stream: CanonStateNotificationStream<PrimitivesTy<Node::Types>>,
72}
73
74impl<Node, Payload, AddOns> NodeTestContext<Node, AddOns>
75where
76 Payload: PayloadTypes,
77 Node: FullNodeComponents,
78 Node::Types: NodeTypes<ChainSpec: EthereumHardforks, Payload = Payload>,
79 Node::Network: PeersHandleProvider,
80 AddOns: RethRpcAddOns<Node>,
81{
82 pub async fn new(
88 node: FullNode<Node, AddOns>,
89 attributes_generator: impl Fn(u64) -> Payload::PayloadAttributes + Send + Sync + 'static,
90 ) -> eyre::Result<Self> {
91 let mut payload =
92 PayloadTestContext::new(node.payload_builder_handle.clone(), attributes_generator)
93 .await?;
94 if let Some(latest) =
95 node.provider.sealed_header_by_number_or_tag(BlockNumberOrTag::Latest)?
96 {
97 payload.timestamp = payload.timestamp.max(latest.timestamp());
98 }
99 Ok(Self {
100 inner: node.clone(),
101 payload,
102 network: NetworkTestContext::new(node.network.clone()),
103 rpc: RpcTestContext { inner: node.add_ons_handle.rpc_registry },
104 canonical_stream: node.provider.canonical_state_stream(),
105 })
106 }
107
108 pub async fn connect(&mut self, node: &mut Self) {
110 self.network.add_peer(node.network.record()).await;
111 node.network.next_session_established().await;
112 self.network.next_session_established().await;
113 }
114
115 pub async fn advance(
119 &mut self,
120 length: u64,
121 tx_generator: impl Fn(u64) -> Pin<Box<dyn Future<Output = Bytes>>>,
122 ) -> eyre::Result<Vec<Payload::BuiltPayload>>
123 where
124 AddOns::EthApi: EthApiSpec<Provider: BlockReader<Block = BlockTy<Node::Types>>>
125 + EthTransactions
126 + TraceExt,
127 {
128 let mut chain = Vec::with_capacity(length as usize);
129 for i in 0..length {
130 let (_, payload) = self.inject_and_advance(tx_generator(i).await).await?;
131 chain.push(payload);
132 }
133 Ok(chain)
134 }
135
136 pub async fn inject_and_advance(
142 &mut self,
143 raw_tx: Bytes,
144 ) -> eyre::Result<(B256, Payload::BuiltPayload)>
145 where
146 AddOns::EthApi: EthApiSpec<Provider: BlockReader<Block = BlockTy<Node::Types>>>
147 + EthTransactions
148 + TraceExt,
149 {
150 let tx_hash = self.rpc.inject_tx(raw_tx).await?;
151 let payload = self.advance_block().await?;
152 self.assert_new_block(tx_hash, payload.block().hash(), payload.block().number()).await?;
153 Ok((tx_hash, payload))
154 }
155
156 pub fn current_forkchoice_state(&self) -> eyre::Result<ForkchoiceState> {
158 let latest_header =
159 self.inner.provider.sealed_header_by_number_or_tag(BlockNumberOrTag::Latest)?.unwrap();
160
161 if latest_header.number() == 0 {
162 return Ok(ForkchoiceState::same_hash(latest_header.hash()));
163 }
164
165 Ok(ForkchoiceState {
166 head_block_hash: latest_header.hash(),
167 safe_block_hash: self
168 .inner
169 .provider
170 .sealed_header_by_number_or_tag(BlockNumberOrTag::Safe)?
171 .unwrap()
172 .hash(),
173 finalized_block_hash: self
174 .inner
175 .provider
176 .sealed_header_by_number_or_tag(BlockNumberOrTag::Finalized)?
177 .unwrap()
178 .hash(),
179 })
180 }
181
182 pub fn set_next_payload_timestamp(&mut self, timestamp: u64) -> eyre::Result<()> {
191 let latest = self
192 .inner
193 .provider
194 .sealed_header_by_number_or_tag(BlockNumberOrTag::Latest)?
195 .ok_or_else(|| eyre!("latest block not found"))?;
196 ensure!(
197 timestamp > latest.timestamp(),
198 "next payload timestamp {timestamp} must be greater than the latest block timestamp {}",
199 latest.timestamp()
200 );
201 self.payload.timestamp = timestamp - 1;
203 Ok(())
204 }
205
206 pub async fn new_payload(&mut self) -> eyre::Result<Payload::BuiltPayload> {
211 let eth_attr = self.payload.next_attributes();
212 let payload_id = self
213 .inner
214 .add_ons_handle
215 .beacon_engine_handle
216 .fork_choice_updated(self.current_forkchoice_state()?, Some(eth_attr.clone()))
217 .await?
218 .payload_id
219 .unwrap();
220 self.payload.expect_attr_event(eth_attr).await?;
222 self.payload.wait_for_built_payload(payload_id).await;
224 Ok(self.payload.expect_built_payload().await?)
226 }
227
228 pub async fn build_and_submit_payload(&mut self) -> eyre::Result<Payload::BuiltPayload> {
230 let payload = self.new_payload().await?;
231
232 self.submit_payload(payload.clone()).await?;
233
234 Ok(payload)
235 }
236
237 pub async fn advance_block(&mut self) -> eyre::Result<Payload::BuiltPayload> {
240 let payload = self.new_payload().await?;
241
242 self.import_payload(payload.clone()).await?;
243
244 Ok(payload)
245 }
246
247 pub async fn advance_block_synced(&mut self) -> eyre::Result<Payload::BuiltPayload> {
250 let payload = self.advance_block().await?;
251 self.wait_for_pool_head(payload.block().hash()).await?;
252 Ok(payload)
253 }
254
255 pub async fn advance_blocks(
262 &mut self,
263 length: u64,
264 ) -> eyre::Result<Vec<Payload::BuiltPayload>> {
265 let mut chain = Vec::with_capacity(length as usize);
266 for _ in 0..length {
267 chain.push(self.advance_block().await?);
268 }
269 Ok(chain)
270 }
271
272 pub async fn advance_until_receipt(
278 &mut self,
279 hash: B256,
280 ) -> eyre::Result<RpcReceipt<<AddOns::EthApi as EthApiTypes>::NetworkTypes>>
281 where
282 AddOns::EthApi: EthApiSpec<Provider: BlockReader<Block = BlockTy<Node::Types>>>
283 + EthTransactions
284 + TraceExt
285 + LoadReceipt
286 + 'static,
287 {
288 let wait = async {
289 loop {
290 if let Some(receipt) = self.rpc.transaction_receipt(hash).await? {
291 return Ok(receipt)
292 }
293 self.advance_block().await?;
294 }
295 };
296 tokio::time::timeout(WAIT_TIMEOUT, wait)
297 .await
298 .map_err(|_| eyre!("timed out waiting for the receipt of transaction {hash}"))?
299 }
300
301 pub async fn advance_while<F: Future>(&mut self, fut: F) -> eyre::Result<F::Output> {
311 let wait = async {
312 let mut fut = pin!(fut);
313 loop {
314 if let Result::Ok(output) = tokio::time::timeout(POLL_INTERVAL, fut.as_mut()).await
315 {
316 return Ok(output)
317 }
318 let mut advance = pin!(self.advance_block());
321 match select(advance.as_mut(), fut.as_mut()).await {
322 Either::Left((payload, _)) => {
323 payload?;
324 }
325 Either::Right((output, _)) => {
326 advance.await?;
327 return Ok(output)
328 }
329 }
330 }
331 };
332 tokio::time::timeout(WAIT_TIMEOUT, wait)
333 .await
334 .map_err(|_| eyre!("timed out advancing the chain until the future completed"))?
335 }
336
337 pub async fn wait_block(
341 &self,
342 number: BlockNumber,
343 expected_block_hash: BlockHash,
344 wait_finish_checkpoint: bool,
345 ) -> eyre::Result<()> {
346 let provider = &self.inner.provider;
347 poll_until(format!("block {number}"), move || async move {
348 if wait_finish_checkpoint &&
349 provider
350 .get_stage_checkpoint(StageId::Finish)?
351 .is_none_or(|checkpoint| checkpoint.block_number < number)
352 {
353 return Ok(None)
354 }
355 let Some(header) = provider.header_by_number(number)? else {
356 ensure!(
357 !wait_finish_checkpoint,
358 "Finish checkpoint matches, but could not fetch block {number}"
359 );
360 return Ok(None)
361 };
362 let hash = header.hash_slow();
363 ensure!(
364 hash == expected_block_hash,
365 "block {number} is {hash}, expected {expected_block_hash}"
366 );
367 Ok(Some(()))
368 })
369 .await
370 }
371
372 pub async fn wait_unwind(&self, number: BlockNumber) -> eyre::Result<()> {
376 let provider = &self.inner.provider;
377 poll_until(format!("unwind to block {number}"), move || async move {
378 let checkpoint = provider.get_stage_checkpoint(StageId::Headers)?;
379 Ok(checkpoint.is_some_and(|checkpoint| checkpoint.block_number == number).then_some(()))
380 })
381 .await
382 }
383
384 pub async fn wait_for_pool(
388 &self,
389 mut condition: impl FnMut(&Node::Pool) -> bool,
390 ) -> eyre::Result<()> {
391 let pool = &self.inner.pool;
392 poll_until("transaction pool condition", move || {
393 let ready = condition(pool);
394 async move { Ok(ready.then_some(())) }
395 })
396 .await
397 }
398
399 pub async fn wait_for_pool_head(&self, hash: B256) -> eyre::Result<()> {
414 let pool = &self.inner.pool;
415 poll_until(format!("transaction pool to process block {hash}"), move || {
416 let ready = pool.block_info().last_seen_block_hash == hash;
417 async move { Ok(ready.then_some(())) }
418 })
419 .await
420 }
421
422 pub async fn assert_new_block(
428 &mut self,
429 tip_tx_hash: B256,
430 block_hash: B256,
431 block_number: BlockNumber,
432 ) -> eyre::Result<()> {
433 let wait = async {
437 loop {
438 let notification = self
439 .canonical_stream
440 .next()
441 .await
442 .ok_or_else(|| eyre!("canonical state stream closed"))?;
443 let committed = notification.committed();
444 if let Some(block) = committed.blocks().get(&block_number) &&
445 block.hash() == block_hash
446 {
447 return eyre::Ok(Arc::clone(block))
448 }
449 }
450 };
451 let block = tokio::time::timeout(WAIT_TIMEOUT, wait)
452 .await
453 .map_err(|_| eyre!("timed out waiting for block {block_number}"))??;
454 ensure!(
455 block.body().transactions().iter().any(|tx| *tx.tx_hash() == tip_tx_hash),
456 "transaction {tip_tx_hash} is not included in block {block_number}"
457 );
458
459 let provider = &self.inner.provider;
462 poll_until(format!("block {block_number} to be committed"), move || async move {
463 let Some(latest) = provider.block_by_number_or_tag(BlockNumberOrTag::Latest)? else {
464 return Ok(None)
465 };
466 if latest.header().number() != block_number {
467 return Ok(None)
468 }
469 let hash = latest.header().hash_slow();
470 ensure!(
471 hash == block_hash,
472 "latest block {block_number} is {hash}, expected {block_hash}"
473 );
474 Ok(Some(()))
475 })
476 .await
477 }
478
479 pub fn block_hash(&self, number: u64) -> BlockHash {
481 self.inner
482 .provider
483 .sealed_header_by_number_or_tag(BlockNumberOrTag::Number(number))
484 .unwrap()
485 .unwrap()
486 .hash()
487 }
488
489 pub async fn sync_to(&self, block: BlockHash) -> eyre::Result<()> {
493 let sync = async {
494 while self
495 .inner
496 .provider
497 .sealed_header_by_id(BlockId::latest())?
498 .is_none_or(|h| h.hash() != block)
499 {
500 tokio::time::sleep(Duration::from_millis(100)).await;
501 self.update_forkchoice(block, block).await?;
502 }
503 Ok(())
504 };
505 tokio::time::timeout(WAIT_TIMEOUT, sync)
506 .await
507 .map_err(|_| eyre!("timed out syncing to block {block}"))??;
508
509 let _ = tokio::time::timeout(Duration::from_secs(1), self.wait_for_pool_head(block)).await;
513
514 Ok(())
515 }
516
517 pub async fn update_forkchoice(
519 &self,
520 current_head: B256,
521 new_head: B256,
522 ) -> eyre::Result<ForkchoiceUpdated> {
523 Ok(self
524 .inner
525 .add_ons_handle
526 .beacon_engine_handle
527 .fork_choice_updated(
528 ForkchoiceState {
529 head_block_hash: new_head,
530 safe_block_hash: current_head,
531 finalized_block_hash: current_head,
532 },
533 None,
534 )
535 .await?)
536 }
537
538 pub async fn update_optimistic_forkchoice(
541 &self,
542 hash: B256,
543 ) -> eyre::Result<ForkchoiceUpdated> {
544 self.update_forkchoice(B256::ZERO, hash).await
545 }
546
547 pub async fn submit_payload(&self, payload: Payload::BuiltPayload) -> eyre::Result<B256> {
549 let block_hash = payload.block().hash();
550 self.inner.add_ons_handle.beacon_engine_handle.new_payload(payload.into()).await?;
551
552 Ok(block_hash)
553 }
554
555 pub async fn import_payload(&self, payload: Payload::BuiltPayload) -> eyre::Result<B256> {
568 let block_hash = self.submit_payload(payload).await?;
569 let updated = self.update_forkchoice(block_hash, block_hash).await?;
570 ensure!(
571 updated.is_valid(),
572 "forkchoice update to block {block_hash} is not valid: {}",
573 updated.payload_status.status
574 );
575
576 Ok(block_hash)
577 }
578
579 pub fn rpc_url(&self) -> Url {
581 let addr = self.inner.rpc_server_handle().http_local_addr().unwrap();
582 format!("http://{addr}").parse().unwrap()
583 }
584
585 pub fn rpc_client(&self) -> Option<HttpClient> {
587 self.inner.rpc_server_handle().http_client()
588 }
589
590 pub fn rpc_provider(
596 &self,
597 ) -> FillProvider<impl TxFiller<Ethereum> + use<Node, Payload, AddOns>, RootProvider> {
598 self.rpc_provider_for::<Ethereum>()
599 }
600
601 pub fn rpc_provider_with_wallet<W>(
612 &self,
613 wallet: W,
614 ) -> FillProvider<impl TxFiller<Ethereum> + use<W, Node, Payload, AddOns>, RootProvider>
615 where
616 W: IntoWallet<Ethereum, NetworkWallet: Clone>,
617 {
618 self.rpc_provider_with_wallet_for::<Ethereum, W>(wallet)
619 }
620
621 pub fn rpc_provider_for<Net: RecommendedFillers>(
630 &self,
631 ) -> FillProvider<impl TxFiller<Net> + use<Net, Node, Payload, AddOns>, RootProvider<Net>, Net>
632 {
633 let provider = ProviderBuilder::new_with_network::<Net>().connect_http(self.rpc_url());
634 provider.client().set_poll_interval(RPC_PROVIDER_POLL_INTERVAL);
635 provider
636 }
637
638 pub fn rpc_provider_with_wallet_for<Net, W>(
648 &self,
649 wallet: W,
650 ) -> FillProvider<impl TxFiller<Net> + use<Net, W, Node, Payload, AddOns>, RootProvider<Net>, Net>
651 where
652 Net: RecommendedFillers,
653 W: IntoWallet<Net, NetworkWallet: Clone>,
654 {
655 let provider =
656 ProviderBuilder::new_with_network::<Net>().wallet(wallet).connect_http(self.rpc_url());
657 provider.client().set_poll_interval(RPC_PROVIDER_POLL_INTERVAL);
658 provider
659 }
660
661 pub fn auth_server_handle(&self) -> AuthServerHandle {
663 self.inner.auth_server_handle().clone()
664 }
665
666 pub fn to_node_client(&self) -> eyre::Result<crate::testsuite::NodeClient<Payload>> {
673 let rpc = self
674 .rpc_client()
675 .ok_or_else(|| eyre::eyre!("Failed to create HTTP RPC client for node"))?;
676 let auth = self.auth_server_handle();
677 let url = self.rpc_url();
678 let beacon_handle = self.inner.add_ons_handle.beacon_engine_handle.clone();
679
680 let mut client =
681 crate::testsuite::NodeClient::new_with_beacon_engine(rpc, auth, url, beacon_handle);
682 client.payload_builder = Some(self.inner.payload_builder_handle.clone());
683 let provider = self.inner.provider.clone();
684 client.database = Some(Arc::new(move || {
685 provider.database_provider_ro().map(|db| Box::new(db) as Box<dyn BlockNumReader>)
686 }));
687 Ok(client)
688 }
689
690 pub async fn testing_build_block_v1(
695 &self,
696 request: TestingBuildBlockRequestV1,
697 ) -> eyre::Result<ExecutionPayloadEnvelopeV5> {
698 let client =
699 self.rpc_client().ok_or_else(|| eyre::eyre!("HTTP RPC client not available"))?;
700
701 let res: ExecutionPayloadEnvelopeV5 =
702 client.request("testing_buildBlockV1", request.into_params()).await?;
703 eyre::Ok(res)
704 }
705}
706
707#[cfg(test)]
708mod tests {
709 use super::*;
710 use crate::NodeHelperType;
711 use reth_node_ethereum::{EthEngineTypes, EthereumNode};
712
713 fn assert_send<T: Send>(_: T) {}
714
715 #[expect(dead_code)]
718 fn test_helper_futures_are_send(
719 node: &mut NodeHelperType<EthereumNode>,
720 payload: <EthEngineTypes as PayloadTypes>::BuiltPayload,
721 ) {
722 assert_send(node.advance_block());
723 assert_send(node.advance_block_synced());
724 assert_send(node.inject_and_advance(Bytes::new()));
725 assert_send(node.advance_blocks(0));
726 assert_send(node.advance_until_receipt(B256::ZERO));
727 assert_send(node.advance_while(async {}));
728 assert_send(node.wait_block(0, B256::ZERO, false));
729 assert_send(node.wait_unwind(0));
730 assert_send(node.wait_for_pool(|_| true));
731 assert_send(node.wait_for_pool_head(B256::ZERO));
732 assert_send(node.assert_new_block(B256::ZERO, B256::ZERO, 0));
733 assert_send(node.sync_to(B256::ZERO));
734 assert_send(node.import_payload(payload));
735 }
736}