1use 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#[derive(Debug)]
24pub struct AssertMineBlock<Engine>
25where
26 Engine: PayloadTypes,
27{
28 pub node_idx: usize,
30 pub transactions: Vec<Bytes>,
32 pub expected_hash: Option<B256>,
34 pub payload_attributes: Engine::PayloadAttributes,
37 _phantom: PhantomData<Engine>,
39}
40
41impl<Engine> AssertMineBlock<Engine>
42where
43 Engine: PayloadTypes,
44{
45 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 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 let fork_choice_state = ForkchoiceState::same_hash(parent_hash);
96
97 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 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#[derive(Debug, Default)]
164pub struct PickNextBlockProducer {}
165
166impl PickNextBlockProducer {
167 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 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#[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#[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 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 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 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 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 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#[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 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 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#[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 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 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#[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 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#[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 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 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 if !accepted_check {
602 accepted_check = true;
603 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 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#[derive(Debug, Default)]
630pub struct BroadcastNextNewPayload {
631 active_node_only: bool,
633}
634
635impl BroadcastNextNewPayload {
636 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 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 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 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 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 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 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 if first_valid_seen {
732 env.active_node_state_mut()?.latest_payload_executed = Some(next_new_payload);
733 }
734
735 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#[derive(Debug)]
753pub struct ProduceBlocks<Engine> {
754 pub num_blocks: u64,
756 _phantom: PhantomData<Engine>,
758}
759
760impl<Engine> ProduceBlocks<Engine> {
761 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 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#[derive(Debug)]
801pub struct TestFcuToTag {
802 pub tag: String,
804 pub expected_status: PayloadStatusEnum,
806}
807
808impl TestFcuToTag {
809 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 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 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#[derive(Debug)]
866pub struct ExpectFcuStatus {
867 pub target_tag: String,
869 pub expected_status: PayloadStatusEnum,
871}
872
873impl ExpectFcuStatus {
874 pub fn valid(target_tag: impl Into<String>) -> Self {
876 Self { target_tag: target_tag.into(), expected_status: PayloadStatusEnum::Valid }
877 }
878
879 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 pub fn syncing(target_tag: impl Into<String>) -> Self {
891 Self { target_tag: target_tag.into(), expected_status: PayloadStatusEnum::Syncing }
892 }
893
894 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#[derive(Debug)]
914pub struct ValidateCanonicalTag {
915 pub tag: String,
917}
918
919impl ValidateCanonicalTag {
920 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#[derive(Debug)]
944pub struct ProduceBlocksLocally<Engine> {
945 pub num_blocks: u64,
947 _phantom: PhantomData<Engine>,
949}
950
951impl<Engine> ProduceBlocksLocally<Engine> {
952 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 let producer_idx = env.active_node_idx;
974
975 for _ in 0..self.num_blocks {
976 env.last_producer_idx = Some(producer_idx);
978
979 let mut sequence = Sequence::new(vec![
981 Box::new(GeneratePayloadAttributes::default()),
983 Box::new(GenerateNextPayload::default()),
984 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#[derive(Debug)]
997pub struct ProduceInvalidBlocks<Engine> {
998 pub num_blocks: u64,
1000 pub invalid_indices: HashSet<u64>,
1002 _phantom: PhantomData<Engine>,
1004}
1005
1006impl<Engine> ProduceInvalidBlocks<Engine> {
1007 pub fn new(num_blocks: u64, invalid_indices: HashSet<u64>) -> Self {
1009 Self { num_blocks, invalid_indices, _phantom: Default::default() }
1010 }
1011
1012 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 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 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 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 let engine_client = env.node_clients[0].engine.http_client();
1062 let versioned_hashes = Vec::new();
1064 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 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 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 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
1117pub(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 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_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}