1use crate::{
2 capabilities::EngineCapabilities, metrics::EngineApiMetrics, EngineApiError, EngineApiResult,
3};
4use alloy_eips::{
5 eip1898::BlockHashOrNumber,
6 eip4844::{BlobAndProofV1, BlobAndProofV2, BlobCellsAndProofsV1},
7 eip4895::Withdrawals,
8 eip7685::RequestsOrHash,
9};
10use alloy_primitives::{BlockHash, BlockNumber, Bytes, Sealable, B128, B256, U64};
11use alloy_rpc_types_engine::{
12 BogotaPayloadFields, CancunPayloadFields, ClientVersionV1, ExecutionData,
13 ExecutionPayloadBodiesV1, ExecutionPayloadBodiesV2, ExecutionPayloadBodyV1,
14 ExecutionPayloadBodyV2, ExecutionPayloadInputV2, ExecutionPayloadSidecar, ExecutionPayloadV1,
15 ExecutionPayloadV3, ExecutionPayloadV4, ForkchoiceState, ForkchoiceUpdated,
16 ForkchoiceUpdatedResponseV2, PayloadId, PayloadStatus, PayloadStatusV2, PraguePayloadFields,
17 MAX_BYTES_PER_INCLUSION_LIST,
18};
19use async_trait::async_trait;
20use jsonrpsee_core::{server::RpcModule, RpcResult};
21use reth_chainspec::EthereumHardforks;
22use reth_engine_primitives::{ConsensusEngineHandle, EngineApiValidator, EngineTypes};
23use reth_network_api::{CellCustody, NetworkInfo};
24use reth_payload_builder::PayloadStore;
25use reth_payload_primitives::{
26 validate_payload_timestamp, EngineApiMessageVersion, MessageValidationKind,
27 PayloadOrAttributes, PayloadTypes,
28};
29use reth_primitives_traits::{AlloyBlockHeader, Block, BlockBody};
30use reth_rpc_api::{EngineApiServer, IntoEngineApiRpcModule};
31use reth_storage_api::{BalProvider, BlockReader, HeaderProvider, StateProviderFactory};
32use reth_tasks::Runtime;
33use reth_transaction_pool::{BestTransactions, TransactionPool};
34use std::{
35 sync::Arc,
36 time::{Instant, SystemTime},
37};
38use tokio::sync::oneshot;
39use tracing::{debug, trace, warn};
40
41pub type EngineApiSender<Ok> = oneshot::Sender<EngineApiResult<Ok>>;
43
44const MAX_PAYLOAD_BODIES_LIMIT: u64 = 1024;
46
47const MAX_BLOB_LIMIT: usize = 128;
49
50pub struct EngineApi<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
66 inner: Arc<EngineApiInner<Provider, PayloadT, Pool, Validator, ChainSpec>>,
67}
68
69impl<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec>
70 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
71{
72 pub fn chain_spec(&self) -> &Arc<ChainSpec> {
74 &self.inner.chain_spec
75 }
76
77 pub fn client_version(&self) -> &ClientVersionV1 {
79 &self.inner.client
80 }
81}
82
83impl<Provider, PayloadT, Pool, Validator, ChainSpec>
84 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
85where
86 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
87 PayloadT: PayloadTypes,
88 Pool: TransactionPool + 'static,
89 Validator: EngineApiValidator<PayloadT>,
90 ChainSpec: EthereumHardforks + Send + Sync + 'static,
91{
92 #[expect(clippy::too_many_arguments)]
94 pub fn new(
95 provider: Provider,
96 chain_spec: Arc<ChainSpec>,
97 beacon_consensus: ConsensusEngineHandle<PayloadT>,
98 payload_store: PayloadStore<PayloadT>,
99 tx_pool: Pool,
100 task_spawner: Runtime,
101 client: ClientVersionV1,
102 capabilities: EngineCapabilities,
103 validator: Validator,
104 accept_execution_requests_hash: bool,
105 network: impl NetworkInfo + 'static,
106 ) -> Self {
107 let cell_custody = network.cell_custody().clone();
108 let is_syncing = Arc::new(move || network.is_syncing());
109 let inner = Arc::new(EngineApiInner {
110 provider,
111 chain_spec,
112 beacon_consensus,
113 payload_store,
114 task_spawner,
115 metrics: EngineApiMetrics::default(),
116 client,
117 capabilities,
118 tx_pool,
119 validator,
120 accept_execution_requests_hash,
121 cell_custody,
122 is_syncing,
123 });
124 Self { inner }
125 }
126
127 pub fn get_client_version_v1(
129 &self,
130 _client: ClientVersionV1,
131 ) -> EngineApiResult<Vec<ClientVersionV1>> {
132 Ok(vec![self.inner.client.clone()])
133 }
134
135 async fn get_payload_timestamp(&self, payload_id: PayloadId) -> EngineApiResult<u64> {
137 Ok(self
138 .inner
139 .payload_store
140 .payload_timestamp(payload_id)
141 .await
142 .ok_or(EngineApiError::UnknownPayload)??)
143 }
144
145 pub async fn new_payload_v1(
148 &self,
149 payload: PayloadT::ExecutionData,
150 ) -> EngineApiResult<PayloadStatus> {
151 let payload_or_attrs = PayloadOrAttributes::<
152 '_,
153 PayloadT::ExecutionData,
154 PayloadT::PayloadAttributes,
155 >::from_execution_payload(&payload);
156
157 self.inner
158 .validator
159 .validate_version_specific_fields(EngineApiMessageVersion::V1, payload_or_attrs)?;
160
161 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
162 }
163
164 pub async fn new_payload_v1_metered(
166 &self,
167 payload: PayloadT::ExecutionData,
168 ) -> EngineApiResult<PayloadStatus> {
169 let start = Instant::now();
170 let res = Self::new_payload_v1(self, payload).await;
171 let elapsed = start.elapsed();
172 self.inner.metrics.latency.new_payload_v1.record(elapsed);
173 res
174 }
175
176 pub async fn new_payload_v2(
178 &self,
179 payload: PayloadT::ExecutionData,
180 ) -> EngineApiResult<PayloadStatus> {
181 let payload_or_attrs = PayloadOrAttributes::<
182 '_,
183 PayloadT::ExecutionData,
184 PayloadT::PayloadAttributes,
185 >::from_execution_payload(&payload);
186 self.inner
187 .validator
188 .validate_version_specific_fields(EngineApiMessageVersion::V2, payload_or_attrs)?;
189 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
190 }
191
192 pub async fn new_payload_v2_metered(
194 &self,
195 payload: PayloadT::ExecutionData,
196 ) -> EngineApiResult<PayloadStatus> {
197 let start = Instant::now();
198 let res = Self::new_payload_v2(self, payload).await;
199 let elapsed = start.elapsed();
200 self.inner.metrics.latency.new_payload_v2.record(elapsed);
201 res
202 }
203
204 pub async fn new_payload_v3(
206 &self,
207 payload: PayloadT::ExecutionData,
208 ) -> EngineApiResult<PayloadStatus> {
209 let payload_or_attrs = PayloadOrAttributes::<
210 '_,
211 PayloadT::ExecutionData,
212 PayloadT::PayloadAttributes,
213 >::from_execution_payload(&payload);
214 self.inner
215 .validator
216 .validate_version_specific_fields(EngineApiMessageVersion::V3, payload_or_attrs)?;
217
218 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
219 }
220
221 pub async fn new_payload_v3_metered(
223 &self,
224 payload: PayloadT::ExecutionData,
225 ) -> RpcResult<PayloadStatus> {
226 let start = Instant::now();
227
228 let res = Self::new_payload_v3(self, payload).await;
229 let elapsed = start.elapsed();
230 self.inner.metrics.latency.new_payload_v3.record(elapsed);
231 Ok(res?)
232 }
233
234 pub async fn new_payload_v4(
236 &self,
237 payload: PayloadT::ExecutionData,
238 ) -> EngineApiResult<PayloadStatus> {
239 let payload_or_attrs = PayloadOrAttributes::<
240 '_,
241 PayloadT::ExecutionData,
242 PayloadT::PayloadAttributes,
243 >::from_execution_payload(&payload);
244 self.inner
245 .validator
246 .validate_version_specific_fields(EngineApiMessageVersion::V4, payload_or_attrs)?;
247
248 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
249 }
250
251 pub async fn new_payload_v4_metered(
253 &self,
254 payload: PayloadT::ExecutionData,
255 ) -> RpcResult<PayloadStatus> {
256 let start = Instant::now();
257 let res = Self::new_payload_v4(self, payload).await;
258
259 let elapsed = start.elapsed();
260 self.inner.metrics.latency.new_payload_v4.record(elapsed);
261 Ok(res?)
262 }
263
264 pub async fn new_payload_v5(
270 &self,
271 payload: PayloadT::ExecutionData,
272 ) -> EngineApiResult<PayloadStatus> {
273 let payload_or_attrs = PayloadOrAttributes::<
274 '_,
275 PayloadT::ExecutionData,
276 PayloadT::PayloadAttributes,
277 >::from_execution_payload(&payload);
278 self.inner
279 .validator
280 .validate_version_specific_fields(EngineApiMessageVersion::V5, payload_or_attrs)?;
281 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
282 }
283
284 pub async fn new_payload_v5_metered(
286 &self,
287 payload: PayloadT::ExecutionData,
288 ) -> RpcResult<PayloadStatus> {
289 let start = Instant::now();
290 let res = Self::new_payload_v5(self, payload).await;
291 let elapsed = start.elapsed();
292 self.inner.metrics.latency.new_payload_v5.record(elapsed);
293 Ok(res?)
294 }
295
296 pub async fn new_payload_v6(
300 &self,
301 payload: PayloadT::ExecutionData,
302 ) -> EngineApiResult<PayloadStatusV2> {
303 let payload_or_attrs = PayloadOrAttributes::<
304 '_,
305 PayloadT::ExecutionData,
306 PayloadT::PayloadAttributes,
307 >::from_execution_payload(&payload);
308 self.inner
309 .validator
310 .validate_version_specific_fields(EngineApiMessageVersion::V6, payload_or_attrs)?;
311
312 Ok(self.inner.beacon_consensus.new_payload(payload).await?.into())
313 }
314
315 pub async fn new_payload_v6_metered(
317 &self,
318 payload: PayloadT::ExecutionData,
319 ) -> EngineApiResult<PayloadStatusV2> {
320 let start = Instant::now();
321 let result = Self::new_payload_v6(self, payload).await;
322 self.inner.metrics.latency.new_payload_v6.record(start.elapsed());
323 result
324 }
325
326 pub fn accept_execution_requests_hash(&self) -> bool {
328 self.inner.accept_execution_requests_hash
329 }
330}
331
332impl<Provider, EngineT, Pool, Validator, ChainSpec>
333 EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
334where
335 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
336 EngineT: EngineTypes,
337 Pool: TransactionPool + 'static,
338 Validator: EngineApiValidator<EngineT>,
339 ChainSpec: EthereumHardforks + Send + Sync + 'static,
340{
341 pub async fn fork_choice_updated_v1(
348 &self,
349 state: ForkchoiceState,
350 payload_attrs: Option<EngineT::PayloadAttributes>,
351 ) -> EngineApiResult<ForkchoiceUpdated> {
352 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V1, state, payload_attrs)
353 .await
354 }
355
356 pub async fn fork_choice_updated_v1_metered(
358 &self,
359 state: ForkchoiceState,
360 payload_attrs: Option<EngineT::PayloadAttributes>,
361 ) -> EngineApiResult<ForkchoiceUpdated> {
362 let start = Instant::now();
363 let res = Self::fork_choice_updated_v1(self, state, payload_attrs).await;
364 self.inner.metrics.latency.fork_choice_updated_v1.record(start.elapsed());
365 res
366 }
367
368 pub async fn fork_choice_updated_v2(
373 &self,
374 state: ForkchoiceState,
375 payload_attrs: Option<EngineT::PayloadAttributes>,
376 ) -> EngineApiResult<ForkchoiceUpdated> {
377 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V2, state, payload_attrs)
378 .await
379 }
380
381 pub async fn fork_choice_updated_v2_metered(
383 &self,
384 state: ForkchoiceState,
385 payload_attrs: Option<EngineT::PayloadAttributes>,
386 ) -> EngineApiResult<ForkchoiceUpdated> {
387 let start = Instant::now();
388 let res = Self::fork_choice_updated_v2(self, state, payload_attrs).await;
389 self.inner.metrics.latency.fork_choice_updated_v2.record(start.elapsed());
390 res
391 }
392
393 pub async fn fork_choice_updated_v3(
398 &self,
399 state: ForkchoiceState,
400 payload_attrs: Option<EngineT::PayloadAttributes>,
401 ) -> EngineApiResult<ForkchoiceUpdated> {
402 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V3, state, payload_attrs)
403 .await
404 }
405
406 pub async fn fork_choice_updated_v3_metered(
408 &self,
409 state: ForkchoiceState,
410 payload_attrs: Option<EngineT::PayloadAttributes>,
411 ) -> EngineApiResult<ForkchoiceUpdated> {
412 let start = Instant::now();
413 let res = Self::fork_choice_updated_v3(self, state, payload_attrs).await;
414 self.inner.metrics.latency.fork_choice_updated_v3.record(start.elapsed());
415 res
416 }
417
418 pub async fn fork_choice_updated_v4(
423 &self,
424 state: ForkchoiceState,
425 payload_attrs: Option<EngineT::PayloadAttributes>,
426 custody_columns: Option<B128>,
427 ) -> EngineApiResult<ForkchoiceUpdated> {
428 if let Some(custody_columns) = custody_columns {
429 self.inner.cell_custody.set_from_engine_api(custody_columns);
430 }
431 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V4, state, payload_attrs)
432 .await
433 }
434
435 pub async fn fork_choice_updated_v4_metered(
437 &self,
438 state: ForkchoiceState,
439 payload_attrs: Option<EngineT::PayloadAttributes>,
440 custody_columns: Option<B128>,
441 ) -> EngineApiResult<ForkchoiceUpdated> {
442 let start = Instant::now();
443 let res = Self::fork_choice_updated_v4(self, state, payload_attrs, custody_columns).await;
444 self.inner.metrics.latency.fork_choice_updated_v4.record(start.elapsed());
445 res
446 }
447
448 pub async fn fork_choice_updated_v5(
452 &self,
453 state: ForkchoiceState,
454 payload_attrs: Option<EngineT::PayloadAttributes>,
455 custody_columns: Option<B128>,
456 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
457 if let Some(custody_columns) = custody_columns {
458 self.inner.cell_custody.set_from_engine_api(custody_columns);
459 }
460
461 Ok(self
463 .validate_and_execute_forkchoice(EngineApiMessageVersion::V5, state, payload_attrs)
464 .await?
465 .into())
466 }
467
468 pub async fn fork_choice_updated_v5_metered(
470 &self,
471 state: ForkchoiceState,
472 payload_attrs: Option<EngineT::PayloadAttributes>,
473 custody_columns: Option<B128>,
474 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
475 let start = Instant::now();
476 let res = Self::fork_choice_updated_v5(self, state, payload_attrs, custody_columns).await;
477 self.inner.metrics.latency.fork_choice_updated_v5.record(start.elapsed());
478 res
479 }
480
481 pub fn get_inclusion_list_v1(&self) -> EngineApiResult<Vec<Bytes>> {
483 let mut total_size = 0;
484 let mut inclusion_list = Vec::new();
485
486 for pool_tx in self.inner.tx_pool.best_transactions().without_blobs().without_updates() {
487 let encoded = pool_tx.encoded_2718_consensus();
488 let new_size = total_size + encoded.len();
489 if new_size > MAX_BYTES_PER_INCLUSION_LIST as usize {
490 break
491 }
492
493 total_size = new_size;
494 inclusion_list.push(encoded);
495 }
496
497 Ok(inclusion_list)
498 }
499
500 pub fn get_inclusion_list_v1_metered(&self) -> EngineApiResult<Vec<Bytes>> {
502 let start = Instant::now();
503 let result = Self::get_inclusion_list_v1(self);
504 self.inner.metrics.latency.get_inclusion_list_v1.record(start.elapsed());
505 result
506 }
507
508 async fn get_built_payload(
510 &self,
511 payload_id: PayloadId,
512 ) -> EngineApiResult<EngineT::BuiltPayload> {
513 self.inner
514 .payload_store
515 .resolve(payload_id)
516 .await
517 .ok_or(EngineApiError::UnknownPayload)?
518 .map_err(|_| EngineApiError::UnknownPayload)
519 }
520
521 async fn get_payload_inner<R>(
524 &self,
525 payload_id: PayloadId,
526 version: EngineApiMessageVersion,
527 ) -> EngineApiResult<R>
528 where
529 EngineT::BuiltPayload: TryInto<R>,
530 {
531 let timestamp = self.get_payload_timestamp(payload_id).await?;
534 validate_payload_timestamp(
535 &self.inner.chain_spec,
536 version,
537 timestamp,
538 MessageValidationKind::GetPayload,
539 )?;
540
541 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
543 warn!(?version, "could not transform built payload");
544 EngineApiError::UnknownPayload
545 })
546 }
547
548 pub async fn get_payload_v1(
558 &self,
559 payload_id: PayloadId,
560 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
561 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
562 warn!(version = ?EngineApiMessageVersion::V1, "could not transform built payload");
563 EngineApiError::UnknownPayload
564 })
565 }
566
567 pub async fn get_payload_v1_metered(
569 &self,
570 payload_id: PayloadId,
571 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
572 let start = Instant::now();
573 let res = Self::get_payload_v1(self, payload_id).await;
574 self.inner.metrics.latency.get_payload_v1.record(start.elapsed());
575 res
576 }
577
578 pub async fn get_payload_v2(
586 &self,
587 payload_id: PayloadId,
588 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
589 self.get_payload_inner(payload_id, EngineApiMessageVersion::V2).await
590 }
591
592 pub async fn get_payload_v2_metered(
594 &self,
595 payload_id: PayloadId,
596 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
597 let start = Instant::now();
598 let res = Self::get_payload_v2(self, payload_id).await;
599 self.inner.metrics.latency.get_payload_v2.record(start.elapsed());
600 res
601 }
602
603 pub async fn get_payload_v3(
611 &self,
612 payload_id: PayloadId,
613 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
614 self.get_payload_inner(payload_id, EngineApiMessageVersion::V3).await
615 }
616
617 pub async fn get_payload_v3_metered(
619 &self,
620 payload_id: PayloadId,
621 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
622 let start = Instant::now();
623 let res = Self::get_payload_v3(self, payload_id).await;
624 self.inner.metrics.latency.get_payload_v3.record(start.elapsed());
625 res
626 }
627
628 pub async fn get_payload_v4(
636 &self,
637 payload_id: PayloadId,
638 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
639 self.get_payload_inner(payload_id, EngineApiMessageVersion::V4).await
640 }
641
642 pub async fn get_payload_v4_metered(
644 &self,
645 payload_id: PayloadId,
646 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
647 let start = Instant::now();
648 let res = Self::get_payload_v4(self, payload_id).await;
649 self.inner.metrics.latency.get_payload_v4.record(start.elapsed());
650 res
651 }
652
653 pub async fn get_payload_v5(
663 &self,
664 payload_id: PayloadId,
665 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
666 self.get_payload_inner(payload_id, EngineApiMessageVersion::V5).await
667 }
668
669 pub async fn get_payload_v5_metered(
671 &self,
672 payload_id: PayloadId,
673 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
674 let start = Instant::now();
675 let res = Self::get_payload_v5(self, payload_id).await;
676 self.inner.metrics.latency.get_payload_v5.record(start.elapsed());
677 res
678 }
679
680 pub async fn get_payload_v6(
686 &self,
687 payload_id: PayloadId,
688 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
689 self.get_payload_inner(payload_id, EngineApiMessageVersion::V6).await
690 }
691
692 pub async fn get_payload_v6_metered(
694 &self,
695 payload_id: PayloadId,
696 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
697 let start = Instant::now();
698 let res = Self::get_payload_v6(self, payload_id).await;
699 self.inner.metrics.latency.get_payload_v6.record(start.elapsed());
700 res
701 }
702
703 pub async fn get_payload_bodies_by_range_with<F, R>(
706 &self,
707 start: BlockNumber,
708 count: u64,
709 f: F,
710 ) -> EngineApiResult<Vec<Option<R>>>
711 where
712 F: Fn(Provider::Block) -> R + Send + 'static,
713 R: Send + 'static,
714 {
715 let (tx, rx) = oneshot::channel();
716 let inner = self.inner.clone();
717
718 self.inner.task_spawner.spawn_blocking_task(async move {
719 if count > MAX_PAYLOAD_BODIES_LIMIT {
720 tx.send(Err(EngineApiError::PayloadRequestTooLarge { len: count })).ok();
721 return;
722 }
723
724 if start == 0 || count == 0 {
725 tx.send(Err(EngineApiError::InvalidBodiesRange { start, count })).ok();
726 return;
727 }
728
729 let mut result = Vec::with_capacity(count as usize);
730
731 let mut end = start.saturating_add(count - 1);
733
734 if let Ok(best_block) = inner.provider.best_block_number()
737 && end > best_block {
738 end = best_block;
739 }
740
741 let earliest_block = inner.provider.earliest_block_number().unwrap_or(0);
743 for num in start..=end {
744 if tx.is_closed() {
745 return;
746 }
747
748 if num < earliest_block {
749 result.push(None);
750 continue;
751 }
752 let block_result = inner.provider.block(BlockHashOrNumber::Number(num));
753 match block_result {
754 Ok(block) => {
755 result.push(block.map(&f));
756 }
757 Err(err) => {
758 tx.send(Err(EngineApiError::Internal(Box::new(err)))).ok();
759 return;
760 }
761 };
762 }
763 tx.send(Ok(result)).ok();
764 });
765
766 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
767 }
768
769 pub async fn get_payload_bodies_by_range_with_timestamps(
774 &self,
775 start: BlockNumber,
776 count: u64,
777 include_bal: bool,
778 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
779 let bodies = self
780 .get_payload_bodies_by_range_with(start, count, Self::payload_body_with_timestamp)
781 .await?;
782 self.attach_payload_body_bals(bodies, include_bal).await
783 }
784
785 pub async fn get_payload_bodies_by_range_with_timestamps_metered(
787 &self,
788 start: BlockNumber,
789 count: u64,
790 include_bal: bool,
791 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
792 let start_time = Instant::now();
793 let result =
794 self.get_payload_bodies_by_range_with_timestamps(start, count, include_bal).await;
795 let latency = &self.inner.metrics.latency;
796 if include_bal {
797 latency.get_payload_bodies_by_range_v2.record(start_time.elapsed());
798 } else {
799 latency.get_payload_bodies_by_range_v1.record(start_time.elapsed());
800 }
801 result
802 }
803
804 pub async fn get_payload_bodies_by_range_v1(
815 &self,
816 start: BlockNumber,
817 count: u64,
818 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
819 self.get_payload_bodies_by_range_with(start, count, |block| ExecutionPayloadBodyV1 {
820 transactions: block.body().encoded_2718_transactions(),
821 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
822 })
823 .await
824 }
825
826 pub async fn get_payload_bodies_by_range_v1_metered(
828 &self,
829 start: BlockNumber,
830 count: u64,
831 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
832 let start_time = Instant::now();
833 let res = Self::get_payload_bodies_by_range_v1(self, start, count).await;
834 self.inner.metrics.latency.get_payload_bodies_by_range_v1.record(start_time.elapsed());
835 res
836 }
837
838 pub async fn get_payload_bodies_by_range_v2(
842 &self,
843 start: BlockNumber,
844 count: u64,
845 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
846 Ok(self
847 .get_payload_bodies_by_range_with_timestamps(start, count, true)
848 .await?
849 .into_iter()
850 .map(|body| body.map(|(_, body)| body))
851 .collect())
852 }
853
854 pub async fn get_payload_bodies_by_range_v2_metered(
856 &self,
857 start: BlockNumber,
858 count: u64,
859 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
860 let start_time = Instant::now();
861 let res = Self::get_payload_bodies_by_range_v2(self, start, count).await;
862 self.inner.metrics.latency.get_payload_bodies_by_range_v2.record(start_time.elapsed());
863 res
864 }
865
866 pub async fn get_payload_bodies_by_hash_with<F, R>(
868 &self,
869 hashes: Vec<BlockHash>,
870 f: F,
871 ) -> EngineApiResult<Vec<Option<R>>>
872 where
873 F: Fn(Provider::Block) -> R + Send + 'static,
874 R: Send + 'static,
875 {
876 let len = hashes.len() as u64;
877 if len > MAX_PAYLOAD_BODIES_LIMIT {
878 return Err(EngineApiError::PayloadRequestTooLarge { len });
879 }
880
881 let (tx, rx) = oneshot::channel();
882 let inner = self.inner.clone();
883
884 self.inner.task_spawner.spawn_blocking_task(async move {
885 let mut result = Vec::with_capacity(hashes.len());
886 for hash in hashes {
887 if tx.is_closed() {
888 return;
889 }
890
891 let block_result = inner.provider.block(BlockHashOrNumber::Hash(hash));
892 match block_result {
893 Ok(block) => {
894 result.push(block.map(&f));
895 }
896 Err(err) => {
897 let _ = tx.send(Err(EngineApiError::Internal(Box::new(err))));
898 return;
899 }
900 }
901 }
902 tx.send(Ok(result)).ok();
903 });
904
905 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
906 }
907
908 pub async fn get_payload_bodies_by_hash_with_timestamps(
910 &self,
911 hashes: Vec<BlockHash>,
912 include_bal: bool,
913 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
914 let bodies =
915 self.get_payload_bodies_by_hash_with(hashes, Self::payload_body_with_timestamp).await?;
916 self.attach_payload_body_bals(bodies, include_bal).await
917 }
918
919 pub async fn get_payload_bodies_by_hash_with_timestamps_metered(
921 &self,
922 hashes: Vec<BlockHash>,
923 include_bal: bool,
924 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
925 let start = Instant::now();
926 let result = self.get_payload_bodies_by_hash_with_timestamps(hashes, include_bal).await;
927 let latency = &self.inner.metrics.latency;
928 if include_bal {
929 latency.get_payload_bodies_by_hash_v2.record(start.elapsed());
930 } else {
931 latency.get_payload_bodies_by_hash_v1.record(start.elapsed());
932 }
933 result
934 }
935
936 fn payload_body_with_timestamp(
937 block: Provider::Block,
938 ) -> (BlockHash, u64, ExecutionPayloadBodyV2) {
939 (
940 block.header().hash_slow(),
941 block.header().timestamp(),
942 ExecutionPayloadBodyV2 {
943 transactions: block.body().encoded_2718_transactions(),
944 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
945 block_access_list: None,
946 },
947 )
948 }
949
950 async fn attach_payload_body_bals(
951 &self,
952 mut bodies: Vec<Option<(BlockHash, u64, ExecutionPayloadBodyV2)>>,
953 include_bal: bool,
954 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
955 if include_bal {
956 let hashes = bodies.iter().flatten().map(|(hash, _, _)| *hash).collect();
957 let bals = self.get_block_access_lists_by_hashes(hashes).await?;
958 for ((_, _, body), bal) in bodies.iter_mut().flatten().zip(bals) {
959 body.block_access_list = bal;
960 }
961 }
962 Ok(bodies
963 .into_iter()
964 .map(|body| body.map(|(_, timestamp, body)| (timestamp, body)))
965 .collect())
966 }
967
968 async fn get_block_access_lists_by_hashes(
969 &self,
970 hashes: Vec<BlockHash>,
971 ) -> EngineApiResult<Vec<Option<Bytes>>> {
972 let len = hashes.len() as u64;
973 if len > MAX_PAYLOAD_BODIES_LIMIT {
974 return Err(EngineApiError::PayloadRequestTooLarge { len });
975 }
976
977 let (tx, rx) = oneshot::channel();
978 let inner = self.inner.clone();
979
980 self.inner.task_spawner.spawn_blocking_task(async move {
981 if tx.is_closed() {
982 return;
983 }
984
985 tx.send(
986 inner
987 .provider
988 .get_bals_by_hashes(&hashes)
989 .map_err(|err| EngineApiError::Internal(Box::new(err))),
990 )
991 .ok();
992 });
993
994 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
995 }
996
997 pub async fn get_payload_bodies_by_hash_v1(
999 &self,
1000 hashes: Vec<BlockHash>,
1001 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
1002 self.get_payload_bodies_by_hash_with(hashes, |block| ExecutionPayloadBodyV1 {
1003 transactions: block.body().encoded_2718_transactions(),
1004 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
1005 })
1006 .await
1007 }
1008
1009 pub async fn get_payload_bodies_by_hash_v1_metered(
1011 &self,
1012 hashes: Vec<BlockHash>,
1013 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
1014 let start = Instant::now();
1015 let res = Self::get_payload_bodies_by_hash_v1(self, hashes).await;
1016 self.inner.metrics.latency.get_payload_bodies_by_hash_v1.record(start.elapsed());
1017 res
1018 }
1019
1020 pub async fn get_payload_bodies_by_hash_v2(
1024 &self,
1025 hashes: Vec<BlockHash>,
1026 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
1027 let payload_bodies =
1028 self.get_payload_bodies_by_hash_with(hashes.clone(), |block| ExecutionPayloadBodyV2 {
1029 transactions: block.body().encoded_2718_transactions(),
1030 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
1031 block_access_list: None,
1032 });
1033 let block_access_lists = self.get_block_access_lists_by_hashes(hashes);
1034 let (mut payload_bodies, block_access_lists) =
1035 tokio::try_join!(payload_bodies, block_access_lists)?;
1036
1037 for (payload_body, block_access_list) in payload_bodies.iter_mut().zip(block_access_lists) {
1038 if let Some(payload_body) = payload_body {
1039 payload_body.block_access_list = block_access_list;
1040 }
1041 }
1042
1043 Ok(payload_bodies)
1044 }
1045
1046 pub async fn get_payload_bodies_by_hash_v2_metered(
1048 &self,
1049 hashes: Vec<BlockHash>,
1050 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
1051 let start = Instant::now();
1052 let res = Self::get_payload_bodies_by_hash_v2(self, hashes).await;
1053 self.inner.metrics.latency.get_payload_bodies_by_hash_v2.record(start.elapsed());
1054 res
1055 }
1056
1057 async fn validate_and_execute_forkchoice(
1071 &self,
1072 version: EngineApiMessageVersion,
1073 state: ForkchoiceState,
1074 payload_attrs: Option<EngineT::PayloadAttributes>,
1075 ) -> EngineApiResult<ForkchoiceUpdated> {
1076 if let Some(ref attrs) = payload_attrs {
1077 let attr_validation_res =
1078 self.inner.validator.ensure_well_formed_attributes(version, attrs);
1079
1080 if let Err(err) = attr_validation_res {
1090 let fcu_res = self.inner.beacon_consensus.fork_choice_updated(state, None).await?;
1091 if fcu_res.is_invalid() || fcu_res.payload_status.is_syncing() {
1092 return Ok(fcu_res)
1093 }
1094 return Err(err.into())
1095 }
1096 }
1097
1098 Ok(self.inner.beacon_consensus.fork_choice_updated(state, payload_attrs).await?)
1099 }
1100
1101 pub fn capabilities(&self) -> &EngineCapabilities {
1103 &self.inner.capabilities
1104 }
1105
1106 fn has_blobs(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1107 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1108 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1109 }
1110
1111 self.inner
1112 .tx_pool
1113 .has_blobs_for_versioned_hashes(&versioned_hashes)
1114 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1115 }
1116
1117 pub fn has_blobs_metered(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1119 let start = Instant::now();
1120 let res = Self::has_blobs(self, versioned_hashes);
1121 self.inner.metrics.latency.has_blobs.record(start.elapsed());
1122 res
1123 }
1124
1125 fn get_blobs_v1(
1126 &self,
1127 versioned_hashes: Vec<B256>,
1128 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1129 let current_timestamp =
1131 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1132 if self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1133 return Err(EngineApiError::EngineObjectValidationError(
1134 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1135 ));
1136 }
1137
1138 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1139 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1140 }
1141
1142 self.inner
1143 .tx_pool
1144 .get_blobs_for_versioned_hashes_v1(&versioned_hashes)
1145 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1146 }
1147
1148 pub fn get_blobs_v1_metered(
1150 &self,
1151 versioned_hashes: Vec<B256>,
1152 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1153 let hashes_len = versioned_hashes.len();
1154 let start = Instant::now();
1155 let res = Self::get_blobs_v1(self, versioned_hashes);
1156 self.inner.metrics.latency.get_blobs_v1.record(start.elapsed());
1157
1158 if let Ok(blobs) = &res {
1159 let blobs_found = blobs.iter().flatten().count();
1160 let blobs_missed = hashes_len - blobs_found;
1161
1162 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1163 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1164 }
1165
1166 res
1167 }
1168
1169 fn get_blobs_v2(
1170 &self,
1171 versioned_hashes: Vec<B256>,
1172 ) -> EngineApiResult<Option<Vec<BlobAndProofV2>>> {
1173 let current_timestamp =
1175 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1176 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1177 return Err(EngineApiError::EngineObjectValidationError(
1178 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1179 ));
1180 }
1181
1182 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1183 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1184 }
1185
1186 self.inner
1187 .tx_pool
1188 .get_blobs_for_versioned_hashes_v2(&versioned_hashes)
1189 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1190 }
1191
1192 fn get_blobs_v3(
1193 &self,
1194 versioned_hashes: Vec<B256>,
1195 ) -> EngineApiResult<Option<Vec<Option<BlobAndProofV2>>>> {
1196 let current_timestamp =
1198 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1199 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1200 return Err(EngineApiError::EngineObjectValidationError(
1201 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1202 ));
1203 }
1204
1205 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1206 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1207 }
1208
1209 if (*self.inner.is_syncing)() {
1211 return Ok(None)
1212 }
1213
1214 self.inner
1215 .tx_pool
1216 .get_blobs_for_versioned_hashes_v3(&versioned_hashes)
1217 .map(Some)
1218 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1219 }
1220
1221 fn get_blobs_v4(
1222 &self,
1223 versioned_hashes: Vec<B256>,
1224 indices_bitarray: B128,
1225 ) -> EngineApiResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1226 let indices_bitarray = B128::from(u128::from_le_bytes(indices_bitarray.into()));
1229 let current_timestamp =
1230 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1231 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1232 return Err(EngineApiError::EngineObjectValidationError(
1233 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1234 ));
1235 }
1236
1237 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1238 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1239 }
1240
1241 if (*self.inner.is_syncing)() {
1243 return Ok(None)
1244 }
1245
1246 self.inner
1247 .tx_pool
1248 .get_blobs_for_versioned_hashes_v4(&versioned_hashes, indices_bitarray)
1249 .map(Some)
1250 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1251 }
1252
1253 pub fn get_blobs_v2_metered(
1255 &self,
1256 versioned_hashes: Vec<B256>,
1257 ) -> EngineApiResult<Option<Vec<BlobAndProofV2>>> {
1258 let hashes_len = versioned_hashes.len();
1259 let start = Instant::now();
1260 let res = Self::get_blobs_v2(self, versioned_hashes);
1261 self.inner.metrics.latency.get_blobs_v2.record(start.elapsed());
1262
1263 if let Ok(blobs) = &res {
1264 let blobs_found = blobs.iter().flatten().count();
1265
1266 self.inner
1267 .metrics
1268 .blob_metrics
1269 .get_blobs_requests_blobs_total
1270 .increment(hashes_len as u64);
1271 self.inner
1272 .metrics
1273 .blob_metrics
1274 .get_blobs_requests_blobs_in_blobpool_total
1275 .increment(blobs_found as u64);
1276
1277 if blobs_found == hashes_len {
1278 self.inner.metrics.blob_metrics.get_blobs_requests_success_total.increment(1);
1279 } else {
1280 self.inner.metrics.blob_metrics.get_blobs_requests_failure_total.increment(1);
1281 }
1282 } else {
1283 self.inner.metrics.blob_metrics.get_blobs_requests_failure_total.increment(1);
1284 }
1285
1286 res
1287 }
1288
1289 pub fn get_blobs_v3_metered(
1291 &self,
1292 versioned_hashes: Vec<B256>,
1293 ) -> EngineApiResult<Option<Vec<Option<BlobAndProofV2>>>> {
1294 let hashes_len = versioned_hashes.len();
1295 let start = Instant::now();
1296 let res = Self::get_blobs_v3(self, versioned_hashes);
1297 self.inner.metrics.latency.get_blobs_v3.record(start.elapsed());
1298
1299 if let Ok(Some(blobs)) = &res {
1300 let blobs_found = blobs.iter().flatten().count();
1301 let blobs_missed = hashes_len - blobs_found;
1302
1303 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1304 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1305 }
1306
1307 res
1308 }
1309
1310 pub fn get_blobs_v4_metered(
1312 &self,
1313 versioned_hashes: Vec<B256>,
1314 indices_bitarray: B128,
1315 ) -> EngineApiResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1316 let hashes_len = versioned_hashes.len();
1317 let start = Instant::now();
1318 let res = Self::get_blobs_v4(self, versioned_hashes, indices_bitarray);
1319 self.inner.metrics.latency.get_blobs_v4.record(start.elapsed());
1320
1321 if let Ok(Some(blobs)) = &res {
1322 let blobs_found = blobs.iter().flatten().count();
1323 let blobs_missed = hashes_len - blobs_found;
1324
1325 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1326 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1327 }
1328
1329 res
1330 }
1331}
1332
1333#[async_trait]
1335impl<Provider, EngineT, Pool, Validator, ChainSpec> EngineApiServer<EngineT>
1336 for EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
1337where
1338 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
1339 EngineT: EngineTypes<ExecutionData = ExecutionData>,
1340 Pool: TransactionPool + 'static,
1341 Validator: EngineApiValidator<EngineT>,
1342 ChainSpec: EthereumHardforks + Send + Sync + 'static,
1343{
1344 async fn new_payload_v1(&self, payload: ExecutionPayloadV1) -> RpcResult<PayloadStatus> {
1348 trace!(target: "rpc::engine", "Serving engine_newPayloadV1");
1349 let payload =
1350 ExecutionData { payload: payload.into(), sidecar: ExecutionPayloadSidecar::none() };
1351 Ok(self.new_payload_v1_metered(payload).await?)
1352 }
1353
1354 async fn new_payload_v2(&self, payload: ExecutionPayloadInputV2) -> RpcResult<PayloadStatus> {
1357 trace!(target: "rpc::engine", "Serving engine_newPayloadV2");
1358 let payload = ExecutionData {
1359 payload: payload.into_payload(),
1360 sidecar: ExecutionPayloadSidecar::none(),
1361 };
1362
1363 Ok(self.new_payload_v2_metered(payload).await?)
1364 }
1365
1366 async fn new_payload_v3(
1369 &self,
1370 payload: ExecutionPayloadV3,
1371 versioned_hashes: Vec<B256>,
1372 parent_beacon_block_root: B256,
1373 ) -> RpcResult<PayloadStatus> {
1374 trace!(target: "rpc::engine", "Serving engine_newPayloadV3");
1375 let payload = ExecutionData {
1376 payload: payload.into(),
1377 sidecar: ExecutionPayloadSidecar::v3(CancunPayloadFields {
1378 versioned_hashes,
1379 parent_beacon_block_root,
1380 }),
1381 };
1382
1383 Ok(self.new_payload_v3_metered(payload).await?)
1384 }
1385
1386 async fn new_payload_v4(
1389 &self,
1390 payload: ExecutionPayloadV3,
1391 versioned_hashes: Vec<B256>,
1392 parent_beacon_block_root: B256,
1393 requests: RequestsOrHash,
1394 ) -> RpcResult<PayloadStatus> {
1395 trace!(target: "rpc::engine", "Serving engine_newPayloadV4");
1396
1397 if requests.is_hash() && !self.inner.accept_execution_requests_hash {
1399 return Err(EngineApiError::UnexpectedRequestsHash.into());
1400 }
1401
1402 let payload = ExecutionData {
1403 payload: payload.into(),
1404 sidecar: ExecutionPayloadSidecar::v4(
1405 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1406 PraguePayloadFields { requests },
1407 ),
1408 };
1409
1410 Ok(self.new_payload_v4_metered(payload).await?)
1411 }
1412
1413 async fn new_payload_v5(
1419 &self,
1420 payload: ExecutionPayloadV4,
1421 versioned_hashes: Vec<B256>,
1422 parent_beacon_block_root: B256,
1423 requests: RequestsOrHash,
1424 ) -> RpcResult<PayloadStatus> {
1425 trace!(target: "rpc::engine", "Serving engine_newPayloadV5");
1426 if requests.is_hash() && !self.inner.accept_execution_requests_hash {
1428 return Err(EngineApiError::UnexpectedRequestsHash.into());
1429 }
1430
1431 let payload = ExecutionData {
1432 payload: payload.into(),
1433 sidecar: ExecutionPayloadSidecar::v4(
1434 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1435 PraguePayloadFields { requests },
1436 ),
1437 };
1438
1439 Ok(self.new_payload_v5_metered(payload).await?)
1440 }
1441
1442 async fn new_payload_v6(
1446 &self,
1447 payload: ExecutionPayloadV4,
1448 versioned_hashes: Vec<B256>,
1449 parent_beacon_block_root: B256,
1450 execution_requests: RequestsOrHash,
1451 inclusion_list_transactions: Vec<Bytes>,
1452 ) -> RpcResult<PayloadStatusV2> {
1453 trace!(target: "rpc::engine", "Serving engine_newPayloadV6");
1454 if execution_requests.is_hash() && !self.inner.accept_execution_requests_hash {
1455 return Err(EngineApiError::UnexpectedRequestsHash.into());
1456 }
1457
1458 let payload = ExecutionData {
1459 payload: payload.into(),
1460 sidecar: ExecutionPayloadSidecar::v6(
1461 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1462 PraguePayloadFields { requests: execution_requests },
1463 BogotaPayloadFields { inclusion_list_transactions },
1464 ),
1465 };
1466
1467 Ok(self.new_payload_v6_metered(payload).await?)
1470 }
1471
1472 async fn fork_choice_updated_v1(
1477 &self,
1478 fork_choice_state: ForkchoiceState,
1479 payload_attributes: Option<EngineT::PayloadAttributes>,
1480 ) -> RpcResult<ForkchoiceUpdated> {
1481 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV1");
1482 Ok(self.fork_choice_updated_v1_metered(fork_choice_state, payload_attributes).await?)
1483 }
1484
1485 async fn fork_choice_updated_v2(
1488 &self,
1489 fork_choice_state: ForkchoiceState,
1490 payload_attributes: Option<EngineT::PayloadAttributes>,
1491 ) -> RpcResult<ForkchoiceUpdated> {
1492 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV2");
1493 Ok(self.fork_choice_updated_v2_metered(fork_choice_state, payload_attributes).await?)
1494 }
1495
1496 async fn fork_choice_updated_v3(
1500 &self,
1501 fork_choice_state: ForkchoiceState,
1502 payload_attributes: Option<EngineT::PayloadAttributes>,
1503 ) -> RpcResult<ForkchoiceUpdated> {
1504 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV3");
1505 Ok(self.fork_choice_updated_v3_metered(fork_choice_state, payload_attributes).await?)
1506 }
1507
1508 async fn fork_choice_updated_v4(
1512 &self,
1513 fork_choice_state: ForkchoiceState,
1514 payload_attributes: Option<EngineT::PayloadAttributes>,
1515 custody_columns: Option<B128>,
1516 ) -> RpcResult<ForkchoiceUpdated> {
1517 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV4");
1518 Ok(self
1519 .fork_choice_updated_v4_metered(fork_choice_state, payload_attributes, custody_columns)
1520 .await?)
1521 }
1522
1523 async fn fork_choice_updated_v5(
1527 &self,
1528 fork_choice_state: ForkchoiceState,
1529 payload_attributes: Option<EngineT::PayloadAttributes>,
1530 custody_columns: Option<B128>,
1531 ) -> RpcResult<ForkchoiceUpdatedResponseV2> {
1532 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV5");
1533 Ok(self
1534 .fork_choice_updated_v5_metered(fork_choice_state, payload_attributes, custody_columns)
1535 .await?)
1536 }
1537
1538 async fn get_payload_v1(
1550 &self,
1551 payload_id: PayloadId,
1552 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV1> {
1553 trace!(target: "rpc::engine", "Serving engine_getPayloadV1");
1554 Ok(self.get_payload_v1_metered(payload_id).await?)
1555 }
1556
1557 async fn get_payload_v2(
1567 &self,
1568 payload_id: PayloadId,
1569 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV2> {
1570 debug!(target: "rpc::engine", id = %payload_id, "Serving engine_getPayloadV2");
1571 Ok(self.get_payload_v2_metered(payload_id).await?)
1572 }
1573
1574 async fn get_payload_v3(
1584 &self,
1585 payload_id: PayloadId,
1586 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV3> {
1587 trace!(target: "rpc::engine", "Serving engine_getPayloadV3");
1588 Ok(self.get_payload_v3_metered(payload_id).await?)
1589 }
1590
1591 async fn get_payload_v4(
1601 &self,
1602 payload_id: PayloadId,
1603 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV4> {
1604 trace!(target: "rpc::engine", "Serving engine_getPayloadV4");
1605 Ok(self.get_payload_v4_metered(payload_id).await?)
1606 }
1607
1608 async fn get_payload_v5(
1618 &self,
1619 payload_id: PayloadId,
1620 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV5> {
1621 trace!(target: "rpc::engine", "Serving engine_getPayloadV5");
1622 Ok(self.get_payload_v5_metered(payload_id).await?)
1623 }
1624
1625 async fn get_payload_v6(
1631 &self,
1632 payload_id: PayloadId,
1633 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV6> {
1634 trace!(target: "rpc::engine", "Serving engine_getPayloadV6");
1635 Ok(self.get_payload_v6_metered(payload_id).await?)
1636 }
1637
1638 async fn get_inclusion_list_v1(&self) -> RpcResult<Vec<Bytes>> {
1642 trace!(target: "rpc::engine", "Serving engine_getInclusionListV1");
1643 Ok(self.get_inclusion_list_v1_metered()?)
1644 }
1645
1646 async fn get_payload_bodies_by_hash_v1(
1649 &self,
1650 block_hashes: Vec<BlockHash>,
1651 ) -> RpcResult<ExecutionPayloadBodiesV1> {
1652 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByHashV1");
1653 Ok(self.get_payload_bodies_by_hash_v1_metered(block_hashes).await?)
1654 }
1655
1656 async fn get_payload_bodies_by_hash_v2(
1662 &self,
1663 block_hashes: Vec<BlockHash>,
1664 ) -> RpcResult<ExecutionPayloadBodiesV2> {
1665 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByHashV2");
1666 Ok(self.get_payload_bodies_by_hash_v2_metered(block_hashes).await?)
1667 }
1668
1669 async fn get_payload_bodies_by_range_v1(
1686 &self,
1687 start: U64,
1688 count: U64,
1689 ) -> RpcResult<ExecutionPayloadBodiesV1> {
1690 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByRangeV1");
1691 Ok(self.get_payload_bodies_by_range_v1_metered(start.to(), count.to()).await?)
1692 }
1693
1694 async fn get_payload_bodies_by_range_v2(
1700 &self,
1701 start: U64,
1702 count: U64,
1703 ) -> RpcResult<ExecutionPayloadBodiesV2> {
1704 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByRangeV2");
1705 Ok(self.get_payload_bodies_by_range_v2_metered(start.to(), count.to()).await?)
1706 }
1707
1708 async fn get_client_version_v1(
1712 &self,
1713 client: ClientVersionV1,
1714 ) -> RpcResult<Vec<ClientVersionV1>> {
1715 trace!(target: "rpc::engine", "Serving engine_getClientVersionV1");
1716 Ok(Self::get_client_version_v1(self, client)?)
1717 }
1718
1719 async fn exchange_capabilities(&self, capabilities: Vec<String>) -> RpcResult<Vec<String>> {
1722 trace!(target: "rpc::engine", "Serving engine_exchangeCapabilities");
1723
1724 let el_caps = self.capabilities();
1725 el_caps.log_capability_mismatches(&capabilities);
1726
1727 Ok(el_caps.list())
1728 }
1729
1730 async fn has_blobs(&self, versioned_hashes: Vec<B256>) -> RpcResult<Vec<bool>> {
1731 trace!(target: "rpc::engine", "Serving engine_hasBlobs");
1732 Ok(self.has_blobs_metered(versioned_hashes)?)
1733 }
1734
1735 async fn get_blobs_v1(
1736 &self,
1737 versioned_hashes: Vec<B256>,
1738 ) -> RpcResult<Vec<Option<BlobAndProofV1>>> {
1739 trace!(target: "rpc::engine", "Serving engine_getBlobsV1");
1740 Ok(self.get_blobs_v1_metered(versioned_hashes)?)
1741 }
1742
1743 async fn get_blobs_v2(
1744 &self,
1745 versioned_hashes: Vec<B256>,
1746 ) -> RpcResult<Option<Vec<BlobAndProofV2>>> {
1747 trace!(target: "rpc::engine", "Serving engine_getBlobsV2");
1748 Ok(self.get_blobs_v2_metered(versioned_hashes)?)
1749 }
1750
1751 async fn get_blobs_v3(
1752 &self,
1753 versioned_hashes: Vec<B256>,
1754 ) -> RpcResult<Option<Vec<Option<BlobAndProofV2>>>> {
1755 trace!(target: "rpc::engine", "Serving engine_getBlobsV3");
1756 Ok(self.get_blobs_v3_metered(versioned_hashes)?)
1757 }
1758
1759 async fn get_blobs_v4(
1760 &self,
1761 versioned_hashes: Vec<B256>,
1762 indices_bitarray: B128,
1763 ) -> RpcResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1764 trace!(target: "rpc::engine", "Serving engine_getBlobsV4");
1765 Ok(self.get_blobs_v4_metered(versioned_hashes, indices_bitarray)?)
1766 }
1767}
1768
1769impl<Provider, EngineT, Pool, Validator, ChainSpec> IntoEngineApiRpcModule
1770 for EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
1771where
1772 EngineT: EngineTypes,
1773 Self: EngineApiServer<EngineT>,
1774{
1775 fn into_rpc_module(self) -> RpcModule<()> {
1776 EngineApiServer::<EngineT>::into_rpc(self).remove_context()
1777 }
1778}
1779
1780impl<Provider, PayloadT, Pool, Validator, ChainSpec> std::fmt::Debug
1781 for EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
1782where
1783 PayloadT: PayloadTypes,
1784{
1785 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1786 f.debug_struct("EngineApi").finish_non_exhaustive()
1787 }
1788}
1789
1790impl<Provider, PayloadT, Pool, Validator, ChainSpec> Clone
1791 for EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
1792where
1793 PayloadT: PayloadTypes,
1794{
1795 fn clone(&self) -> Self {
1796 Self { inner: Arc::clone(&self.inner) }
1797 }
1798}
1799
1800struct EngineApiInner<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
1802 provider: Provider,
1804 chain_spec: Arc<ChainSpec>,
1806 beacon_consensus: ConsensusEngineHandle<PayloadT>,
1808 payload_store: PayloadStore<PayloadT>,
1810 task_spawner: Runtime,
1812 metrics: EngineApiMetrics,
1814 client: ClientVersionV1,
1816 capabilities: EngineCapabilities,
1818 tx_pool: Pool,
1820 validator: Validator,
1822 accept_execution_requests_hash: bool,
1823 cell_custody: CellCustody,
1825 is_syncing: Arc<dyn Fn() -> bool + Send + Sync>,
1827}
1828
1829#[cfg(test)]
1830mod tests {
1831 use super::*;
1832 use alloy_eips::{eip7685::Requests, Encodable2718, NumHash};
1833 use alloy_primitives::{Address, Bytes, B256};
1834 use alloy_rpc_types_engine::{
1835 ClientCode, ClientVersionV1, ExecutionPayloadV2, PayloadAttributes, PayloadStatusEnum,
1836 };
1837 use assert_matches::assert_matches;
1838 use reth_chainspec::{ChainSpec, ChainSpecBuilder, MAINNET};
1839 use reth_engine_primitives::{BeaconEngineMessage, OnForkChoiceUpdated};
1840 use reth_ethereum_engine_primitives::EthEngineTypes;
1841 use reth_ethereum_primitives::{Block, TransactionSigned};
1842 use reth_network_api::{
1843 noop::NoopNetwork, EthProtocolInfo, NetworkError, NetworkInfo, NetworkStatus,
1844 };
1845 use reth_node_ethereum::EthereumEngineValidator;
1846 use reth_payload_builder::test_utils::spawn_test_payload_service;
1847 use reth_primitives_traits::SignedTransaction;
1848 use reth_provider::{test_utils::MockEthProvider, BalStoreHandle, InMemoryBalStore, RawBal};
1849 use reth_tasks::Runtime;
1850 use reth_transaction_pool::{
1851 blobstore::InMemoryBlobStore,
1852 noop::NoopTransactionPool,
1853 test_utils::{OkValidator, TransactionBuilder},
1854 CoinbaseTipOrdering, EthPooledTransaction, Pool, PoolTransaction, TransactionOrigin,
1855 };
1856 use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver};
1857
1858 type EthTestPool = Pool<
1859 OkValidator<EthPooledTransaction>,
1860 CoinbaseTipOrdering<EthPooledTransaction>,
1861 InMemoryBlobStore,
1862 >;
1863
1864 fn eth_test_pool() -> EthTestPool {
1865 Pool::new(
1866 OkValidator::default(),
1867 CoinbaseTipOrdering::default(),
1868 InMemoryBlobStore::default(),
1869 Default::default(),
1870 )
1871 }
1872
1873 fn pooled_transaction(transaction: TransactionSigned) -> EthPooledTransaction {
1874 let transaction = transaction.try_into_recovered().unwrap();
1875 let encoded_length = transaction.encode_2718_len();
1876 EthPooledTransaction::new(transaction, encoded_length)
1877 }
1878
1879 fn setup_engine_api() -> (
1880 EngineApiTestHandle,
1881 EngineApi<
1882 Arc<MockEthProvider>,
1883 EthEngineTypes,
1884 NoopTransactionPool,
1885 EthereumEngineValidator,
1886 ChainSpec,
1887 >,
1888 ) {
1889 setup_engine_api_with_pool(NoopTransactionPool::default())
1890 }
1891
1892 fn setup_engine_api_with_pool<Pool>(
1893 tx_pool: Pool,
1894 ) -> (
1895 EngineApiTestHandle,
1896 EngineApi<Arc<MockEthProvider>, EthEngineTypes, Pool, EthereumEngineValidator, ChainSpec>,
1897 )
1898 where
1899 Pool: TransactionPool + 'static,
1900 {
1901 setup_engine_api_with_provider(tx_pool, MockEthProvider::default())
1902 }
1903
1904 fn setup_engine_api_with_provider<Pool>(
1905 tx_pool: Pool,
1906 provider: MockEthProvider,
1907 ) -> (
1908 EngineApiTestHandle,
1909 EngineApi<Arc<MockEthProvider>, EthEngineTypes, Pool, EthereumEngineValidator, ChainSpec>,
1910 )
1911 where
1912 Pool: TransactionPool + 'static,
1913 {
1914 let client = ClientVersionV1 {
1915 code: ClientCode::RH,
1916 name: "Reth".to_string(),
1917 version: "v0.2.0-beta.5".to_string(),
1918 commit: "defa64b2".to_string(),
1919 };
1920
1921 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
1922 let provider = Arc::new(provider);
1923 let payload_store = spawn_test_payload_service();
1924 let (to_engine, engine_rx) = unbounded_channel();
1925 let task_executor = Runtime::test();
1926 let api = EngineApi::new(
1927 provider.clone(),
1928 chain_spec.clone(),
1929 ConsensusEngineHandle::new(to_engine),
1930 payload_store.into(),
1931 tx_pool,
1932 task_executor,
1933 client,
1934 EngineCapabilities::default(),
1935 EthereumEngineValidator::new(chain_spec.clone()),
1936 false,
1937 NoopNetwork::default(),
1938 );
1939 let handle = EngineApiTestHandle { chain_spec, provider, from_api: engine_rx };
1940 (handle, api)
1941 }
1942
1943 #[tokio::test]
1944 async fn engine_client_version_v1() {
1945 let client = ClientVersionV1 {
1946 code: ClientCode::RH,
1947 name: "Reth".to_string(),
1948 version: "v0.2.0-beta.5".to_string(),
1949 commit: "defa64b2".to_string(),
1950 };
1951 let (_, api) = setup_engine_api();
1952 let res = api.get_client_version_v1(client.clone());
1953 assert_eq!(res.unwrap(), vec![client]);
1954 }
1955
1956 #[tokio::test]
1957 async fn get_inclusion_list_v1_returns_empty_list() {
1958 let (_, api) = setup_engine_api();
1959
1960 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
1961 assert!(res.is_empty());
1962 }
1963
1964 #[tokio::test]
1965 async fn get_inclusion_list_v1_stops_at_size_limit() {
1966 let pool = eth_test_pool();
1967 let first = pooled_transaction(
1968 TransactionBuilder::default()
1969 .max_fee_per_gas(3_000_000_000u128)
1970 .input(vec![0; 4_200])
1971 .into_legacy(),
1972 );
1973 let second = pooled_transaction(
1974 TransactionBuilder::default()
1975 .max_fee_per_gas(2_000_000_000u128)
1976 .input(vec![0; 4_200])
1977 .into_legacy(),
1978 );
1979 let third = pooled_transaction(
1980 TransactionBuilder::default().max_fee_per_gas(1_000_000_000u128).into_legacy(),
1981 );
1982 let expected = first.encoded_2718_consensus();
1983
1984 pool.add_transaction(TransactionOrigin::External, first).await.unwrap();
1985 pool.add_transaction(TransactionOrigin::External, second).await.unwrap();
1986 pool.add_transaction(TransactionOrigin::External, third).await.unwrap();
1987 let (_, api) = setup_engine_api_with_pool(pool);
1988
1989 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
1990 assert_eq!(res, vec![expected]);
1991 assert!(
1992 res.iter().map(|tx| tx.len()).sum::<usize>() <= MAX_BYTES_PER_INCLUSION_LIST as usize
1993 );
1994 }
1995
1996 fn inclusion_list_transaction(size: usize, typed: bool, fee: u128) -> EthPooledTransaction {
1998 for input_size in size.saturating_sub(256)..size {
1999 let builder = TransactionBuilder::default()
2000 .signer(B256::with_last_byte(if typed { 2 } else { 1 }))
2001 .max_fee_per_gas(fee)
2002 .max_priority_fee_per_gas(fee)
2003 .input(vec![0; input_size]);
2004 let transaction = if typed { builder.into_eip1559() } else { builder.into_legacy() };
2005 if transaction.encode_2718_len() == size {
2006 return pooled_transaction(transaction)
2007 }
2008 }
2009 panic!("could not construct a transaction of {size} bytes")
2010 }
2011
2012 #[tokio::test]
2013 async fn get_inclusion_list_v1_transaction_byte_boundaries() {
2014 let limit = MAX_BYTES_PER_INCLUSION_LIST as usize;
2015 for typed in [false, true] {
2016 for size in [limit - 3, limit, limit + 1] {
2017 let pool = eth_test_pool();
2018 let transaction = inclusion_list_transaction(size, typed, 3_000_000_000);
2019 let encoded = transaction.encoded_2718_consensus();
2020 assert_eq!(encoded.len(), size);
2021 pool.add_transaction(TransactionOrigin::External, transaction).await.unwrap();
2022 let (_, api) = setup_engine_api_with_pool(pool);
2023
2024 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
2025 let expected = if size <= limit { vec![encoded] } else { vec![] };
2026 assert_eq!(res, expected, "typed={typed}, size={size}");
2027 assert!(res.iter().all(|tx| !tx.is_empty()));
2028 }
2029 }
2030 }
2031
2032 #[tokio::test]
2033 async fn get_inclusion_list_v1_cumulative_transaction_byte_boundaries() {
2034 let limit = MAX_BYTES_PER_INCLUSION_LIST as usize;
2035 for extra in [0, 1] {
2036 let pool = eth_test_pool();
2037 let first = inclusion_list_transaction(limit / 2, false, 3_000_000_000);
2038 let second = inclusion_list_transaction(limit / 2 + extra, true, 2_000_000_000);
2039 let mut expected =
2040 vec![first.encoded_2718_consensus(), second.encoded_2718_consensus()];
2041 assert_eq!(expected.iter().map(|tx| tx.len()).sum::<usize>(), limit + extra);
2042 if extra != 0 {
2043 expected.pop();
2044 }
2045 pool.add_transaction(TransactionOrigin::External, first).await.unwrap();
2046 pool.add_transaction(TransactionOrigin::External, second).await.unwrap();
2047 let (_, api) = setup_engine_api_with_pool(pool);
2048
2049 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
2050 assert_eq!(res, expected);
2051 assert!(res.iter().all(|tx| !tx.is_empty()));
2052 }
2053 }
2054
2055 #[tokio::test]
2056 async fn get_inclusion_list_v1_excludes_blob_transactions() {
2057 let pool = eth_test_pool();
2058 let blob = pooled_transaction(
2059 TransactionBuilder::default()
2060 .max_fee_per_gas(2_000_000_000u128)
2061 .max_priority_fee_per_gas(1_000_000_000u128)
2062 .into_eip4844(),
2063 );
2064 let non_blob = pooled_transaction(
2065 TransactionBuilder::default()
2066 .max_fee_per_gas(1_000_000_000u128)
2067 .max_priority_fee_per_gas(1_000_000_000u128)
2068 .into_eip1559(),
2069 );
2070 let expected = non_blob.encoded_2718_consensus();
2071
2072 pool.add_transaction(TransactionOrigin::External, blob).await.unwrap();
2073 pool.add_transaction(TransactionOrigin::External, non_blob).await.unwrap();
2074 let (_, api) = setup_engine_api_with_pool(pool);
2075
2076 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
2077 assert_eq!(res, vec![expected]);
2078 }
2079
2080 #[tokio::test]
2081 async fn has_blobs_returns_ordered_availability() {
2082 let (_, api) = setup_engine_api();
2083
2084 let res = api.has_blobs_metered(vec![B256::ZERO, B256::with_last_byte(1)]).unwrap();
2085 assert_eq!(res, vec![false, false]);
2086 }
2087
2088 #[tokio::test]
2089 async fn has_blobs_rejects_large_requests() {
2090 let (_, api) = setup_engine_api();
2091
2092 let res = api.has_blobs_metered(vec![B256::ZERO; MAX_BLOB_LIMIT + 1]);
2093 assert_matches!(
2094 res,
2095 Err(EngineApiError::BlobRequestTooLarge { len }) if len == MAX_BLOB_LIMIT + 1
2096 );
2097 }
2098
2099 #[tokio::test]
2100 async fn payload_bodies_with_timestamps_preserve_missing_entries_and_head_truncation() {
2101 let mut provider = MockEthProvider::default();
2102 provider.bal_store = BalStoreHandle::new(InMemoryBalStore::default());
2103 let (handle, api) =
2104 setup_engine_api_with_provider(NoopTransactionPool::default(), provider);
2105 let mut block = Block::default();
2106 block.header.number = 1;
2107 block.header.timestamp = 123;
2108 let hash = block.header.hash_slow();
2109 handle.provider.add_block(hash, block);
2110 let raw_bal = Bytes::from_static(&[alloy_rlp::EMPTY_LIST_CODE]);
2111 handle
2112 .provider
2113 .bal_store
2114 .insert(NumHash::new(1, hash), RawBal::new(raw_bal.clone()))
2115 .unwrap();
2116 for include_bal in [false, true] {
2117 let by_hash = api
2118 .get_payload_bodies_by_hash_with_timestamps(
2119 vec![B256::ZERO, hash, hash],
2120 include_bal,
2121 )
2122 .await
2123 .unwrap();
2124 assert_eq!(by_hash.len(), 3);
2125 assert!(by_hash[0].is_none());
2126 assert_eq!(by_hash[1], by_hash[2]);
2127 let (timestamp, body) = by_hash[1].as_ref().unwrap();
2128 assert_eq!(*timestamp, 123);
2129 assert_eq!(body.block_access_list, include_bal.then(|| raw_bal.clone()));
2130 let by_range =
2131 api.get_payload_bodies_by_range_with_timestamps(1, 3, include_bal).await.unwrap();
2132 assert_eq!(by_range, vec![by_hash[1].clone()]);
2133 assert!(api
2134 .get_payload_bodies_by_range_with_timestamps(2, 3, include_bal)
2135 .await
2136 .unwrap()
2137 .is_empty());
2138 }
2139 }
2140
2141 #[tokio::test]
2142 async fn get_payload_bodies_by_hash_v2_returns_block_access_list_from_store() {
2143 let bal_store = BalStoreHandle::new(InMemoryBalStore::default());
2144 let mut provider = MockEthProvider::default();
2145 provider.bal_store = bal_store.clone();
2146 let provider = Arc::new(provider);
2147
2148 let client = ClientVersionV1 {
2149 code: ClientCode::RH,
2150 name: "Reth".to_string(),
2151 version: "v0.2.0-beta.5".to_string(),
2152 commit: "defa64b2".to_string(),
2153 };
2154 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
2155 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2156 let (to_engine, _engine_rx) = unbounded_channel();
2157 let api = EngineApi::new(
2158 provider.clone(),
2159 chain_spec.clone(),
2160 ConsensusEngineHandle::new(to_engine),
2161 payload_store.into(),
2162 NoopTransactionPool::default(),
2163 Runtime::test(),
2164 client,
2165 EngineCapabilities::default(),
2166 EthereumEngineValidator::new(chain_spec),
2167 false,
2168 NoopNetwork::default(),
2169 );
2170
2171 let mut block = Block::default();
2172 block.header.number = 1;
2173 let block_hash = block.header.hash_slow();
2174 provider.add_block(block_hash, block);
2175
2176 let mut block_without_bal = Block::default();
2177 block_without_bal.header.number = 2;
2178 let block_without_bal_hash = block_without_bal.header.hash_slow();
2179 provider.add_block(block_without_bal_hash, block_without_bal);
2180
2181 let raw_bal = Bytes::from_static(&[alloy_rlp::EMPTY_LIST_CODE]);
2182 bal_store.insert(NumHash::new(1, block_hash), RawBal::new(raw_bal.clone())).unwrap();
2183
2184 let missing_hash = B256::with_last_byte(3);
2185 let response = api
2186 .get_payload_bodies_by_hash_v2(vec![block_hash, block_without_bal_hash, missing_hash])
2187 .await
2188 .unwrap();
2189
2190 assert_eq!(response.len(), 3);
2191 assert_eq!(response[0].as_ref().unwrap().block_access_list, Some(raw_bal));
2192 assert_eq!(response[1].as_ref().unwrap().block_access_list, None);
2193 assert!(response[2].is_none());
2194 }
2195
2196 #[tokio::test]
2197 async fn get_payload_bodies_by_range_v2_returns_block_access_lists_from_store() {
2198 let bal_store = BalStoreHandle::new(InMemoryBalStore::default());
2199 let mut provider = MockEthProvider::default();
2200 provider.bal_store = bal_store.clone();
2201 let provider = Arc::new(provider);
2202
2203 let client = ClientVersionV1 {
2204 code: ClientCode::RH,
2205 name: "Reth".to_string(),
2206 version: "v0.2.0-beta.5".to_string(),
2207 commit: "defa64b2".to_string(),
2208 };
2209 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
2210 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2211 let (to_engine, _engine_rx) = unbounded_channel();
2212 let api = EngineApi::new(
2213 provider.clone(),
2214 chain_spec.clone(),
2215 ConsensusEngineHandle::new(to_engine),
2216 payload_store.into(),
2217 NoopTransactionPool::default(),
2218 Runtime::test(),
2219 client,
2220 EngineCapabilities::default(),
2221 EthereumEngineValidator::new(chain_spec),
2222 false,
2223 NoopNetwork::default(),
2224 );
2225
2226 let mut block = Block::default();
2227 block.header.number = 1;
2228 let block_hash = block.header.hash_slow();
2229 provider.add_block(block_hash, block);
2230
2231 let mut block_without_bal = Block::default();
2232 block_without_bal.header.number = 3;
2233 let block_without_bal_hash = block_without_bal.header.hash_slow();
2234 provider.add_block(block_without_bal_hash, block_without_bal);
2235
2236 let raw_bal = Bytes::from_static(&[alloy_rlp::EMPTY_LIST_CODE]);
2237 bal_store.insert(NumHash::new(1, block_hash), RawBal::new(raw_bal.clone())).unwrap();
2238
2239 let response = api.get_payload_bodies_by_range_v2(1, 3).await.unwrap();
2240
2241 assert_eq!(response.len(), 3);
2242 assert_eq!(response[0].as_ref().unwrap().block_access_list, Some(raw_bal));
2243 assert!(response[1].is_none());
2244 assert_eq!(response[2].as_ref().unwrap().block_access_list, None);
2245 }
2246
2247 struct EngineApiTestHandle {
2248 #[allow(dead_code)]
2249 chain_spec: Arc<ChainSpec>,
2250 provider: Arc<MockEthProvider>,
2251 from_api: UnboundedReceiver<BeaconEngineMessage<EthEngineTypes>>,
2252 }
2253
2254 #[tokio::test]
2255 async fn forwards_responses_to_consensus_engine() {
2256 let (mut handle, api) = setup_engine_api();
2257
2258 tokio::spawn(async move {
2259 let payload_v1 = ExecutionPayloadV1::from_block_slow(&Block::default());
2260 let execution_data = ExecutionData {
2261 payload: payload_v1.into(),
2262 sidecar: ExecutionPayloadSidecar::none(),
2263 };
2264
2265 api.new_payload_v1(execution_data).await.unwrap();
2266 });
2267 assert_matches!(handle.from_api.recv().await, Some(BeaconEngineMessage::NewPayload { .. }));
2268 }
2269
2270 #[tokio::test]
2271 async fn new_payload_v5_accepts_amsterdam_payloads() {
2272 let chain_spec = Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2273 let provider = Arc::new(MockEthProvider::default());
2274 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2275 let (to_engine, mut engine_rx) = unbounded_channel();
2276
2277 let api = EngineApi::new(
2278 provider,
2279 chain_spec.clone(),
2280 ConsensusEngineHandle::new(to_engine),
2281 payload_store.into(),
2282 NoopTransactionPool::default(),
2283 Runtime::test(),
2284 ClientVersionV1 {
2285 code: ClientCode::RH,
2286 name: "Reth".to_string(),
2287 version: "v0.0.0-test".to_string(),
2288 commit: "test".to_string(),
2289 },
2290 EngineCapabilities::default(),
2291 EthereumEngineValidator::new(chain_spec),
2292 false,
2293 NoopNetwork::default(),
2294 );
2295
2296 tokio::spawn(async move {
2297 let payload_v1 = ExecutionPayloadV1::from_block_slow(&Block::default());
2298 let payload = ExecutionPayloadV4 {
2299 payload_inner: ExecutionPayloadV3 {
2300 payload_inner: ExecutionPayloadV2 {
2301 payload_inner: payload_v1,
2302 withdrawals: Vec::new(),
2303 },
2304 blob_gas_used: 0,
2305 excess_blob_gas: 0,
2306 },
2307 block_access_list: Bytes::from_static(b"bal"),
2308 slot_number: 1,
2309 };
2310 let execution_data = ExecutionData {
2311 payload: payload.into(),
2312 sidecar: ExecutionPayloadSidecar::v4(
2313 CancunPayloadFields {
2314 versioned_hashes: Vec::new(),
2315 parent_beacon_block_root: B256::ZERO,
2316 },
2317 PraguePayloadFields { requests: RequestsOrHash::Requests(Requests::default()) },
2318 ),
2319 };
2320
2321 api.new_payload_v5(execution_data).await.unwrap();
2322 });
2323
2324 assert_matches!(engine_rx.recv().await, Some(BeaconEngineMessage::NewPayload { .. }));
2325 }
2326
2327 #[derive(Clone)]
2328 struct TestNetworkInfo {
2329 syncing: bool,
2330 }
2331
2332 impl NetworkInfo for TestNetworkInfo {
2333 fn local_addr(&self) -> std::net::SocketAddr {
2334 (std::net::Ipv4Addr::UNSPECIFIED, 0).into()
2335 }
2336
2337 async fn network_status(&self) -> Result<NetworkStatus, NetworkError> {
2338 #[allow(deprecated)]
2339 Ok(NetworkStatus {
2340 client_version: "test".to_string(),
2341 protocol_version: 5,
2342 eth_protocol_info: EthProtocolInfo {
2343 network: 1,
2344 difficulty: None,
2345 genesis: Default::default(),
2346 config: Default::default(),
2347 head: Default::default(),
2348 },
2349 capabilities: vec![],
2350 })
2351 }
2352
2353 fn chain_id(&self) -> u64 {
2354 1
2355 }
2356
2357 fn cell_custody(&self) -> &CellCustody {
2358 static CELL_CUSTODY: std::sync::OnceLock<CellCustody> = std::sync::OnceLock::new();
2359 CELL_CUSTODY.get_or_init(CellCustody::default)
2360 }
2361
2362 fn is_syncing(&self) -> bool {
2363 self.syncing
2364 }
2365
2366 fn is_initially_syncing(&self) -> bool {
2367 self.syncing
2368 }
2369 }
2370
2371 #[tokio::test]
2372 async fn get_blobs_v3_returns_null_when_syncing() {
2373 let chain_spec: Arc<ChainSpec> =
2374 Arc::new(ChainSpecBuilder::mainnet().osaka_activated().build());
2375 let provider = Arc::new(MockEthProvider::default());
2376 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2377 let (to_engine, _engine_rx) = unbounded_channel::<BeaconEngineMessage<EthEngineTypes>>();
2378
2379 let api = EngineApi::new(
2380 provider,
2381 chain_spec.clone(),
2382 ConsensusEngineHandle::new(to_engine),
2383 payload_store.into(),
2384 NoopTransactionPool::default(),
2385 Runtime::test(),
2386 ClientVersionV1 {
2387 code: ClientCode::RH,
2388 name: "Reth".to_string(),
2389 version: "v0.0.0-test".to_string(),
2390 commit: "test".to_string(),
2391 },
2392 EngineCapabilities::default(),
2393 EthereumEngineValidator::new(chain_spec),
2394 false,
2395 TestNetworkInfo { syncing: true },
2396 );
2397
2398 let res = api.get_blobs_v3_metered(vec![B256::ZERO]);
2399 assert_matches!(res, Ok(None));
2400 }
2401
2402 #[tokio::test]
2403 async fn get_blobs_v4_returns_null_when_syncing() {
2404 let chain_spec: Arc<ChainSpec> =
2405 Arc::new(ChainSpecBuilder::mainnet().osaka_activated().build());
2406 let provider = Arc::new(MockEthProvider::default());
2407 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2408 let (to_engine, _engine_rx) = unbounded_channel::<BeaconEngineMessage<EthEngineTypes>>();
2409
2410 let api = EngineApi::new(
2411 provider,
2412 chain_spec.clone(),
2413 ConsensusEngineHandle::new(to_engine),
2414 payload_store.into(),
2415 NoopTransactionPool::default(),
2416 Runtime::test(),
2417 ClientVersionV1 {
2418 code: ClientCode::RH,
2419 name: "Reth".to_string(),
2420 version: "v0.0.0-test".to_string(),
2421 commit: "test".to_string(),
2422 },
2423 EngineCapabilities::default(),
2424 EthereumEngineValidator::new(chain_spec),
2425 false,
2426 TestNetworkInfo { syncing: true },
2427 );
2428
2429 let res = api.get_blobs_v4_metered(vec![B256::ZERO], B128::from(1u128.to_le_bytes()));
2430 assert_matches!(res, Ok(None));
2431 }
2432
2433 #[test]
2434 fn engine_bitvector_uses_little_endian_cell_indices() {
2435 assert_eq!(
2436 B128::from(u128::from_le_bytes(B128::from(1u128.to_le_bytes()).into())),
2437 B128::from(1u128)
2438 );
2439 assert_eq!(
2440 B128::from(u128::from_le_bytes(B128::from((1u128 << 127).to_le_bytes()).into())),
2441 B128::from(1u128 << 127)
2442 );
2443 assert_eq!(
2444 B128::from(u128::from_le_bytes(B128::from(((1u128 << 64) - 1).to_le_bytes()).into(),)),
2445 B128::from((1u128 << 64) - 1)
2446 );
2447 }
2448
2449 #[tokio::test]
2450 async fn fcu_v4_updates_shared_cell_custody_before_forkchoice_result() {
2451 let chain_spec: Arc<ChainSpec> =
2452 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2453 let provider = Arc::new(MockEthProvider::default());
2454 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2455 let (to_engine, mut engine_rx) = unbounded_channel();
2456 let network = NoopNetwork::default();
2457 let cell_custody = network.cell_custody().clone();
2458
2459 let api = EngineApi::new(
2460 provider,
2461 chain_spec.clone(),
2462 ConsensusEngineHandle::new(to_engine),
2463 payload_store.into(),
2464 NoopTransactionPool::default(),
2465 Runtime::test(),
2466 ClientVersionV1 {
2467 code: ClientCode::RH,
2468 name: "Reth".to_string(),
2469 version: "v0.0.0-test".to_string(),
2470 commit: "test".to_string(),
2471 },
2472 EngineCapabilities::default(),
2473 EthereumEngineValidator::new(chain_spec),
2474 false,
2475 network,
2476 );
2477
2478 let state = ForkchoiceState {
2479 head_block_hash: B256::from([0x33; 32]),
2480 safe_block_hash: B256::ZERO,
2481 finalized_block_hash: B256::ZERO,
2482 };
2483 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2484 let expected_custody_columns = B128::from(0b1010u128);
2485
2486 let api_task = tokio::spawn(async move {
2487 api.fork_choice_updated_v4(state, None, Some(custody_columns)).await
2488 });
2489
2490 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2491 .await
2492 .expect("timed out waiting for forkchoiceUpdated request")
2493 .expect("expected forkchoiceUpdated request");
2494 let response_tx = match request {
2495 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2496 assert!(payload_attrs.is_none());
2497 tx
2498 }
2499 other => panic!("unexpected engine message: {other:?}"),
2500 };
2501 assert_eq!(cell_custody.get(), expected_custody_columns);
2502
2503 response_tx
2504 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2505 PayloadStatusEnum::Valid,
2506 ))))
2507 .expect("send valid response");
2508
2509 api_task
2510 .await
2511 .expect("api task should not panic")
2512 .expect("forkchoiceUpdatedV4 should succeed");
2513 assert_eq!(cell_custody.get(), expected_custody_columns);
2514 }
2515
2516 #[tokio::test]
2517 async fn fcu_v4_updates_shared_cell_custody_when_payload_attrs_invalid() {
2518 let chain_spec: Arc<ChainSpec> =
2519 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2520 let provider = Arc::new(MockEthProvider::default());
2521 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2522 let (to_engine, mut engine_rx) = unbounded_channel();
2523 let network = NoopNetwork::default();
2524 let cell_custody = network.cell_custody().clone();
2525
2526 let api = EngineApi::new(
2527 provider,
2528 chain_spec.clone(),
2529 ConsensusEngineHandle::new(to_engine),
2530 payload_store.into(),
2531 NoopTransactionPool::default(),
2532 Runtime::test(),
2533 ClientVersionV1 {
2534 code: ClientCode::RH,
2535 name: "Reth".to_string(),
2536 version: "v0.0.0-test".to_string(),
2537 commit: "test".to_string(),
2538 },
2539 EngineCapabilities::default(),
2540 EthereumEngineValidator::new(chain_spec),
2541 false,
2542 network,
2543 );
2544
2545 let state = ForkchoiceState {
2546 head_block_hash: B256::from([0x44; 32]),
2547 safe_block_hash: B256::ZERO,
2548 finalized_block_hash: B256::ZERO,
2549 };
2550 let payload_attributes = PayloadAttributes {
2551 timestamp: 1,
2552 prev_randao: B256::ZERO,
2553 suggested_fee_recipient: Address::ZERO,
2554 withdrawals: Some(vec![]),
2555 parent_beacon_block_root: None,
2556 slot_number: None,
2557 ..Default::default()
2558 };
2559 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2560 let expected_custody_columns = B128::from(0b1010u128);
2561
2562 let api_task = tokio::spawn(async move {
2563 api.fork_choice_updated_v4(state, Some(payload_attributes), Some(custody_columns)).await
2564 });
2565
2566 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2567 .await
2568 .expect("timed out waiting for forkchoiceUpdated request")
2569 .expect("expected forkchoiceUpdated request");
2570 let response_tx = match request {
2571 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2572 assert!(
2573 payload_attrs.is_none(),
2574 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2575 );
2576 tx
2577 }
2578 other => panic!("unexpected engine message: {other:?}"),
2579 };
2580 assert_eq!(cell_custody.get(), expected_custody_columns);
2581
2582 response_tx
2583 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2584 PayloadStatusEnum::Valid,
2585 ))))
2586 .expect("send valid response");
2587
2588 let response = api_task.await.expect("api task should not panic");
2589 assert_matches!(
2590 response,
2591 Err(EngineApiError::EngineObjectValidationError(
2592 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2593 ))
2594 );
2595 }
2596
2597 #[tokio::test]
2598 async fn fcu_v3_syncing_precedes_invalid_payload_attributes_validation() {
2599 let (mut handle, api) = setup_engine_api();
2600
2601 let state = ForkchoiceState {
2602 head_block_hash: B256::from([0x11; 32]),
2603 safe_block_hash: B256::ZERO,
2604 finalized_block_hash: B256::ZERO,
2605 };
2606 let payload_attributes = PayloadAttributes {
2607 timestamp: 1,
2608 prev_randao: B256::ZERO,
2609 suggested_fee_recipient: Address::ZERO,
2610 withdrawals: Some(vec![]),
2611 parent_beacon_block_root: None,
2613 slot_number: None,
2614 ..Default::default()
2615 };
2616
2617 let api_task = tokio::spawn(async move {
2618 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2619 });
2620
2621 let request =
2622 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2623 .await
2624 .expect("timed out waiting for forkchoiceUpdated request")
2625 .expect("expected forkchoiceUpdated request");
2626 let response_tx = match request {
2627 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2628 assert!(
2629 payload_attrs.is_none(),
2630 "FCU for syncing state should be evaluated before payload attributes"
2631 );
2632 tx
2633 }
2634 other => panic!("unexpected engine message: {other:?}"),
2635 };
2636
2637 response_tx.send(Ok(OnForkChoiceUpdated::syncing())).expect("send syncing response");
2638
2639 let response = api_task
2640 .await
2641 .expect("api task should not panic")
2642 .expect("forkchoiceUpdatedV3 should return a syncing response");
2643 assert!(response.payload_status.is_syncing());
2644 assert!(response.payload_id.is_none());
2645 }
2646
2647 #[tokio::test]
2648 async fn fcu_v3_valid_forkchoice_missing_beacon_root_returns_invalid_attributes() {
2649 let (mut handle, api) = setup_engine_api();
2650
2651 let state = ForkchoiceState {
2652 head_block_hash: B256::from([0x22; 32]),
2653 safe_block_hash: B256::ZERO,
2654 finalized_block_hash: B256::ZERO,
2655 };
2656 let payload_attributes = PayloadAttributes {
2657 timestamp: 1,
2658 prev_randao: B256::ZERO,
2659 suggested_fee_recipient: Address::ZERO,
2660 withdrawals: Some(vec![]),
2661 parent_beacon_block_root: None,
2662 slot_number: None,
2663 ..Default::default()
2664 };
2665
2666 let api_task = tokio::spawn(async move {
2667 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2668 });
2669
2670 let request =
2671 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2672 .await
2673 .expect("timed out waiting for forkchoiceUpdated request")
2674 .expect("expected forkchoiceUpdated request");
2675
2676 let response_tx = match request {
2677 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2678 assert!(
2679 payload_attrs.is_none(),
2680 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2681 );
2682 tx
2683 }
2684 other => panic!("unexpected engine message: {other:?}"),
2685 };
2686
2687 response_tx
2688 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2689 PayloadStatusEnum::Valid,
2690 ))))
2691 .expect("send valid response");
2692
2693 let response = api_task.await.expect("api task should not panic");
2694 assert_matches!(
2695 response,
2696 Err(EngineApiError::EngineObjectValidationError(
2697 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2698 ))
2699 );
2700
2701 match tokio::time::timeout(std::time::Duration::from_millis(100), handle.from_api.recv())
2702 .await
2703 {
2704 Err(_) | Ok(None) => {}
2705 Ok(Some(BeaconEngineMessage::ForkchoiceUpdated { .. })) => {
2706 panic!("no second forkchoiceUpdated call should be sent when attrs are invalid")
2707 }
2708 Ok(Some(other)) => panic!("unexpected engine message: {other:?}"),
2709 }
2710 }
2711
2712 mod get_payload_bodies {
2714 use super::*;
2715 use alloy_rpc_types_engine::ExecutionPayloadBodyV1;
2716 use reth_testing_utils::generators::{self, random_block_range, BlockRangeParams};
2717
2718 #[tokio::test]
2719 async fn invalid_params() {
2720 let (_, api) = setup_engine_api();
2721
2722 let by_range_tests = [
2723 (0, 0),
2725 (0, 1),
2726 (1, 0),
2727 ];
2728
2729 for (start, count) in by_range_tests {
2731 let res = api.get_payload_bodies_by_range_v1(start, count).await;
2732 assert_matches!(res, Err(EngineApiError::InvalidBodiesRange { .. }));
2733 }
2734 }
2735
2736 #[tokio::test]
2737 async fn request_too_large() {
2738 let (_, api) = setup_engine_api();
2739
2740 let request_count = MAX_PAYLOAD_BODIES_LIMIT + 1;
2741 let res = api.get_payload_bodies_by_range_v1(0, request_count).await;
2742 assert_matches!(res, Err(EngineApiError::PayloadRequestTooLarge { .. }));
2743 }
2744
2745 #[tokio::test]
2746 async fn returns_payload_bodies() {
2747 let mut rng = generators::rng();
2748 let (handle, api) = setup_engine_api();
2749
2750 let (start, count) = (1, 10);
2751 let blocks = random_block_range(
2752 &mut rng,
2753 start..=start + count - 1,
2754 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2755 );
2756 handle
2757 .provider
2758 .extend_blocks(blocks.iter().cloned().map(|b| (b.hash(), b.into_block())));
2759
2760 let expected = blocks
2761 .iter()
2762 .cloned()
2763 .map(|b| Some(ExecutionPayloadBodyV1::from_block(b.into_block())))
2764 .collect::<Vec<_>>();
2765
2766 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2767 assert_eq!(res, expected);
2768 }
2769
2770 #[tokio::test]
2771 async fn returns_payload_bodies_with_gaps() {
2772 let mut rng = generators::rng();
2773 let (handle, api) = setup_engine_api();
2774
2775 let (start, count) = (1, 100);
2776 let blocks = random_block_range(
2777 &mut rng,
2778 start..=start + count - 1,
2779 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2780 );
2781
2782 let first_missing_range = 26..=50;
2784 let second_missing_range = 76..=100;
2785 handle.provider.extend_blocks(
2786 blocks
2787 .iter()
2788 .filter(|b| {
2789 !first_missing_range.contains(&b.number) &&
2790 !second_missing_range.contains(&b.number)
2791 })
2792 .map(|b| (b.hash(), b.clone().into_block())),
2793 );
2794
2795 let expected = blocks
2796 .iter()
2797 .filter(|b| !second_missing_range.contains(&b.number))
2800 .cloned()
2801 .map(|b| {
2802 if first_missing_range.contains(&b.number) {
2803 None
2804 } else {
2805 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2806 }
2807 })
2808 .collect::<Vec<_>>();
2809
2810 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2811 assert_eq!(res, expected);
2812
2813 let expected = blocks
2814 .iter()
2815 .cloned()
2816 .map(|b| {
2819 if first_missing_range.contains(&b.number) ||
2820 second_missing_range.contains(&b.number)
2821 {
2822 None
2823 } else {
2824 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2825 }
2826 })
2827 .collect::<Vec<_>>();
2828
2829 let hashes = blocks.iter().map(|b| b.hash()).collect();
2830 let res = api.get_payload_bodies_by_hash_v1(hashes).await.unwrap();
2831 assert_eq!(res, expected);
2832 }
2833 }
2834}