Skip to main content

reth_e2e_test_utils/testsuite/actions/
produce_blocks.rs

1//! Block production actions for the e2e testing framework.
2
3use crate::testsuite::{
4    actions::{expect_fcu_not_syncing_or_accepted, validate_fcu_response, Action, Sequence},
5    BlockInfo, Environment,
6};
7use alloy_primitives::{Bytes, B256};
8use alloy_rpc_types_engine::{
9    payload::ExecutionPayloadEnvelopeV3, ForkchoiceState, PayloadAttributes, PayloadStatusEnum,
10};
11use alloy_rpc_types_eth::{Block, Header, Receipt, Transaction, TransactionRequest};
12use eyre::Result;
13use futures_util::future::BoxFuture;
14use reth_ethereum_primitives::TransactionSigned;
15use reth_node_api::{EngineTypes, PayloadKind, PayloadTypes};
16use reth_rpc_api::clients::{EngineApiClient, EthApiClient};
17use std::{collections::HashSet, marker::PhantomData, time::Duration};
18use tokio::time::sleep;
19use tracing::debug;
20
21/// Mine a single block with the given transactions and verify the block was created
22/// successfully.
23#[derive(Debug)]
24pub struct AssertMineBlock<Engine>
25where
26    Engine: PayloadTypes,
27{
28    /// The node index to mine
29    pub node_idx: usize,
30    /// Transactions to include in the block
31    pub transactions: Vec<Bytes>,
32    /// Expected block hash (optional)
33    pub expected_hash: Option<B256>,
34    /// Block's payload attributes
35    // TODO: refactor once we have actions to generate payload attributes.
36    pub payload_attributes: Engine::PayloadAttributes,
37    /// Tracks engine type
38    _phantom: PhantomData<Engine>,
39}
40
41impl<Engine> AssertMineBlock<Engine>
42where
43    Engine: PayloadTypes,
44{
45    /// Create a new `AssertMineBlock` action
46    pub fn new(
47        node_idx: usize,
48        transactions: Vec<Bytes>,
49        expected_hash: Option<B256>,
50        payload_attributes: Engine::PayloadAttributes,
51    ) -> Self {
52        Self {
53            node_idx,
54            transactions,
55            expected_hash,
56            payload_attributes,
57            _phantom: Default::default(),
58        }
59    }
60}
61
62impl<Engine> Action<Engine> for AssertMineBlock<Engine>
63where
64    Engine: EngineTypes,
65{
66    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
67        Box::pin(async move {
68            if self.node_idx >= env.node_clients.len() {
69                return Err(eyre::eyre!("Node index out of bounds: {}", self.node_idx));
70            }
71
72            let node_client = &env.node_clients[self.node_idx];
73            let rpc_client = &node_client.rpc;
74            let engine_client = node_client.engine.http_client();
75
76            // get the latest block to use as parent
77            let latest_block = EthApiClient::<
78                TransactionRequest,
79                Transaction,
80                Block,
81                Receipt,
82                Header,
83                TransactionSigned,
84            >::block_by_number(
85                rpc_client, alloy_eips::BlockNumberOrTag::Latest, false
86            )
87            .await?;
88
89            let latest_block = latest_block.ok_or_else(|| eyre::eyre!("Latest block not found"))?;
90            let parent_hash = latest_block.header.hash;
91
92            debug!("Latest block hash: {parent_hash}");
93
94            // create a simple forkchoice state with the latest block as head
95            let fork_choice_state = ForkchoiceState::same_hash(parent_hash);
96
97            // Try v2 first for backwards compatibility, fall back to v3 on error.
98            match EngineApiClient::<Engine>::fork_choice_updated_v2(
99                &engine_client,
100                fork_choice_state,
101                Some(self.payload_attributes.clone()),
102            )
103            .await
104            {
105                Ok(fcu_result) => {
106                    debug!(?fcu_result, "FCU v2 result");
107                    match fcu_result.payload_status.status {
108                        PayloadStatusEnum::Valid => {
109                            if let Some(payload_id) = fcu_result.payload_id {
110                                debug!(id=%payload_id, "Got payload");
111                                let _engine_payload = EngineApiClient::<Engine>::get_payload_v2(
112                                    &engine_client,
113                                    payload_id,
114                                )
115                                .await?;
116                                Ok(())
117                            } else {
118                                Err(eyre::eyre!("No payload ID returned from forkchoiceUpdated"))
119                            }
120                        }
121                        _ => Err(eyre::eyre!(
122                            "Payload status not valid: {:?}",
123                            fcu_result.payload_status
124                        ))?,
125                    }
126                }
127                Err(_) => {
128                    // If v2 fails due to unsupported fork/missing fields, try v3
129                    let fcu_result = EngineApiClient::<Engine>::fork_choice_updated_v3(
130                        &engine_client,
131                        fork_choice_state,
132                        Some(self.payload_attributes.clone()),
133                    )
134                    .await?;
135
136                    debug!(?fcu_result, "FCU v3 result");
137                    match fcu_result.payload_status.status {
138                        PayloadStatusEnum::Valid => {
139                            if let Some(payload_id) = fcu_result.payload_id {
140                                debug!(id=%payload_id, "Got payload");
141                                let _engine_payload = EngineApiClient::<Engine>::get_payload_v3(
142                                    &engine_client,
143                                    payload_id,
144                                )
145                                .await?;
146                                Ok(())
147                            } else {
148                                Err(eyre::eyre!("No payload ID returned from forkchoiceUpdated"))
149                            }
150                        }
151                        _ => Err(eyre::eyre!(
152                            "Payload status not valid: {:?}",
153                            fcu_result.payload_status
154                        )),
155                    }
156                }
157            }
158        })
159    }
160}
161
162/// Pick the next block producer based on the latest block information.
163#[derive(Debug, Default)]
164pub struct PickNextBlockProducer {}
165
166impl PickNextBlockProducer {
167    /// Create a new `PickNextBlockProducer` action
168    pub const fn new() -> Self {
169        Self {}
170    }
171}
172
173impl<Engine> Action<Engine> for PickNextBlockProducer
174where
175    Engine: EngineTypes,
176{
177    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
178        Box::pin(async move {
179            let num_clients = env.node_clients.len();
180            if num_clients == 0 {
181                return Err(eyre::eyre!("No node clients available"));
182            }
183
184            let latest_info = env
185                .current_block_info()
186                .ok_or_else(|| eyre::eyre!("No latest block information available"))?;
187
188            // simple round-robin selection based on next block number
189            let next_producer_idx = ((latest_info.number + 1) % num_clients as u64) as usize;
190
191            env.last_producer_idx = Some(next_producer_idx);
192            debug!(
193                "Selected node {} as the next block producer for block {}",
194                next_producer_idx,
195                latest_info.number + 1
196            );
197
198            Ok(())
199        })
200    }
201}
202
203/// Store payload attributes for the next block.
204#[derive(Debug, Default)]
205pub struct GeneratePayloadAttributes {}
206
207impl<Engine> Action<Engine> for GeneratePayloadAttributes
208where
209    Engine: EngineTypes + PayloadTypes,
210    Engine::PayloadAttributes: From<PayloadAttributes>,
211{
212    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
213        Box::pin(async move {
214            let latest_block = env
215                .current_block_info()
216                .ok_or_else(|| eyre::eyre!("No latest block information available"))?;
217            let block_number = latest_block.number;
218            let timestamp =
219                env.active_node_state()?.latest_header_time + env.block_timestamp_increment;
220            let payload_attributes = PayloadAttributes {
221                timestamp,
222                prev_randao: B256::random(),
223                suggested_fee_recipient: alloy_primitives::Address::random(),
224                withdrawals: Some(vec![]),
225                parent_beacon_block_root: Some(B256::ZERO),
226                slot_number: None,
227                ..Default::default()
228            };
229
230            env.active_node_state_mut()?
231                .payload_attributes
232                .insert(latest_block.number + 1, payload_attributes);
233            debug!("Stored payload attributes for block {}", block_number + 1);
234            Ok(())
235        })
236    }
237}
238
239/// Action that generates the next payload
240#[derive(Debug, Default)]
241pub struct GenerateNextPayload {}
242
243impl<Engine> Action<Engine> for GenerateNextPayload
244where
245    Engine: EngineTypes + PayloadTypes,
246    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
247{
248    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
249        Box::pin(async move {
250            let latest_block = env
251                .current_block_info()
252                .ok_or_else(|| eyre::eyre!("No latest block information available"))?;
253
254            let parent_hash = latest_block.hash;
255            debug!("Latest block hash: {parent_hash}");
256
257            let fork_choice_state = ForkchoiceState::same_hash(parent_hash);
258
259            let payload_attributes = env
260                .active_node_state()?
261                .payload_attributes
262                .get(&(latest_block.number + 1))
263                .cloned()
264                .ok_or_else(|| eyre::eyre!("No payload attributes found for next block"))?;
265
266            let producer_idx =
267                env.last_producer_idx.ok_or_else(|| eyre::eyre!("No block producer selected"))?;
268
269            let fcu_result = EngineApiClient::<Engine>::fork_choice_updated_v3(
270                &env.node_clients[producer_idx].engine.http_client(),
271                fork_choice_state,
272                Some(payload_attributes.clone().into()),
273            )
274            .await?;
275
276            debug!("FCU result: {:?}", fcu_result);
277
278            // validate the FCU status before proceeding
279            // Note: In the context of GenerateNextPayload, Syncing usually means the engine
280            // doesn't have the requested head block, which should be an error
281            expect_fcu_not_syncing_or_accepted(&fcu_result, "GenerateNextPayload")?;
282
283            let payload_id = if let Some(payload_id) = fcu_result.payload_id {
284                debug!("Received new payload ID: {:?}", payload_id);
285                payload_id
286            } else {
287                debug!("No payload ID returned, generating fresh payload attributes for forking");
288
289                let fresh_payload_attributes = PayloadAttributes {
290                    timestamp: env.active_node_state()?.latest_header_time +
291                        env.block_timestamp_increment,
292                    prev_randao: B256::random(),
293                    suggested_fee_recipient: alloy_primitives::Address::random(),
294                    withdrawals: Some(vec![]),
295                    parent_beacon_block_root: Some(B256::ZERO),
296                    slot_number: None,
297                    ..Default::default()
298                };
299
300                let fresh_fcu_result = EngineApiClient::<Engine>::fork_choice_updated_v3(
301                    &env.node_clients[producer_idx].engine.http_client(),
302                    fork_choice_state,
303                    Some(fresh_payload_attributes.clone().into()),
304                )
305                .await?;
306
307                debug!("Fresh FCU result: {:?}", fresh_fcu_result);
308
309                // validate the fresh FCU status
310                expect_fcu_not_syncing_or_accepted(
311                    &fresh_fcu_result,
312                    "GenerateNextPayload (fresh)",
313                )?;
314
315                if let Some(payload_id) = fresh_fcu_result.payload_id {
316                    payload_id
317                } else {
318                    debug!("Engine considers the fork base already canonical, skipping payload generation");
319                    return Ok(());
320                }
321            };
322
323            env.active_node_state_mut()?.next_payload_id = Some(payload_id);
324
325            if let Some(builder) = &env.node_clients[producer_idx].payload_builder {
326                // Wait for the pending build rather than racing it with an empty fallback payload.
327                tokio::time::timeout(
328                    Duration::from_secs(30),
329                    builder.resolve_kind(payload_id, PayloadKind::WaitForPending),
330                )
331                .await?
332                .ok_or_else(|| eyre::eyre!("Unknown payload {payload_id}"))??;
333            } else {
334                // RPC-only clients do not expose the local payload builder.
335                sleep(Duration::from_secs(1)).await;
336            }
337
338            let built_payload_envelope = EngineApiClient::<Engine>::get_payload_v3(
339                &env.node_clients[producer_idx].engine.http_client(),
340                payload_id,
341            )
342            .await?;
343
344            // Store the payload attributes that were used to generate this payload
345            let built_payload = payload_attributes.clone();
346            env.active_node_state_mut()?
347                .payload_id_history
348                .insert(latest_block.number + 1, payload_id);
349            env.active_node_state_mut()?.latest_payload_built = Some(built_payload);
350            env.active_node_state_mut()?.latest_payload_envelope = Some(built_payload_envelope);
351
352            Ok(())
353        })
354    }
355}
356
357/// Action that broadcasts the latest fork choice state to all clients
358#[derive(Debug, Default)]
359pub struct BroadcastLatestForkchoice {}
360
361impl<Engine> Action<Engine> for BroadcastLatestForkchoice
362where
363    Engine: EngineTypes + PayloadTypes,
364    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
365    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
366{
367    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
368        Box::pin(async move {
369            if env.node_clients.is_empty() {
370                return Err(eyre::eyre!("No node clients available"));
371            }
372
373            // use the hash of the newly executed payload if available
374            let head_hash = if let Some(payload_envelope) =
375                &env.active_node_state()?.latest_payload_envelope
376            {
377                let execution_payload_envelope: ExecutionPayloadEnvelopeV3 =
378                    payload_envelope.clone().into();
379                let new_block_hash = execution_payload_envelope
380                    .execution_payload
381                    .payload_inner
382                    .payload_inner
383                    .block_hash;
384                debug!("Using newly executed block hash as head: {new_block_hash}");
385                new_block_hash
386            } else {
387                // fallback to RPC query
388                let rpc_client = &env.node_clients[0].rpc;
389                let current_head_block = EthApiClient::<
390                    TransactionRequest,
391                    Transaction,
392                    Block,
393                    Receipt,
394                    Header,
395                    TransactionSigned,
396                >::block_by_number(
397                    rpc_client, alloy_eips::BlockNumberOrTag::Latest, false
398                )
399                .await?
400                .ok_or_else(|| eyre::eyre!("No latest block found from RPC"))?;
401                debug!("Using RPC latest block hash as head: {}", current_head_block.header.hash);
402                current_head_block.header.hash
403            };
404
405            broadcast_forkchoice(env, head_hash).await
406        })
407    }
408}
409
410/// Action that syncs environment state with the node's canonical chain via RPC.
411///
412/// This queries the latest canonical block from the node and updates the environment
413/// to match. Typically used after forkchoice operations to ensure the environment
414/// is in sync with the node's view of the canonical chain.
415#[derive(Debug, Default)]
416pub struct UpdateBlockInfo {}
417
418impl<Engine> Action<Engine> for UpdateBlockInfo
419where
420    Engine: EngineTypes,
421{
422    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
423        Box::pin(async move {
424            // get the latest block from the first client to update environment state
425            let rpc_client = &env.node_clients[0].rpc;
426            let latest_block = EthApiClient::<
427                TransactionRequest,
428                Transaction,
429                Block,
430                Receipt,
431                Header,
432                TransactionSigned,
433            >::block_by_number(
434                rpc_client, alloy_eips::BlockNumberOrTag::Latest, false
435            )
436            .await?
437            .ok_or_else(|| eyre::eyre!("No latest block found from RPC"))?;
438
439            // update environment with the new block information
440            env.set_current_block_info(BlockInfo {
441                hash: latest_block.header.hash,
442                number: latest_block.header.number,
443                timestamp: latest_block.header.timestamp,
444            })?;
445
446            env.active_node_state_mut()?.latest_header_time = latest_block.header.timestamp;
447            env.active_node_state_mut()?.latest_fork_choice_state.head_block_hash =
448                latest_block.header.hash;
449
450            debug!(
451                "Updated environment to block {} (hash: {})",
452                latest_block.header.number, latest_block.header.hash
453            );
454
455            Ok(())
456        })
457    }
458}
459
460/// Action that updates environment state using the locally produced payload.
461///
462/// This uses the execution payload stored in the environment rather than querying RPC,
463/// making it more efficient and reliable during block production. Preferred over
464/// `UpdateBlockInfo` when we have just produced a block and have the payload available.
465#[derive(Debug, Default)]
466pub struct UpdateBlockInfoToLatestPayload {}
467
468impl<Engine> Action<Engine> for UpdateBlockInfoToLatestPayload
469where
470    Engine: EngineTypes + PayloadTypes,
471    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
472{
473    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
474        Box::pin(async move {
475            let payload_envelope = env
476                .active_node_state()?
477                .latest_payload_envelope
478                .as_ref()
479                .ok_or_else(|| eyre::eyre!("No execution payload envelope available"))?;
480
481            let execution_payload_envelope: ExecutionPayloadEnvelopeV3 =
482                payload_envelope.clone().into();
483            let execution_payload = execution_payload_envelope.execution_payload;
484
485            let block_hash = execution_payload.payload_inner.payload_inner.block_hash;
486            let block_number = execution_payload.payload_inner.payload_inner.block_number;
487            let block_timestamp = execution_payload.payload_inner.payload_inner.timestamp;
488
489            // update environment with the new block information from the payload
490            env.set_current_block_info(BlockInfo {
491                hash: block_hash,
492                number: block_number,
493                timestamp: block_timestamp,
494            })?;
495
496            env.active_node_state_mut()?.latest_header_time = block_timestamp;
497            env.active_node_state_mut()?.latest_fork_choice_state.head_block_hash = block_hash;
498
499            debug!(
500                "Updated environment to newly produced block {} (hash: {})",
501                block_number, block_hash
502            );
503
504            Ok(())
505        })
506    }
507}
508
509/// Action that checks whether the broadcasted new payload has been accepted
510#[derive(Debug, Default)]
511pub struct CheckPayloadAccepted {}
512
513impl<Engine> Action<Engine> for CheckPayloadAccepted
514where
515    Engine: EngineTypes,
516    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
517{
518    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
519        Box::pin(async move {
520            let mut accepted_check: bool = false;
521
522            let latest_block = env
523                .current_block_info()
524                .ok_or_else(|| eyre::eyre!("No latest block information available"))?;
525
526            let payload_id = *env
527                .active_node_state()?
528                .payload_id_history
529                .get(&(latest_block.number + 1))
530                .ok_or_else(|| eyre::eyre!("Cannot find payload_id"))?;
531
532            let node_clients = env.node_clients.clone();
533            for (idx, client) in node_clients.iter().enumerate() {
534                let rpc_client = &client.rpc;
535
536                // get the last header by number using latest_head_number
537                let rpc_latest_header = EthApiClient::<
538                    TransactionRequest,
539                    Transaction,
540                    Block,
541                    Receipt,
542                    Header,
543                    TransactionSigned,
544                >::header_by_number(
545                    rpc_client, alloy_eips::BlockNumberOrTag::Latest
546                )
547                .await?
548                .ok_or_else(|| eyre::eyre!("No latest header found from rpc"))?;
549
550                // perform several checks
551                let next_new_payload = env
552                    .active_node_state()?
553                    .latest_payload_built
554                    .as_ref()
555                    .ok_or_else(|| eyre::eyre!("No next built payload found"))?;
556
557                let built_payload = EngineApiClient::<Engine>::get_payload_v3(
558                    &client.engine.http_client(),
559                    payload_id,
560                )
561                .await?;
562
563                let execution_payload_envelope: ExecutionPayloadEnvelopeV3 = built_payload.into();
564                let new_payload_block_hash = execution_payload_envelope
565                    .execution_payload
566                    .payload_inner
567                    .payload_inner
568                    .block_hash;
569
570                if rpc_latest_header.hash != new_payload_block_hash {
571                    debug!(
572                        "Client {}: The hash is not matched: {:?} {:?}",
573                        idx, rpc_latest_header.hash, new_payload_block_hash
574                    );
575                    continue;
576                }
577
578                if rpc_latest_header.inner.difficulty != alloy_primitives::U256::ZERO {
579                    debug!(
580                        "Client {}: difficulty != 0: {:?}",
581                        idx, rpc_latest_header.inner.difficulty
582                    );
583                    continue;
584                }
585
586                if rpc_latest_header.inner.mix_hash != next_new_payload.prev_randao {
587                    debug!(
588                        "Client {}: The mix_hash and prev_randao is not same: {:?} {:?}",
589                        idx, rpc_latest_header.inner.mix_hash, next_new_payload.prev_randao
590                    );
591                    continue;
592                }
593
594                let extra_len = rpc_latest_header.inner.extra_data.len();
595                if extra_len <= 32 {
596                    debug!("Client {}: extra_len is fewer than 32. extra_len: {}", idx, extra_len);
597                    continue;
598                }
599
600                // at least one client passes all the check, save the header in Env
601                if !accepted_check {
602                    accepted_check = true;
603                    // save the current block info in Env
604                    env.set_current_block_info(BlockInfo {
605                        hash: rpc_latest_header.hash,
606                        number: rpc_latest_header.inner.number,
607                        timestamp: rpc_latest_header.inner.timestamp,
608                    })?;
609
610                    // align latest header time and forkchoice state with the accepted canonical
611                    // head
612                    env.active_node_state_mut()?.latest_header_time =
613                        rpc_latest_header.inner.timestamp;
614                    env.active_node_state_mut()?.latest_fork_choice_state.head_block_hash =
615                        rpc_latest_header.hash;
616                }
617            }
618
619            if accepted_check {
620                Ok(())
621            } else {
622                Err(eyre::eyre!("No clients passed payload acceptance checks"))
623            }
624        })
625    }
626}
627
628/// Action that broadcasts the next new payload
629#[derive(Debug, Default)]
630pub struct BroadcastNextNewPayload {
631    /// If true, only send to the active node. If false, broadcast to all nodes.
632    active_node_only: bool,
633}
634
635impl BroadcastNextNewPayload {
636    /// Create a new `BroadcastNextNewPayload` action that only sends to the active node
637    pub const fn with_active_node() -> Self {
638        Self { active_node_only: true }
639    }
640}
641
642impl<Engine> Action<Engine> for BroadcastNextNewPayload
643where
644    Engine: EngineTypes + PayloadTypes,
645    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
646    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
647{
648    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
649        Box::pin(async move {
650            // Get the next new payload to broadcast
651            let next_new_payload = env
652                .active_node_state()?
653                .latest_payload_built
654                .as_ref()
655                .ok_or_else(|| eyre::eyre!("No next built payload found"))?
656                .clone();
657            let parent_beacon_block_root = next_new_payload
658                .parent_beacon_block_root
659                .ok_or_else(|| eyre::eyre!("No parent beacon block root for next new payload"))?;
660
661            let payload_envelope = env
662                .active_node_state()?
663                .latest_payload_envelope
664                .as_ref()
665                .ok_or_else(|| eyre::eyre!("No execution payload envelope available"))?
666                .clone();
667
668            let execution_payload_envelope: ExecutionPayloadEnvelopeV3 = payload_envelope.into();
669            let execution_payload = execution_payload_envelope.execution_payload;
670
671            if self.active_node_only {
672                // Send only to the active node
673                let active_idx = env.active_node_idx;
674                let engine = env.node_clients[active_idx].engine.http_client();
675
676                let result = EngineApiClient::<Engine>::new_payload_v3(
677                    &engine,
678                    execution_payload.clone(),
679                    vec![],
680                    parent_beacon_block_root,
681                )
682                .await?;
683
684                debug!("Active node {}: new_payload status: {:?}", active_idx, result.status);
685
686                // Validate the response
687                match result.status {
688                    PayloadStatusEnum::Valid => {
689                        env.active_node_state_mut()?.latest_payload_executed =
690                            Some(next_new_payload);
691                        Ok(())
692                    }
693                    other => Err(eyre::eyre!(
694                        "Active node {}: Unexpected payload status: {:?}",
695                        active_idx,
696                        other
697                    )),
698                }
699            } else {
700                // Loop through all clients and broadcast the next new payload
701                let mut broadcast_results = Vec::new();
702                let mut first_valid_seen = false;
703
704                for (idx, client) in env.node_clients.iter().enumerate() {
705                    let engine = client.engine.http_client();
706
707                    // Broadcast the execution payload
708                    let result = EngineApiClient::<Engine>::new_payload_v3(
709                        &engine,
710                        execution_payload.clone(),
711                        vec![],
712                        parent_beacon_block_root,
713                    )
714                    .await?;
715
716                    broadcast_results.push((idx, result.status.clone()));
717                    debug!("Node {}: new_payload broadcast status: {:?}", idx, result.status);
718
719                    // Check if this node accepted the payload
720                    if result.is_valid() && !first_valid_seen {
721                        first_valid_seen = true;
722                    } else if let PayloadStatusEnum::Invalid { validation_error } = result.status {
723                        debug!(
724                            "Node {}: Invalid payload status returned from broadcast: {:?}",
725                            idx, validation_error
726                        );
727                    }
728                }
729
730                // Update the executed payload state after broadcasting to all nodes
731                if first_valid_seen {
732                    env.active_node_state_mut()?.latest_payload_executed = Some(next_new_payload);
733                }
734
735                // Check if at least one node accepted the payload
736                let any_valid = broadcast_results.iter().any(|(_, status)| status.is_valid());
737                if !any_valid {
738                    return Err(eyre::eyre!(
739                        "Failed to successfully broadcast payload to any client"
740                    ));
741                }
742
743                debug!("Broadcast complete. Results: {:?}", broadcast_results);
744
745                Ok(())
746            }
747        })
748    }
749}
750
751/// Action that produces a sequence of blocks using the available clients
752#[derive(Debug)]
753pub struct ProduceBlocks<Engine> {
754    /// Number of blocks to produce
755    pub num_blocks: u64,
756    /// Tracks engine type
757    _phantom: PhantomData<Engine>,
758}
759
760impl<Engine> ProduceBlocks<Engine> {
761    /// Create a new `ProduceBlocks` action
762    pub fn new(num_blocks: u64) -> Self {
763        Self { num_blocks, _phantom: Default::default() }
764    }
765}
766
767impl<Engine> Default for ProduceBlocks<Engine> {
768    fn default() -> Self {
769        Self::new(0)
770    }
771}
772
773impl<Engine> Action<Engine> for ProduceBlocks<Engine>
774where
775    Engine: EngineTypes + PayloadTypes,
776    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
777    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
778{
779    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
780        Box::pin(async move {
781            for _ in 0..self.num_blocks {
782                // create a fresh sequence for each block to avoid state pollution
783                // Note: This produces blocks but does NOT make them canonical
784                // Use MakeCanonical action explicitly if canonicalization is needed
785                let mut sequence = Sequence::new(vec![
786                    Box::new(PickNextBlockProducer::default()),
787                    Box::new(GeneratePayloadAttributes::default()),
788                    Box::new(GenerateNextPayload::default()),
789                    Box::new(BroadcastNextNewPayload::default()),
790                    Box::new(UpdateBlockInfoToLatestPayload::default()),
791                ]);
792                sequence.execute(env).await?;
793            }
794            Ok(())
795        })
796    }
797}
798
799/// Action to test forkchoice update to a tagged block with expected status
800#[derive(Debug)]
801pub struct TestFcuToTag {
802    /// Tag name of the target block
803    pub tag: String,
804    /// Expected payload status
805    pub expected_status: PayloadStatusEnum,
806}
807
808impl TestFcuToTag {
809    /// Create a new `TestFcuToTag` action
810    pub fn new(tag: impl Into<String>, expected_status: PayloadStatusEnum) -> Self {
811        Self { tag: tag.into(), expected_status }
812    }
813}
814
815impl<Engine> Action<Engine> for TestFcuToTag
816where
817    Engine: EngineTypes,
818{
819    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
820        Box::pin(async move {
821            // get the target block from the registry
822            let (target_block, _node_idx) = env
823                .block_registry
824                .get(&self.tag)
825                .copied()
826                .ok_or_else(|| eyre::eyre!("Block tag '{}' not found in registry", self.tag))?;
827
828            let engine_client = env.node_clients[0].engine.http_client();
829            let fcu_state = ForkchoiceState::same_hash(target_block.hash);
830
831            let fcu_response =
832                EngineApiClient::<Engine>::fork_choice_updated_v2(&engine_client, fcu_state, None)
833                    .await?;
834
835            // validate the response matches expected status
836            match (&fcu_response.payload_status.status, &self.expected_status) {
837                (PayloadStatusEnum::Valid, PayloadStatusEnum::Valid) => {
838                    debug!("FCU to '{}' returned VALID as expected", self.tag);
839                }
840                (PayloadStatusEnum::Invalid { .. }, PayloadStatusEnum::Invalid { .. }) => {
841                    debug!("FCU to '{}' returned INVALID as expected", self.tag);
842                }
843                (PayloadStatusEnum::Syncing, PayloadStatusEnum::Syncing) => {
844                    debug!("FCU to '{}' returned SYNCING as expected", self.tag);
845                }
846                (PayloadStatusEnum::Accepted, PayloadStatusEnum::Accepted) => {
847                    debug!("FCU to '{}' returned ACCEPTED as expected", self.tag);
848                }
849                (actual, expected) => {
850                    return Err(eyre::eyre!(
851                        "FCU to '{}': expected status {:?}, but got {:?}",
852                        self.tag,
853                        expected,
854                        actual
855                    ));
856                }
857            }
858
859            Ok(())
860        })
861    }
862}
863
864/// Action to expect a specific FCU status when targeting a tagged block
865#[derive(Debug)]
866pub struct ExpectFcuStatus {
867    /// Tag name of the target block
868    pub target_tag: String,
869    /// Expected payload status
870    pub expected_status: PayloadStatusEnum,
871}
872
873impl ExpectFcuStatus {
874    /// Create a new `ExpectFcuStatus` action expecting VALID status
875    pub fn valid(target_tag: impl Into<String>) -> Self {
876        Self { target_tag: target_tag.into(), expected_status: PayloadStatusEnum::Valid }
877    }
878
879    /// Create a new `ExpectFcuStatus` action expecting INVALID status
880    pub fn invalid(target_tag: impl Into<String>) -> Self {
881        Self {
882            target_tag: target_tag.into(),
883            expected_status: PayloadStatusEnum::Invalid {
884                validation_error: "corrupted block".to_string(),
885            },
886        }
887    }
888
889    /// Create a new `ExpectFcuStatus` action expecting SYNCING status
890    pub fn syncing(target_tag: impl Into<String>) -> Self {
891        Self { target_tag: target_tag.into(), expected_status: PayloadStatusEnum::Syncing }
892    }
893
894    /// Create a new `ExpectFcuStatus` action expecting ACCEPTED status
895    pub fn accepted(target_tag: impl Into<String>) -> Self {
896        Self { target_tag: target_tag.into(), expected_status: PayloadStatusEnum::Accepted }
897    }
898}
899
900impl<Engine> Action<Engine> for ExpectFcuStatus
901where
902    Engine: EngineTypes,
903{
904    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
905        Box::pin(async move {
906            let mut test_fcu = TestFcuToTag::new(&self.target_tag, self.expected_status.clone());
907            test_fcu.execute(env).await
908        })
909    }
910}
911
912/// Action to validate that a tagged block remains canonical by performing FCU to it
913#[derive(Debug)]
914pub struct ValidateCanonicalTag {
915    /// Tag name of the block to validate as canonical
916    pub tag: String,
917}
918
919impl ValidateCanonicalTag {
920    /// Create a new `ValidateCanonicalTag` action
921    pub fn new(tag: impl Into<String>) -> Self {
922        Self { tag: tag.into() }
923    }
924}
925
926impl<Engine> Action<Engine> for ValidateCanonicalTag
927where
928    Engine: EngineTypes,
929{
930    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
931        Box::pin(async move {
932            let mut expect_valid = ExpectFcuStatus::valid(&self.tag);
933            expect_valid.execute(env).await?;
934
935            debug!("Successfully validated that '{}' remains canonical", self.tag);
936            Ok(())
937        })
938    }
939}
940
941/// Action that produces blocks locally without broadcasting to other nodes
942/// This sends the payload only to the active node to ensure it's available locally
943#[derive(Debug)]
944pub struct ProduceBlocksLocally<Engine> {
945    /// Number of blocks to produce
946    pub num_blocks: u64,
947    /// Tracks engine type
948    _phantom: PhantomData<Engine>,
949}
950
951impl<Engine> ProduceBlocksLocally<Engine> {
952    /// Create a new `ProduceBlocksLocally` action
953    pub fn new(num_blocks: u64) -> Self {
954        Self { num_blocks, _phantom: Default::default() }
955    }
956}
957
958impl<Engine> Default for ProduceBlocksLocally<Engine> {
959    fn default() -> Self {
960        Self::new(0)
961    }
962}
963
964impl<Engine> Action<Engine> for ProduceBlocksLocally<Engine>
965where
966    Engine: EngineTypes + PayloadTypes,
967    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
968    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
969{
970    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
971        Box::pin(async move {
972            // Remember the active node to ensure all blocks are produced on the same node
973            let producer_idx = env.active_node_idx;
974
975            for _ in 0..self.num_blocks {
976                // Ensure we always use the same producer
977                env.last_producer_idx = Some(producer_idx);
978
979                // create a sequence that produces blocks and sends only to active node
980                let mut sequence = Sequence::new(vec![
981                    // Skip PickNextBlockProducer to maintain the same producer
982                    Box::new(GeneratePayloadAttributes::default()),
983                    Box::new(GenerateNextPayload::default()),
984                    // Send payload only to the active node to make it available
985                    Box::new(BroadcastNextNewPayload::with_active_node()),
986                    Box::new(UpdateBlockInfoToLatestPayload::default()),
987                ]);
988                sequence.execute(env).await?;
989            }
990            Ok(())
991        })
992    }
993}
994
995/// Action that produces a sequence of blocks where some blocks are intentionally invalid
996#[derive(Debug)]
997pub struct ProduceInvalidBlocks<Engine> {
998    /// Number of blocks to produce
999    pub num_blocks: u64,
1000    /// Set of indices (0-based) where blocks should be made invalid
1001    pub invalid_indices: HashSet<u64>,
1002    /// Tracks engine type
1003    _phantom: PhantomData<Engine>,
1004}
1005
1006impl<Engine> ProduceInvalidBlocks<Engine> {
1007    /// Create a new `ProduceInvalidBlocks` action
1008    pub fn new(num_blocks: u64, invalid_indices: HashSet<u64>) -> Self {
1009        Self { num_blocks, invalid_indices, _phantom: Default::default() }
1010    }
1011
1012    /// Create a new `ProduceInvalidBlocks` action with a single invalid block at the specified
1013    /// index
1014    pub fn with_invalid_at(num_blocks: u64, invalid_index: u64) -> Self {
1015        let mut invalid_indices = HashSet::new();
1016        invalid_indices.insert(invalid_index);
1017        Self::new(num_blocks, invalid_indices)
1018    }
1019}
1020
1021impl<Engine> Action<Engine> for ProduceInvalidBlocks<Engine>
1022where
1023    Engine: EngineTypes + PayloadTypes,
1024    Engine::PayloadAttributes: From<PayloadAttributes> + Clone,
1025    Engine::ExecutionPayloadEnvelopeV3: Into<ExecutionPayloadEnvelopeV3>,
1026{
1027    fn execute<'a>(&'a mut self, env: &'a mut Environment<Engine>) -> BoxFuture<'a, Result<()>> {
1028        Box::pin(async move {
1029            for block_index in 0..self.num_blocks {
1030                let is_invalid = self.invalid_indices.contains(&block_index);
1031
1032                if is_invalid {
1033                    debug!("Producing invalid block at index {}", block_index);
1034
1035                    // produce a valid block first, then corrupt it
1036                    let mut sequence = Sequence::new(vec![
1037                        Box::new(PickNextBlockProducer::default()),
1038                        Box::new(GeneratePayloadAttributes::default()),
1039                        Box::new(GenerateNextPayload::default()),
1040                    ]);
1041                    sequence.execute(env).await?;
1042
1043                    // get the latest payload and corrupt it
1044                    let latest_envelope =
1045                        env.active_node_state()?.latest_payload_envelope.as_ref().ok_or_else(
1046                            || eyre::eyre!("No payload envelope available to corrupt"),
1047                        )?;
1048
1049                    let envelope_v3: ExecutionPayloadEnvelopeV3 = latest_envelope.clone().into();
1050                    let mut corrupted_payload = envelope_v3.execution_payload;
1051
1052                    // corrupt the state root to make the block invalid
1053                    corrupted_payload.payload_inner.payload_inner.state_root = B256::random();
1054
1055                    debug!(
1056                        "Corrupted state root for block {} to: {}",
1057                        block_index, corrupted_payload.payload_inner.payload_inner.state_root
1058                    );
1059
1060                    // send the corrupted payload via newPayload
1061                    let engine_client = env.node_clients[0].engine.http_client();
1062                    // for simplicity, we'll use empty versioned hashes for invalid block testing
1063                    let versioned_hashes = Vec::new();
1064                    // use a random parent beacon block root since this is for invalid block testing
1065                    let parent_beacon_block_root = B256::random();
1066
1067                    let new_payload_response = EngineApiClient::<Engine>::new_payload_v3(
1068                        &engine_client,
1069                        corrupted_payload.clone(),
1070                        versioned_hashes,
1071                        parent_beacon_block_root,
1072                    )
1073                    .await?;
1074
1075                    // expect the payload to be rejected as invalid
1076                    match new_payload_response.status {
1077                        PayloadStatusEnum::Invalid { validation_error } => {
1078                            debug!(
1079                                "Block {} correctly rejected as invalid: {:?}",
1080                                block_index, validation_error
1081                            );
1082                        }
1083                        other_status => {
1084                            return Err(eyre::eyre!(
1085                                "Expected block {} to be rejected as INVALID, but got: {:?}",
1086                                block_index,
1087                                other_status
1088                            ));
1089                        }
1090                    }
1091
1092                    // update block info with the corrupted block (for potential future reference)
1093                    env.set_current_block_info(BlockInfo {
1094                        hash: corrupted_payload.payload_inner.payload_inner.block_hash,
1095                        number: corrupted_payload.payload_inner.payload_inner.block_number,
1096                        timestamp: corrupted_payload.timestamp(),
1097                    })?;
1098                } else {
1099                    debug!("Producing valid block at index {}", block_index);
1100
1101                    // produce a valid block normally
1102                    let mut sequence = Sequence::new(vec![
1103                        Box::new(PickNextBlockProducer::default()),
1104                        Box::new(GeneratePayloadAttributes::default()),
1105                        Box::new(GenerateNextPayload::default()),
1106                        Box::new(BroadcastNextNewPayload::default()),
1107                        Box::new(UpdateBlockInfoToLatestPayload::default()),
1108                    ]);
1109                    sequence.execute(env).await?;
1110                }
1111            }
1112            Ok(())
1113        })
1114    }
1115}
1116
1117/// Broadcasts a forkchoice update to all clients, using `head_hash` as the head and safe block.
1118pub(super) async fn broadcast_forkchoice<Engine>(
1119    env: &Environment<Engine>,
1120    head_hash: B256,
1121) -> Result<()>
1122where
1123    Engine: EngineTypes,
1124{
1125    let fork_choice_state = ForkchoiceState {
1126        head_block_hash: head_hash,
1127        safe_block_hash: head_hash,
1128        // Making a block canonical does not imply finality: tests advance the finalized block
1129        // explicitly via `FinalizeBlock`, and a finalized tip would reject any later forkchoice
1130        // update below it as a too deep reorg.
1131        finalized_block_hash: B256::ZERO,
1132    };
1133    debug!(
1134        "Broadcasting forkchoice update to {} clients. Head: {:?}",
1135        env.node_clients.len(),
1136        fork_choice_state.head_block_hash
1137    );
1138
1139    for (idx, client) in env.node_clients.iter().enumerate() {
1140        match EngineApiClient::<Engine>::fork_choice_updated_v3(
1141            &client.engine.http_client(),
1142            fork_choice_state,
1143            None,
1144        )
1145        .await
1146        {
1147            Ok(resp) => {
1148                debug!(
1149                    "Client {}: Forkchoice update status: {:?}",
1150                    idx, resp.payload_status.status
1151                );
1152                // validate that the forkchoice update was accepted
1153                validate_fcu_response(&resp, &format!("Client {idx}"))?;
1154            }
1155            Err(err) => {
1156                return Err(eyre::eyre!(
1157                    "Client {}: Failed to broadcast forkchoice: {:?}",
1158                    idx,
1159                    err
1160                ));
1161            }
1162        }
1163    }
1164    debug!("Forkchoice update broadcasted successfully");
1165    Ok(())
1166}