reth_e2e_test_utils/testsuite/
setup.rs1use crate::{testsuite::Environment, E2ETestSetupExt, NodeBuilderHelper};
4use alloy_eips::BlockNumberOrTag;
5use alloy_rpc_types_engine::ForkchoiceState;
6use eyre::{eyre, Result};
7use reth_chainspec::ChainSpec;
8use reth_ethereum_primitives::Block;
9use reth_node_api::{EngineTypes, PayloadTypes, TreeConfig};
10use reth_node_core::{args::StorageArgs, primitives::RecoveredBlock};
11use revm::state::EvmState;
12use std::{marker::PhantomData, path::Path, sync::Arc};
13use tokio::{
14 sync::mpsc,
15 time::{sleep, Duration},
16};
17use tracing::debug;
18
19#[derive(Debug)]
21pub struct Setup<I> {
22 pub chain_spec: Option<Arc<ChainSpec>>,
24 pub genesis: Option<Genesis>,
26 pub blocks: Vec<RecoveredBlock<Block>>,
28 pub state: Option<EvmState>,
30 pub network: NetworkSetup,
32 pub tree_config: TreeConfig,
34 shutdown_tx: Option<mpsc::Sender<()>>,
36 pub is_dev: bool,
38 pub storage_v2: bool,
42 _phantom: PhantomData<I>,
44 import_result_holder: Option<crate::setup_import::ChainImportResult>,
47 pub import_rlp_path: Option<std::path::PathBuf>,
49}
50
51impl<I> Default for Setup<I> {
52 fn default() -> Self {
53 Self {
54 chain_spec: None,
55 genesis: None,
56 blocks: Vec::new(),
57 state: None,
58 network: NetworkSetup::default(),
59 tree_config: TreeConfig::default(),
60 shutdown_tx: None,
61 is_dev: true,
62 storage_v2: StorageArgs::default().v2,
63 _phantom: Default::default(),
64 import_result_holder: None,
65 import_rlp_path: None,
66 }
67 }
68}
69
70impl<I> Drop for Setup<I> {
71 fn drop(&mut self) {
72 if let Some(tx) = self.shutdown_tx.take() {
74 let _ = tx.try_send(());
75 }
76 }
77}
78
79impl<I> Setup<I>
80where
81 I: EngineTypes,
82{
83 pub fn with_chain_spec(mut self, chain_spec: Arc<ChainSpec>) -> Self {
85 self.chain_spec = Some(chain_spec);
86 self
87 }
88
89 pub const fn with_genesis(mut self, genesis: Genesis) -> Self {
91 self.genesis = Some(genesis);
92 self
93 }
94
95 pub fn with_block(mut self, block: RecoveredBlock<Block>) -> Self {
97 self.blocks.push(block);
98 self
99 }
100
101 pub fn with_blocks(mut self, blocks: Vec<RecoveredBlock<Block>>) -> Self {
103 self.blocks.extend(blocks);
104 self
105 }
106
107 pub fn with_state(mut self, state: EvmState) -> Self {
109 self.state = Some(state);
110 self
111 }
112
113 pub const fn with_network(mut self, network: NetworkSetup) -> Self {
115 self.network = network;
116 self
117 }
118
119 pub const fn with_dev_mode(mut self, is_dev: bool) -> Self {
121 self.is_dev = is_dev;
122 self
123 }
124
125 pub const fn with_tree_config(mut self, tree_config: TreeConfig) -> Self {
127 self.tree_config = tree_config;
128 self
129 }
130
131 pub const fn with_storage_v2(mut self, storage_v2: bool) -> Self {
133 self.storage_v2 = storage_v2;
134 self
135 }
136
137 pub async fn apply_with_import<N>(
139 &mut self,
140 env: &mut Environment<I>,
141 rlp_path: &Path,
142 ) -> Result<()>
143 where
144 N: NodeBuilderHelper<Payload = I>,
145 {
146 Box::pin(self.apply_with_import_(env, rlp_path)).await
148 }
149
150 async fn apply_with_import_(
152 &mut self,
153 env: &mut Environment<I>,
154 rlp_path: &Path,
155 ) -> Result<()> {
156 let import_result = self.create_nodes_with_import(rlp_path).await?;
158
159 let mut node_clients = Vec::new();
161 let nodes = &import_result.nodes;
162 for node in nodes {
163 let rpc = node
164 .rpc_client()
165 .ok_or_else(|| eyre!("Failed to create HTTP RPC client for node"))?;
166 let auth = node.auth_server_handle();
167 let url = node.rpc_url();
168 node_clients.push(crate::testsuite::NodeClient::new(rpc, auth, url));
170 }
171
172 self.import_result_holder = Some(import_result);
175
176 self.finalize_setup(env, node_clients, true).await
178 }
179
180 pub async fn apply<N>(&mut self, env: &mut Environment<I>) -> Result<()>
182 where
183 N: NodeBuilderHelper<Payload = I, ChainSpec: From<ChainSpec>>,
184 {
185 Box::pin(self.apply_::<N>(env)).await
187 }
188
189 async fn apply_<N>(&mut self, env: &mut Environment<I>) -> Result<()>
191 where
192 N: NodeBuilderHelper<Payload = I, ChainSpec: From<ChainSpec>>,
193 {
194 if let Some(rlp_path) = self.import_rlp_path.take() {
196 return self.apply_with_import::<N>(env, &rlp_path).await;
197 }
198 let chain_spec =
199 self.chain_spec.clone().ok_or_else(|| eyre!("Chain specification is required"))?;
200
201 let (shutdown_tx, mut shutdown_rx) = mpsc::channel(1);
202 self.shutdown_tx = Some(shutdown_tx);
203
204 let tree_config = self.tree_config.clone();
205
206 let result = N::test_setup(
207 self.network.node_count,
208 Arc::<N::ChainSpec>::new((*chain_spec).clone().into()),
209 )
210 .with_tree_config_modifier(move |base| {
211 tree_config.clone().with_cross_block_cache_size(base.cross_block_cache_size())
212 })
213 .with_dev_mode(self.is_dev)
214 .with_storage_v2(self.storage_v2)
215 .with_connect_nodes(self.network.connect_nodes)
216 .build()
217 .await;
218
219 let mut node_clients = Vec::new();
220 match result {
221 Ok((nodes, _wallet)) => {
222 for node in &nodes {
224 node_clients.push(node.to_node_client()?);
225 }
226
227 tokio::spawn(async move {
229 let _nodes = nodes;
231 let _ = shutdown_rx.recv().await;
233 });
235 }
236 Err(e) => {
237 return Err(eyre!("Failed to setup nodes: {}", e));
238 }
239 }
240
241 self.finalize_setup(env, node_clients, false).await
243 }
244
245 async fn create_nodes_with_import(
251 &self,
252 rlp_path: &Path,
253 ) -> Result<crate::setup_import::ChainImportResult> {
254 let chain_spec =
255 self.chain_spec.clone().ok_or_else(|| eyre!("Chain specification is required"))?;
256
257 crate::setup_import::setup_engine_with_chain_import(
258 self.network.node_count,
259 chain_spec,
260 self.is_dev,
261 self.storage_v2,
262 self.tree_config.clone(),
263 rlp_path,
264 )
265 .await
266 }
267
268 async fn finalize_setup(
270 &self,
271 env: &mut Environment<I>,
272 node_clients: Vec<crate::testsuite::NodeClient<I>>,
273 use_latest_block: bool,
274 ) -> Result<()> {
275 if node_clients.is_empty() {
276 return Err(eyre!("No nodes were created"));
277 }
278
279 self.wait_for_nodes_ready(&node_clients).await?;
281
282 env.node_clients = node_clients;
283 env.initialize_node_states(self.network.node_count);
284
285 let (initial_block_info, genesis_block_info) = if use_latest_block {
287 let latest =
289 self.get_block_info(&env.node_clients[0], BlockNumberOrTag::Latest).await?;
290 let genesis =
291 self.get_block_info(&env.node_clients[0], BlockNumberOrTag::Number(0)).await?;
292 (latest, genesis)
293 } else {
294 let genesis =
296 self.get_block_info(&env.node_clients[0], BlockNumberOrTag::Number(0)).await?;
297 (genesis, genesis)
298 };
299
300 for (node_idx, node_state) in env.node_states.iter_mut().enumerate() {
302 node_state.current_block_info = Some(initial_block_info);
303 node_state.latest_header_time = initial_block_info.timestamp;
304 node_state.latest_fork_choice_state = ForkchoiceState {
305 head_block_hash: initial_block_info.hash,
306 safe_block_hash: initial_block_info.hash,
307 finalized_block_hash: genesis_block_info.hash,
308 };
309
310 debug!(
311 "Node {} initialized with block {} (hash: {})",
312 node_idx, initial_block_info.number, initial_block_info.hash
313 );
314 }
315
316 debug!(
317 "Environment initialized with {} nodes, starting from block {} (hash: {})",
318 self.network.node_count, initial_block_info.number, initial_block_info.hash
319 );
320
321 Ok(())
322 }
323
324 async fn wait_for_nodes_ready<P>(
326 &self,
327 node_clients: &[crate::testsuite::NodeClient<P>],
328 ) -> Result<()>
329 where
330 P: PayloadTypes,
331 {
332 for (idx, client) in node_clients.iter().enumerate() {
333 let mut retry_count = 0;
334 const MAX_RETRIES: usize = 10;
335
336 while retry_count < MAX_RETRIES {
337 if client.is_ready().await {
338 debug!("Node {idx} RPC endpoint is ready");
339 break;
340 }
341
342 retry_count += 1;
343 debug!("Node {idx} RPC endpoint not ready, retry {retry_count}/{MAX_RETRIES}");
344 sleep(Duration::from_millis(500)).await;
345 }
346
347 if retry_count == MAX_RETRIES {
348 return Err(eyre!(
349 "Failed to connect to node {idx} RPC endpoint after {MAX_RETRIES} retries"
350 ));
351 }
352 }
353 Ok(())
354 }
355
356 async fn get_block_info<P>(
358 &self,
359 client: &crate::testsuite::NodeClient<P>,
360 block: BlockNumberOrTag,
361 ) -> Result<crate::testsuite::BlockInfo>
362 where
363 P: PayloadTypes,
364 {
365 let block = client
366 .get_block_by_number(block)
367 .await?
368 .ok_or_else(|| eyre!("Block {:?} not found", block))?;
369
370 Ok(crate::testsuite::BlockInfo {
371 hash: block.header.hash,
372 number: block.header.number,
373 timestamp: block.header.timestamp,
374 })
375 }
376}
377
378#[derive(Debug)]
380pub struct Genesis {}
381
382#[derive(Debug, Default)]
384pub struct NetworkSetup {
385 pub node_count: usize,
387 pub connect_nodes: bool,
389}
390
391impl NetworkSetup {
392 pub const fn single_node() -> Self {
394 Self { node_count: 1, connect_nodes: true }
395 }
396
397 pub const fn multi_node(count: usize) -> Self {
399 Self { node_count: count, connect_nodes: true }
400 }
401
402 pub const fn multi_node_unconnected(count: usize) -> Self {
404 Self { node_count: count, connect_nodes: false }
405 }
406}