reth_e2e_test_utils/testsuite/
mod.rs1use 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#[derive(Clone)]
27pub struct NodeClient<Payload>
28where
29 Payload: PayloadTypes,
30{
31 pub rpc: HttpClient,
33 pub engine: AuthServerHandle,
35 pub beacon_engine_handle: Option<ConsensusEngineHandle<Payload>>,
37 pub(crate) payload_builder: Option<PayloadBuilderHandle<Payload>>,
39 provider: Arc<dyn Provider + Send + Sync>,
41 pub(crate) database: Option<DatabaseOpener>,
43}
44
45impl<Payload> NodeClient<Payload>
46where
47 Payload: PayloadTypes,
48{
49 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 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 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 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 pub async fn is_ready(&self) -> bool {
105 self.get_block_by_number(alloy_eips::BlockNumberOrTag::Latest).await.is_ok()
106 }
107
108 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
134pub(crate) type DatabaseOpener =
136 Arc<dyn Fn() -> ProviderResult<Box<dyn BlockNumReader>> + Send + Sync>;
137
138#[derive(Debug, Clone, Copy)]
140pub struct BlockInfo {
141 pub hash: B256,
143 pub number: u64,
145 pub timestamp: u64,
147}
148
149#[derive(Clone)]
151pub struct NodeState<I>
152where
153 I: EngineTypes,
154{
155 pub current_block_info: Option<BlockInfo>,
157 pub payload_attributes: HashMap<u64, PayloadAttributes>,
159 pub latest_header_time: u64,
161 pub payload_id_history: HashMap<u64, PayloadId>,
163 pub next_payload_id: Option<PayloadId>,
165 pub latest_fork_choice_state: ForkchoiceState,
167 pub latest_payload_built: Option<PayloadAttributes>,
169 pub latest_payload_executed: Option<PayloadAttributes>,
171 pub latest_payload_envelope: Option<I::ExecutionPayloadEnvelopeV3>,
173 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#[derive(Debug)]
219pub struct Environment<I>
220where
221 I: EngineTypes,
222{
223 pub node_clients: Vec<NodeClient<I>>,
225 pub node_states: Vec<NodeState<I>>,
227 _phantom: PhantomData<I>,
229 pub last_producer_idx: Option<usize>,
231 pub block_timestamp_increment: u64,
233 pub slots_to_safe: u64,
235 pub slots_to_finalized: u64,
237 pub block_registry: HashMap<String, (BlockInfo, usize)>,
239 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 pub const fn node_count(&self) -> usize {
268 self.node_clients.len()
269 }
270
271 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 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 pub fn active_node_state(&self) -> Result<&NodeState<I>, eyre::Error> {
288 self.node_state(self.active_node_idx)
289 }
290
291 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 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 pub fn initialize_node_states(&mut self, node_count: usize) {
312 self.node_states = (0..node_count).map(|_| NodeState::default()).collect();
313 }
314
315 pub fn current_block_info(&self) -> Option<BlockInfo> {
317 self.active_node_state().ok()?.current_block_info
318 }
319
320 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#[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 pub fn new() -> Self {
353 Self::default()
354 }
355
356 pub fn with_setup(mut self, setup: Setup<I>) -> Self {
358 self.setup = Some(setup);
359 self
360 }
361
362 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 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 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 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 drop(setup);
412
413 Ok(())
414 }
415}