Skip to main content

reth_e2e_test_utils/testsuite/
mod.rs

1//! Utilities for running e2e tests against a node or a network of nodes.
2
3use crate::{
4    testsuite::actions::{Action, ActionBox},
5    NodeBuilderHelper,
6};
7use alloy_primitives::{Bytes, B256};
8use eyre::Result;
9use jsonrpsee::http_client::HttpClient;
10use reth_node_api::{EngineTypes, PayloadTypes};
11use reth_payload_builder::{PayloadBuilderHandle, PayloadId};
12use std::{collections::HashMap, marker::PhantomData};
13pub mod actions;
14pub mod setup;
15use crate::testsuite::setup::Setup;
16use alloy_provider::{Provider, ProviderBuilder};
17use alloy_rpc_types_engine::{ForkchoiceState, PayloadAttributes};
18use reth_chainspec::ChainSpec;
19use reth_engine_primitives::ConsensusEngineHandle;
20use reth_provider::{BlockNumReader, ProviderResult};
21use reth_rpc_builder::auth::AuthServerHandle;
22use std::sync::Arc;
23use url::Url;
24
25/// Client handles for both regular RPC and Engine API endpoints
26#[derive(Clone)]
27pub struct NodeClient<Payload>
28where
29    Payload: PayloadTypes,
30{
31    /// Regular JSON-RPC client
32    pub rpc: HttpClient,
33    /// Engine API client
34    pub engine: AuthServerHandle,
35    /// Beacon consensus engine handle for direct interaction with the consensus engine
36    pub beacon_engine_handle: Option<ConsensusEngineHandle<Payload>>,
37    /// Local payload builder used to wait for an in-progress build before requesting it over RPC.
38    pub(crate) payload_builder: Option<PayloadBuilderHandle<Payload>>,
39    /// Alloy provider for interacting with the node
40    provider: Arc<dyn Provider + Send + Sync>,
41    /// Opens read-only views of the node's database, if the node runs in-process.
42    pub(crate) database: Option<DatabaseOpener>,
43}
44
45impl<Payload> NodeClient<Payload>
46where
47    Payload: PayloadTypes,
48{
49    /// Instantiates a new [`NodeClient`] with the given handles and RPC URL
50    pub fn new(rpc: HttpClient, engine: AuthServerHandle, url: Url) -> Self {
51        let provider =
52            Arc::new(ProviderBuilder::new().connect_http(url)) as Arc<dyn Provider + Send + Sync>;
53        Self {
54            rpc,
55            engine,
56            beacon_engine_handle: None,
57            payload_builder: None,
58            provider,
59            database: None,
60        }
61    }
62
63    /// Instantiates a new [`NodeClient`] with the given handles, RPC URL, and beacon engine handle
64    pub fn new_with_beacon_engine(
65        rpc: HttpClient,
66        engine: AuthServerHandle,
67        url: Url,
68        beacon_engine_handle: ConsensusEngineHandle<Payload>,
69    ) -> Self {
70        let provider =
71            Arc::new(ProviderBuilder::new().connect_http(url)) as Arc<dyn Provider + Send + Sync>;
72        Self {
73            rpc,
74            engine,
75            beacon_engine_handle: Some(beacon_engine_handle),
76            payload_builder: None,
77            provider,
78            database: None,
79        }
80    }
81
82    /// Get a block by number using the alloy provider
83    pub async fn get_block_by_number(
84        &self,
85        number: alloy_eips::BlockNumberOrTag,
86    ) -> Result<Option<alloy_rpc_types_eth::Block>> {
87        self.provider
88            .get_block_by_number(number)
89            .await
90            .map_err(|e| eyre::eyre!("Failed to get block by number: {}", e))
91    }
92
93    /// Submit a raw transaction using the alloy provider.
94    pub async fn send_raw_transaction(&self, raw_tx: Bytes) -> Result<B256> {
95        let pending = self
96            .provider
97            .send_raw_transaction(&raw_tx)
98            .await
99            .map_err(|e| eyre::eyre!("Failed to send raw transaction: {}", e))?;
100        Ok(*pending.tx_hash())
101    }
102
103    /// Check if the node is ready by attempting to get the latest block
104    pub async fn is_ready(&self) -> bool {
105        self.get_block_by_number(alloy_eips::BlockNumberOrTag::Latest).await.is_ok()
106    }
107
108    /// Opens a read-only view of the node's database.
109    ///
110    /// Unlike the RPC endpoints, the view only contains blocks that were persisted to disk, not
111    /// canonical blocks that the engine still holds in memory.
112    pub fn database_provider_ro(&self) -> Result<Box<dyn BlockNumReader>> {
113        let open =
114            self.database.as_ref().ok_or_else(|| eyre::eyre!("Node database is not accessible"))?;
115        Ok(open()?)
116    }
117}
118
119impl<Payload> std::fmt::Debug for NodeClient<Payload>
120where
121    Payload: PayloadTypes,
122{
123    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
124        f.debug_struct("NodeClient")
125            .field("rpc", &self.rpc)
126            .field("engine", &self.engine)
127            .field("beacon_engine_handle", &self.beacon_engine_handle.is_some())
128            .field("provider", &"<Provider>")
129            .field("database", &self.database.is_some())
130            .finish()
131    }
132}
133
134/// Opens a read-only view of a node's database.
135pub(crate) type DatabaseOpener =
136    Arc<dyn Fn() -> ProviderResult<Box<dyn BlockNumReader>> + Send + Sync>;
137
138/// Represents complete block information.
139#[derive(Debug, Clone, Copy)]
140pub struct BlockInfo {
141    /// Hash of the block
142    pub hash: B256,
143    /// Number of the block
144    pub number: u64,
145    /// Timestamp of the block
146    pub timestamp: u64,
147}
148
149/// Per-node state tracking for multi-node environments
150#[derive(Clone)]
151pub struct NodeState<I>
152where
153    I: EngineTypes,
154{
155    /// Current block information for this node
156    pub current_block_info: Option<BlockInfo>,
157    /// Stores payload attributes indexed by block number for this node
158    pub payload_attributes: HashMap<u64, PayloadAttributes>,
159    /// Tracks the latest block header timestamp for this node
160    pub latest_header_time: u64,
161    /// Stores payload IDs returned by this node, indexed by block number
162    pub payload_id_history: HashMap<u64, PayloadId>,
163    /// Stores the next expected payload ID for this node
164    pub next_payload_id: Option<PayloadId>,
165    /// Stores the latest fork choice state for this node
166    pub latest_fork_choice_state: ForkchoiceState,
167    /// Stores the most recent built execution payload for this node
168    pub latest_payload_built: Option<PayloadAttributes>,
169    /// Stores the most recent executed payload for this node
170    pub latest_payload_executed: Option<PayloadAttributes>,
171    /// Stores the most recent built execution payload envelope for this node
172    pub latest_payload_envelope: Option<I::ExecutionPayloadEnvelopeV3>,
173    /// Fork base block number for validation (if this node is currently on a fork)
174    pub current_fork_base: Option<u64>,
175}
176
177impl<I> Default for NodeState<I>
178where
179    I: EngineTypes,
180{
181    fn default() -> Self {
182        Self {
183            current_block_info: None,
184            payload_attributes: HashMap::new(),
185            latest_header_time: 0,
186            payload_id_history: HashMap::new(),
187            next_payload_id: None,
188            latest_fork_choice_state: ForkchoiceState::default(),
189            latest_payload_built: None,
190            latest_payload_executed: None,
191            latest_payload_envelope: None,
192            current_fork_base: None,
193        }
194    }
195}
196
197impl<I> std::fmt::Debug for NodeState<I>
198where
199    I: EngineTypes,
200{
201    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
202        f.debug_struct("NodeState")
203            .field("current_block_info", &self.current_block_info)
204            .field("payload_attributes", &self.payload_attributes)
205            .field("latest_header_time", &self.latest_header_time)
206            .field("payload_id_history", &self.payload_id_history)
207            .field("next_payload_id", &self.next_payload_id)
208            .field("latest_fork_choice_state", &self.latest_fork_choice_state)
209            .field("latest_payload_built", &self.latest_payload_built)
210            .field("latest_payload_executed", &self.latest_payload_executed)
211            .field("latest_payload_envelope", &"<ExecutionPayloadEnvelopeV3>")
212            .field("current_fork_base", &self.current_fork_base)
213            .finish()
214    }
215}
216
217/// Represents a test environment.
218#[derive(Debug)]
219pub struct Environment<I>
220where
221    I: EngineTypes,
222{
223    /// Combined clients with both RPC and Engine API endpoints
224    pub node_clients: Vec<NodeClient<I>>,
225    /// Per-node state tracking
226    pub node_states: Vec<NodeState<I>>,
227    /// Tracks instance generic.
228    _phantom: PhantomData<I>,
229    /// Last producer index
230    pub last_producer_idx: Option<usize>,
231    /// Defines the increment for block timestamps (default: 2 seconds)
232    pub block_timestamp_increment: u64,
233    /// Number of slots until a block is considered safe
234    pub slots_to_safe: u64,
235    /// Number of slots until a block is considered finalized
236    pub slots_to_finalized: u64,
237    /// Registry for tagged blocks, mapping tag names to block info and node index
238    pub block_registry: HashMap<String, (BlockInfo, usize)>,
239    /// Currently active node index for backward compatibility with single-node actions
240    pub active_node_idx: usize,
241}
242
243impl<I> Default for Environment<I>
244where
245    I: EngineTypes,
246{
247    fn default() -> Self {
248        Self {
249            node_clients: vec![],
250            node_states: vec![],
251            _phantom: Default::default(),
252            last_producer_idx: None,
253            block_timestamp_increment: 2,
254            slots_to_safe: 0,
255            slots_to_finalized: 0,
256            block_registry: HashMap::new(),
257            active_node_idx: 0,
258        }
259    }
260}
261
262impl<I> Environment<I>
263where
264    I: EngineTypes,
265{
266    /// Get the number of nodes in the environment
267    pub const fn node_count(&self) -> usize {
268        self.node_clients.len()
269    }
270
271    /// Get mutable reference to a specific node's state
272    pub fn node_state_mut(&mut self, node_idx: usize) -> Result<&mut NodeState<I>, eyre::Error> {
273        let node_count = self.node_count();
274        self.node_states.get_mut(node_idx).ok_or_else(|| {
275            eyre::eyre!("Node index {} out of bounds (have {} nodes)", node_idx, node_count)
276        })
277    }
278
279    /// Get immutable reference to a specific node's state
280    pub fn node_state(&self, node_idx: usize) -> Result<&NodeState<I>, eyre::Error> {
281        self.node_states.get(node_idx).ok_or_else(|| {
282            eyre::eyre!("Node index {} out of bounds (have {} nodes)", node_idx, self.node_count())
283        })
284    }
285
286    /// Get the currently active node's state
287    pub fn active_node_state(&self) -> Result<&NodeState<I>, eyre::Error> {
288        self.node_state(self.active_node_idx)
289    }
290
291    /// Get mutable reference to the currently active node's state
292    pub fn active_node_state_mut(&mut self) -> Result<&mut NodeState<I>, eyre::Error> {
293        let idx = self.active_node_idx;
294        self.node_state_mut(idx)
295    }
296
297    /// Set the active node index
298    pub fn set_active_node(&mut self, node_idx: usize) -> Result<(), eyre::Error> {
299        if node_idx >= self.node_count() {
300            return Err(eyre::eyre!(
301                "Node index {} out of bounds (have {} nodes)",
302                node_idx,
303                self.node_count()
304            ));
305        }
306        self.active_node_idx = node_idx;
307        Ok(())
308    }
309
310    /// Initialize node states when nodes are created
311    pub fn initialize_node_states(&mut self, node_count: usize) {
312        self.node_states = (0..node_count).map(|_| NodeState::default()).collect();
313    }
314
315    /// Get current block info from active node
316    pub fn current_block_info(&self) -> Option<BlockInfo> {
317        self.active_node_state().ok()?.current_block_info
318    }
319
320    /// Set current block info on active node
321    pub fn set_current_block_info(&mut self, block_info: BlockInfo) -> Result<(), eyre::Error> {
322        self.active_node_state_mut()?.current_block_info = Some(block_info);
323        Ok(())
324    }
325}
326
327/// Builder for creating test scenarios
328#[expect(missing_debug_implementations)]
329pub struct TestBuilder<I>
330where
331    I: EngineTypes,
332{
333    setup: Option<Setup<I>>,
334    actions: Vec<ActionBox<I>>,
335    env: Environment<I>,
336}
337
338impl<I> Default for TestBuilder<I>
339where
340    I: EngineTypes,
341{
342    fn default() -> Self {
343        Self { setup: None, actions: Vec::new(), env: Default::default() }
344    }
345}
346
347impl<I> TestBuilder<I>
348where
349    I: EngineTypes + 'static,
350{
351    /// Create a new test builder
352    pub fn new() -> Self {
353        Self::default()
354    }
355
356    /// Set the test setup
357    pub fn with_setup(mut self, setup: Setup<I>) -> Self {
358        self.setup = Some(setup);
359        self
360    }
361
362    /// Set the test setup with chain import from RLP file
363    pub fn with_setup_and_import(
364        mut self,
365        mut setup: Setup<I>,
366        rlp_path: impl Into<std::path::PathBuf>,
367    ) -> Self {
368        setup.import_rlp_path = Some(rlp_path.into());
369        self.setup = Some(setup);
370        self
371    }
372
373    /// Add an action to the test
374    pub fn with_action<A>(mut self, action: A) -> Self
375    where
376        A: Action<I>,
377    {
378        self.actions.push(ActionBox::<I>::new(action));
379        self
380    }
381
382    /// Add multiple actions to the test
383    pub fn with_actions<II, A>(mut self, actions: II) -> Self
384    where
385        II: IntoIterator<Item = A>,
386        A: Action<I>,
387    {
388        self.actions.extend(actions.into_iter().map(ActionBox::new));
389        self
390    }
391
392    /// Run the test scenario
393    pub async fn run<N>(mut self) -> Result<()>
394    where
395        N: NodeBuilderHelper<Payload = I, ChainSpec: From<ChainSpec>>,
396    {
397        let mut setup = self.setup.take();
398
399        if let Some(ref mut s) = setup {
400            s.apply::<N>(&mut self.env).await?;
401        }
402
403        let actions = std::mem::take(&mut self.actions);
404
405        for action in actions {
406            action.execute(&mut self.env).await?;
407        }
408
409        // explicitly drop the setup to shutdown the nodes
410        // after all actions have completed
411        drop(setup);
412
413        Ok(())
414    }
415}