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 eip7594::BlobCellMask,
9 eip7685::RequestsOrHash,
10};
11use alloy_primitives::{BlockHash, BlockNumber, Bytes, Sealable, B128, B256, U64};
12use alloy_rpc_types_engine::{
13 BogotaPayloadFields, CancunPayloadFields, ClientVersionV1, ExecutionData,
14 ExecutionPayloadBodiesV1, ExecutionPayloadBodiesV2, ExecutionPayloadBodyV1,
15 ExecutionPayloadBodyV2, ExecutionPayloadInputV2, ExecutionPayloadSidecar, ExecutionPayloadV1,
16 ExecutionPayloadV3, ExecutionPayloadV4, ForkchoiceState, ForkchoiceUpdated,
17 ForkchoiceUpdatedResponseV2, PayloadId, PayloadStatus, PayloadStatusV2, PraguePayloadFields,
18 MAX_BYTES_PER_INCLUSION_LIST,
19};
20use async_trait::async_trait;
21use jsonrpsee_core::{server::RpcModule, RpcResult};
22use reth_chainspec::EthereumHardforks;
23use reth_engine_primitives::{ConsensusEngineHandle, EngineApiValidator, EngineTypes};
24use reth_network_api::{CellCustody, NetworkInfo};
25use reth_payload_builder::PayloadStore;
26use reth_payload_primitives::{
27 validate_payload_timestamp, EngineApiMessageVersion, MessageValidationKind,
28 PayloadOrAttributes, PayloadTypes,
29};
30use reth_primitives_traits::{AlloyBlockHeader, Block, BlockBody};
31use reth_rpc_api::{EngineApiServer, IntoEngineApiRpcModule};
32use reth_storage_api::{BalProvider, BlockReader, HeaderProvider, StateProviderFactory};
33use reth_tasks::Runtime;
34use reth_transaction_pool::{BestTransactions, TransactionPool};
35use std::{
36 sync::Arc,
37 time::{Instant, SystemTime},
38};
39use tokio::sync::oneshot;
40use tracing::{debug, trace, warn};
41
42pub type EngineApiSender<Ok> = oneshot::Sender<EngineApiResult<Ok>>;
44
45const MAX_PAYLOAD_BODIES_LIMIT: u64 = 1024;
47
48const MAX_BLOB_LIMIT: usize = 128;
50
51pub struct EngineApi<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
67 inner: Arc<EngineApiInner<Provider, PayloadT, Pool, Validator, ChainSpec>>,
68}
69
70impl<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec>
71 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
72{
73 pub fn chain_spec(&self) -> &Arc<ChainSpec> {
75 &self.inner.chain_spec
76 }
77
78 pub fn client_version(&self) -> &ClientVersionV1 {
80 &self.inner.client
81 }
82}
83
84impl<Provider, PayloadT, Pool, Validator, ChainSpec>
85 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
86where
87 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
88 PayloadT: PayloadTypes,
89 Pool: TransactionPool + 'static,
90 Validator: EngineApiValidator<PayloadT>,
91 ChainSpec: EthereumHardforks + Send + Sync + 'static,
92{
93 #[expect(clippy::too_many_arguments)]
95 pub fn new(
96 provider: Provider,
97 chain_spec: Arc<ChainSpec>,
98 beacon_consensus: ConsensusEngineHandle<PayloadT>,
99 payload_store: PayloadStore<PayloadT>,
100 tx_pool: Pool,
101 task_spawner: Runtime,
102 client: ClientVersionV1,
103 capabilities: EngineCapabilities,
104 validator: Validator,
105 accept_execution_requests_hash: bool,
106 network: impl NetworkInfo + 'static,
107 ) -> Self {
108 let cell_custody = network.cell_custody().clone();
109 let is_syncing = Arc::new(move || network.is_syncing());
110 let inner = Arc::new(EngineApiInner {
111 provider,
112 chain_spec,
113 beacon_consensus,
114 payload_store,
115 task_spawner,
116 metrics: EngineApiMetrics::default(),
117 client,
118 capabilities,
119 tx_pool,
120 validator,
121 accept_execution_requests_hash,
122 cell_custody,
123 is_syncing,
124 });
125 Self { inner }
126 }
127
128 pub fn get_client_version_v1(
130 &self,
131 _client: ClientVersionV1,
132 ) -> EngineApiResult<Vec<ClientVersionV1>> {
133 Ok(vec![self.inner.client.clone()])
134 }
135
136 async fn get_payload_timestamp(&self, payload_id: PayloadId) -> EngineApiResult<u64> {
138 Ok(self
139 .inner
140 .payload_store
141 .payload_timestamp(payload_id)
142 .await
143 .ok_or(EngineApiError::UnknownPayload)??)
144 }
145
146 pub async fn new_payload_v1(
149 &self,
150 payload: PayloadT::ExecutionData,
151 ) -> EngineApiResult<PayloadStatus> {
152 let payload_or_attrs = PayloadOrAttributes::<
153 '_,
154 PayloadT::ExecutionData,
155 PayloadT::PayloadAttributes,
156 >::from_execution_payload(&payload);
157
158 self.inner
159 .validator
160 .validate_version_specific_fields(EngineApiMessageVersion::V1, payload_or_attrs)?;
161
162 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
163 }
164
165 pub async fn new_payload_v1_metered(
167 &self,
168 payload: PayloadT::ExecutionData,
169 ) -> EngineApiResult<PayloadStatus> {
170 let start = Instant::now();
171 let res = Self::new_payload_v1(self, payload).await;
172 let elapsed = start.elapsed();
173 self.inner.metrics.latency.new_payload_v1.record(elapsed);
174 res
175 }
176
177 pub async fn new_payload_v2(
179 &self,
180 payload: PayloadT::ExecutionData,
181 ) -> EngineApiResult<PayloadStatus> {
182 let payload_or_attrs = PayloadOrAttributes::<
183 '_,
184 PayloadT::ExecutionData,
185 PayloadT::PayloadAttributes,
186 >::from_execution_payload(&payload);
187 self.inner
188 .validator
189 .validate_version_specific_fields(EngineApiMessageVersion::V2, payload_or_attrs)?;
190 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
191 }
192
193 pub async fn new_payload_v2_metered(
195 &self,
196 payload: PayloadT::ExecutionData,
197 ) -> EngineApiResult<PayloadStatus> {
198 let start = Instant::now();
199 let res = Self::new_payload_v2(self, payload).await;
200 let elapsed = start.elapsed();
201 self.inner.metrics.latency.new_payload_v2.record(elapsed);
202 res
203 }
204
205 pub async fn new_payload_v3(
207 &self,
208 payload: PayloadT::ExecutionData,
209 ) -> EngineApiResult<PayloadStatus> {
210 let payload_or_attrs = PayloadOrAttributes::<
211 '_,
212 PayloadT::ExecutionData,
213 PayloadT::PayloadAttributes,
214 >::from_execution_payload(&payload);
215 self.inner
216 .validator
217 .validate_version_specific_fields(EngineApiMessageVersion::V3, payload_or_attrs)?;
218
219 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
220 }
221
222 pub async fn new_payload_v3_metered(
224 &self,
225 payload: PayloadT::ExecutionData,
226 ) -> RpcResult<PayloadStatus> {
227 let start = Instant::now();
228
229 let res = Self::new_payload_v3(self, payload).await;
230 let elapsed = start.elapsed();
231 self.inner.metrics.latency.new_payload_v3.record(elapsed);
232 Ok(res?)
233 }
234
235 pub async fn new_payload_v4(
237 &self,
238 payload: PayloadT::ExecutionData,
239 ) -> EngineApiResult<PayloadStatus> {
240 let payload_or_attrs = PayloadOrAttributes::<
241 '_,
242 PayloadT::ExecutionData,
243 PayloadT::PayloadAttributes,
244 >::from_execution_payload(&payload);
245 self.inner
246 .validator
247 .validate_version_specific_fields(EngineApiMessageVersion::V4, payload_or_attrs)?;
248
249 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
250 }
251
252 pub async fn new_payload_v4_metered(
254 &self,
255 payload: PayloadT::ExecutionData,
256 ) -> RpcResult<PayloadStatus> {
257 let start = Instant::now();
258 let res = Self::new_payload_v4(self, payload).await;
259
260 let elapsed = start.elapsed();
261 self.inner.metrics.latency.new_payload_v4.record(elapsed);
262 Ok(res?)
263 }
264
265 pub async fn new_payload_v5(
271 &self,
272 payload: PayloadT::ExecutionData,
273 ) -> EngineApiResult<PayloadStatus> {
274 let payload_or_attrs = PayloadOrAttributes::<
275 '_,
276 PayloadT::ExecutionData,
277 PayloadT::PayloadAttributes,
278 >::from_execution_payload(&payload);
279 self.inner
280 .validator
281 .validate_version_specific_fields(EngineApiMessageVersion::V5, payload_or_attrs)?;
282 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
283 }
284
285 pub async fn new_payload_v5_metered(
287 &self,
288 payload: PayloadT::ExecutionData,
289 ) -> RpcResult<PayloadStatus> {
290 let start = Instant::now();
291 let res = Self::new_payload_v5(self, payload).await;
292 let elapsed = start.elapsed();
293 self.inner.metrics.latency.new_payload_v5.record(elapsed);
294 Ok(res?)
295 }
296
297 pub async fn new_payload_v6(
301 &self,
302 payload: PayloadT::ExecutionData,
303 ) -> EngineApiResult<PayloadStatusV2> {
304 let payload_or_attrs = PayloadOrAttributes::<
305 '_,
306 PayloadT::ExecutionData,
307 PayloadT::PayloadAttributes,
308 >::from_execution_payload(&payload);
309 self.inner
310 .validator
311 .validate_version_specific_fields(EngineApiMessageVersion::V6, payload_or_attrs)?;
312
313 Ok(self.inner.beacon_consensus.new_payload(payload).await?.into())
314 }
315
316 pub async fn new_payload_v6_metered(
318 &self,
319 payload: PayloadT::ExecutionData,
320 ) -> EngineApiResult<PayloadStatusV2> {
321 let start = Instant::now();
322 let result = Self::new_payload_v6(self, payload).await;
323 self.inner.metrics.latency.new_payload_v6.record(start.elapsed());
324 result
325 }
326
327 pub fn accept_execution_requests_hash(&self) -> bool {
329 self.inner.accept_execution_requests_hash
330 }
331}
332
333impl<Provider, EngineT, Pool, Validator, ChainSpec>
334 EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
335where
336 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
337 EngineT: EngineTypes,
338 Pool: TransactionPool + 'static,
339 Validator: EngineApiValidator<EngineT>,
340 ChainSpec: EthereumHardforks + Send + Sync + 'static,
341{
342 pub async fn fork_choice_updated_v1(
349 &self,
350 state: ForkchoiceState,
351 payload_attrs: Option<EngineT::PayloadAttributes>,
352 ) -> EngineApiResult<ForkchoiceUpdated> {
353 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V1, state, payload_attrs)
354 .await
355 }
356
357 pub async fn fork_choice_updated_v1_metered(
359 &self,
360 state: ForkchoiceState,
361 payload_attrs: Option<EngineT::PayloadAttributes>,
362 ) -> EngineApiResult<ForkchoiceUpdated> {
363 let start = Instant::now();
364 let res = Self::fork_choice_updated_v1(self, state, payload_attrs).await;
365 self.inner.metrics.latency.fork_choice_updated_v1.record(start.elapsed());
366 res
367 }
368
369 pub async fn fork_choice_updated_v2(
374 &self,
375 state: ForkchoiceState,
376 payload_attrs: Option<EngineT::PayloadAttributes>,
377 ) -> EngineApiResult<ForkchoiceUpdated> {
378 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V2, state, payload_attrs)
379 .await
380 }
381
382 pub async fn fork_choice_updated_v2_metered(
384 &self,
385 state: ForkchoiceState,
386 payload_attrs: Option<EngineT::PayloadAttributes>,
387 ) -> EngineApiResult<ForkchoiceUpdated> {
388 let start = Instant::now();
389 let res = Self::fork_choice_updated_v2(self, state, payload_attrs).await;
390 self.inner.metrics.latency.fork_choice_updated_v2.record(start.elapsed());
391 res
392 }
393
394 pub async fn fork_choice_updated_v3(
399 &self,
400 state: ForkchoiceState,
401 payload_attrs: Option<EngineT::PayloadAttributes>,
402 ) -> EngineApiResult<ForkchoiceUpdated> {
403 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V3, state, payload_attrs)
404 .await
405 }
406
407 pub async fn fork_choice_updated_v3_metered(
409 &self,
410 state: ForkchoiceState,
411 payload_attrs: Option<EngineT::PayloadAttributes>,
412 ) -> EngineApiResult<ForkchoiceUpdated> {
413 let start = Instant::now();
414 let res = Self::fork_choice_updated_v3(self, state, payload_attrs).await;
415 self.inner.metrics.latency.fork_choice_updated_v3.record(start.elapsed());
416 res
417 }
418
419 pub async fn fork_choice_updated_v4(
424 &self,
425 state: ForkchoiceState,
426 payload_attrs: Option<EngineT::PayloadAttributes>,
427 custody_columns: Option<B128>,
428 ) -> EngineApiResult<ForkchoiceUpdated> {
429 if let Some(custody_columns) = custody_columns {
430 self.inner.cell_custody.set_from_engine_api(custody_columns);
431 }
432 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V4, state, payload_attrs)
433 .await
434 }
435
436 pub async fn fork_choice_updated_v4_metered(
438 &self,
439 state: ForkchoiceState,
440 payload_attrs: Option<EngineT::PayloadAttributes>,
441 custody_columns: Option<B128>,
442 ) -> EngineApiResult<ForkchoiceUpdated> {
443 let start = Instant::now();
444 let res = Self::fork_choice_updated_v4(self, state, payload_attrs, custody_columns).await;
445 self.inner.metrics.latency.fork_choice_updated_v4.record(start.elapsed());
446 res
447 }
448
449 pub async fn fork_choice_updated_v5(
453 &self,
454 state: ForkchoiceState,
455 payload_attrs: Option<EngineT::PayloadAttributes>,
456 custody_columns: Option<B128>,
457 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
458 if let Some(custody_columns) = custody_columns {
459 self.inner.cell_custody.set_from_engine_api(custody_columns);
460 }
461
462 Ok(self
464 .validate_and_execute_forkchoice(EngineApiMessageVersion::V5, state, payload_attrs)
465 .await?
466 .into())
467 }
468
469 pub async fn fork_choice_updated_v5_metered(
471 &self,
472 state: ForkchoiceState,
473 payload_attrs: Option<EngineT::PayloadAttributes>,
474 custody_columns: Option<B128>,
475 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
476 let start = Instant::now();
477 let res = Self::fork_choice_updated_v5(self, state, payload_attrs, custody_columns).await;
478 self.inner.metrics.latency.fork_choice_updated_v5.record(start.elapsed());
479 res
480 }
481
482 pub fn get_inclusion_list_v1(&self) -> EngineApiResult<Vec<Bytes>> {
484 let mut total_size = 0;
485 let mut inclusion_list = Vec::new();
486
487 for pool_tx in self.inner.tx_pool.best_transactions().without_blobs().without_updates() {
488 let encoded = pool_tx.encoded_2718_consensus();
489 let new_size = total_size + encoded.len();
490 if new_size > MAX_BYTES_PER_INCLUSION_LIST as usize {
491 break
492 }
493
494 total_size = new_size;
495 inclusion_list.push(encoded);
496 }
497
498 Ok(inclusion_list)
499 }
500
501 pub fn get_inclusion_list_v1_metered(&self) -> EngineApiResult<Vec<Bytes>> {
503 let start = Instant::now();
504 let result = Self::get_inclusion_list_v1(self);
505 self.inner.metrics.latency.get_inclusion_list_v1.record(start.elapsed());
506 result
507 }
508
509 async fn get_built_payload(
511 &self,
512 payload_id: PayloadId,
513 ) -> EngineApiResult<EngineT::BuiltPayload> {
514 self.inner
515 .payload_store
516 .resolve(payload_id)
517 .await
518 .ok_or(EngineApiError::UnknownPayload)?
519 .map_err(|_| EngineApiError::UnknownPayload)
520 }
521
522 async fn get_payload_inner<R>(
525 &self,
526 payload_id: PayloadId,
527 version: EngineApiMessageVersion,
528 ) -> EngineApiResult<R>
529 where
530 EngineT::BuiltPayload: TryInto<R>,
531 {
532 let timestamp = self.get_payload_timestamp(payload_id).await?;
535 validate_payload_timestamp(
536 &self.inner.chain_spec,
537 version,
538 timestamp,
539 MessageValidationKind::GetPayload,
540 )?;
541
542 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
544 warn!(?version, "could not transform built payload");
545 EngineApiError::UnknownPayload
546 })
547 }
548
549 pub async fn get_payload_v1(
559 &self,
560 payload_id: PayloadId,
561 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
562 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
563 warn!(version = ?EngineApiMessageVersion::V1, "could not transform built payload");
564 EngineApiError::UnknownPayload
565 })
566 }
567
568 pub async fn get_payload_v1_metered(
570 &self,
571 payload_id: PayloadId,
572 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
573 let start = Instant::now();
574 let res = Self::get_payload_v1(self, payload_id).await;
575 self.inner.metrics.latency.get_payload_v1.record(start.elapsed());
576 res
577 }
578
579 pub async fn get_payload_v2(
587 &self,
588 payload_id: PayloadId,
589 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
590 self.get_payload_inner(payload_id, EngineApiMessageVersion::V2).await
591 }
592
593 pub async fn get_payload_v2_metered(
595 &self,
596 payload_id: PayloadId,
597 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
598 let start = Instant::now();
599 let res = Self::get_payload_v2(self, payload_id).await;
600 self.inner.metrics.latency.get_payload_v2.record(start.elapsed());
601 res
602 }
603
604 pub async fn get_payload_v3(
612 &self,
613 payload_id: PayloadId,
614 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
615 self.get_payload_inner(payload_id, EngineApiMessageVersion::V3).await
616 }
617
618 pub async fn get_payload_v3_metered(
620 &self,
621 payload_id: PayloadId,
622 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
623 let start = Instant::now();
624 let res = Self::get_payload_v3(self, payload_id).await;
625 self.inner.metrics.latency.get_payload_v3.record(start.elapsed());
626 res
627 }
628
629 pub async fn get_payload_v4(
637 &self,
638 payload_id: PayloadId,
639 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
640 self.get_payload_inner(payload_id, EngineApiMessageVersion::V4).await
641 }
642
643 pub async fn get_payload_v4_metered(
645 &self,
646 payload_id: PayloadId,
647 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
648 let start = Instant::now();
649 let res = Self::get_payload_v4(self, payload_id).await;
650 self.inner.metrics.latency.get_payload_v4.record(start.elapsed());
651 res
652 }
653
654 pub async fn get_payload_v5(
664 &self,
665 payload_id: PayloadId,
666 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
667 self.get_payload_inner(payload_id, EngineApiMessageVersion::V5).await
668 }
669
670 pub async fn get_payload_v5_metered(
672 &self,
673 payload_id: PayloadId,
674 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
675 let start = Instant::now();
676 let res = Self::get_payload_v5(self, payload_id).await;
677 self.inner.metrics.latency.get_payload_v5.record(start.elapsed());
678 res
679 }
680
681 pub async fn get_payload_v6(
687 &self,
688 payload_id: PayloadId,
689 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
690 self.get_payload_inner(payload_id, EngineApiMessageVersion::V6).await
691 }
692
693 pub async fn get_payload_v6_metered(
695 &self,
696 payload_id: PayloadId,
697 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
698 let start = Instant::now();
699 let res = Self::get_payload_v6(self, payload_id).await;
700 self.inner.metrics.latency.get_payload_v6.record(start.elapsed());
701 res
702 }
703
704 pub async fn get_payload_bodies_by_range_with<F, R>(
707 &self,
708 start: BlockNumber,
709 count: u64,
710 f: F,
711 ) -> EngineApiResult<Vec<Option<R>>>
712 where
713 F: Fn(Provider::Block) -> R + Send + 'static,
714 R: Send + 'static,
715 {
716 let (tx, rx) = oneshot::channel();
717 let inner = self.inner.clone();
718
719 self.inner.task_spawner.spawn_blocking_task(async move {
720 if count > MAX_PAYLOAD_BODIES_LIMIT {
721 tx.send(Err(EngineApiError::PayloadRequestTooLarge { len: count })).ok();
722 return;
723 }
724
725 if start == 0 || count == 0 {
726 tx.send(Err(EngineApiError::InvalidBodiesRange { start, count })).ok();
727 return;
728 }
729
730 let mut result = Vec::with_capacity(count as usize);
731
732 let mut end = start.saturating_add(count - 1);
734
735 if let Ok(best_block) = inner.provider.best_block_number()
738 && end > best_block {
739 end = best_block;
740 }
741
742 let earliest_block = inner.provider.earliest_block_number().unwrap_or(0);
744 for num in start..=end {
745 if tx.is_closed() {
746 return;
747 }
748
749 if num < earliest_block {
750 result.push(None);
751 continue;
752 }
753 let block_result = inner.provider.block(BlockHashOrNumber::Number(num));
754 match block_result {
755 Ok(block) => {
756 result.push(block.map(&f));
757 }
758 Err(err) => {
759 tx.send(Err(EngineApiError::Internal(Box::new(err)))).ok();
760 return;
761 }
762 };
763 }
764 tx.send(Ok(result)).ok();
765 });
766
767 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
768 }
769
770 pub async fn get_payload_bodies_by_range_with_timestamps(
775 &self,
776 start: BlockNumber,
777 count: u64,
778 include_bal: bool,
779 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
780 let bodies = self
781 .get_payload_bodies_by_range_with(start, count, Self::payload_body_with_timestamp)
782 .await?;
783 self.attach_payload_body_bals(bodies, include_bal).await
784 }
785
786 pub async fn get_payload_bodies_by_range_with_timestamps_metered(
788 &self,
789 start: BlockNumber,
790 count: u64,
791 include_bal: bool,
792 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
793 let start_time = Instant::now();
794 let result =
795 self.get_payload_bodies_by_range_with_timestamps(start, count, include_bal).await;
796 let latency = &self.inner.metrics.latency;
797 if include_bal {
798 latency.get_payload_bodies_by_range_v2.record(start_time.elapsed());
799 } else {
800 latency.get_payload_bodies_by_range_v1.record(start_time.elapsed());
801 }
802 result
803 }
804
805 pub async fn get_payload_bodies_by_range_v1(
816 &self,
817 start: BlockNumber,
818 count: u64,
819 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
820 self.get_payload_bodies_by_range_with(start, count, |block| ExecutionPayloadBodyV1 {
821 transactions: block.body().encoded_2718_transactions(),
822 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
823 })
824 .await
825 }
826
827 pub async fn get_payload_bodies_by_range_v1_metered(
829 &self,
830 start: BlockNumber,
831 count: u64,
832 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
833 let start_time = Instant::now();
834 let res = Self::get_payload_bodies_by_range_v1(self, start, count).await;
835 self.inner.metrics.latency.get_payload_bodies_by_range_v1.record(start_time.elapsed());
836 res
837 }
838
839 pub async fn get_payload_bodies_by_range_v2(
843 &self,
844 start: BlockNumber,
845 count: u64,
846 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
847 Ok(self
848 .get_payload_bodies_by_range_with_timestamps(start, count, true)
849 .await?
850 .into_iter()
851 .map(|body| body.map(|(_, body)| body))
852 .collect())
853 }
854
855 pub async fn get_payload_bodies_by_range_v2_metered(
857 &self,
858 start: BlockNumber,
859 count: u64,
860 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
861 let start_time = Instant::now();
862 let res = Self::get_payload_bodies_by_range_v2(self, start, count).await;
863 self.inner.metrics.latency.get_payload_bodies_by_range_v2.record(start_time.elapsed());
864 res
865 }
866
867 pub async fn get_payload_bodies_by_hash_with<F, R>(
869 &self,
870 hashes: Vec<BlockHash>,
871 f: F,
872 ) -> EngineApiResult<Vec<Option<R>>>
873 where
874 F: Fn(Provider::Block) -> R + Send + 'static,
875 R: Send + 'static,
876 {
877 let len = hashes.len() as u64;
878 if len > MAX_PAYLOAD_BODIES_LIMIT {
879 return Err(EngineApiError::PayloadRequestTooLarge { len });
880 }
881
882 let (tx, rx) = oneshot::channel();
883 let inner = self.inner.clone();
884
885 self.inner.task_spawner.spawn_blocking_task(async move {
886 let mut result = Vec::with_capacity(hashes.len());
887 for hash in hashes {
888 if tx.is_closed() {
889 return;
890 }
891
892 let block_result = inner.provider.block(BlockHashOrNumber::Hash(hash));
893 match block_result {
894 Ok(block) => {
895 result.push(block.map(&f));
896 }
897 Err(err) => {
898 let _ = tx.send(Err(EngineApiError::Internal(Box::new(err))));
899 return;
900 }
901 }
902 }
903 tx.send(Ok(result)).ok();
904 });
905
906 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
907 }
908
909 pub async fn get_payload_bodies_by_hash_with_timestamps(
911 &self,
912 hashes: Vec<BlockHash>,
913 include_bal: bool,
914 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
915 let bodies =
916 self.get_payload_bodies_by_hash_with(hashes, Self::payload_body_with_timestamp).await?;
917 self.attach_payload_body_bals(bodies, include_bal).await
918 }
919
920 pub async fn get_payload_bodies_by_hash_with_timestamps_metered(
922 &self,
923 hashes: Vec<BlockHash>,
924 include_bal: bool,
925 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
926 let start = Instant::now();
927 let result = self.get_payload_bodies_by_hash_with_timestamps(hashes, include_bal).await;
928 let latency = &self.inner.metrics.latency;
929 if include_bal {
930 latency.get_payload_bodies_by_hash_v2.record(start.elapsed());
931 } else {
932 latency.get_payload_bodies_by_hash_v1.record(start.elapsed());
933 }
934 result
935 }
936
937 fn payload_body_with_timestamp(
938 block: Provider::Block,
939 ) -> (BlockHash, u64, ExecutionPayloadBodyV2) {
940 (
941 block.header().hash_slow(),
942 block.header().timestamp(),
943 ExecutionPayloadBodyV2 {
944 transactions: block.body().encoded_2718_transactions(),
945 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
946 block_access_list: None,
947 },
948 )
949 }
950
951 async fn attach_payload_body_bals(
952 &self,
953 mut bodies: Vec<Option<(BlockHash, u64, ExecutionPayloadBodyV2)>>,
954 include_bal: bool,
955 ) -> EngineApiResult<Vec<Option<(u64, ExecutionPayloadBodyV2)>>> {
956 if include_bal {
957 let hashes = bodies.iter().flatten().map(|(hash, _, _)| *hash).collect();
958 let bals = self.get_block_access_lists_by_hashes(hashes).await?;
959 for ((_, _, body), bal) in bodies.iter_mut().flatten().zip(bals) {
960 body.block_access_list = bal;
961 }
962 }
963 Ok(bodies
964 .into_iter()
965 .map(|body| body.map(|(_, timestamp, body)| (timestamp, body)))
966 .collect())
967 }
968
969 async fn get_block_access_lists_by_hashes(
970 &self,
971 hashes: Vec<BlockHash>,
972 ) -> EngineApiResult<Vec<Option<Bytes>>> {
973 let len = hashes.len() as u64;
974 if len > MAX_PAYLOAD_BODIES_LIMIT {
975 return Err(EngineApiError::PayloadRequestTooLarge { len });
976 }
977
978 let (tx, rx) = oneshot::channel();
979 let inner = self.inner.clone();
980
981 self.inner.task_spawner.spawn_blocking_task(async move {
982 if tx.is_closed() {
983 return;
984 }
985
986 tx.send(
987 inner
988 .provider
989 .get_bals_by_hashes(&hashes)
990 .map_err(|err| EngineApiError::Internal(Box::new(err))),
991 )
992 .ok();
993 });
994
995 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
996 }
997
998 pub async fn get_payload_bodies_by_hash_v1(
1000 &self,
1001 hashes: Vec<BlockHash>,
1002 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
1003 self.get_payload_bodies_by_hash_with(hashes, |block| ExecutionPayloadBodyV1 {
1004 transactions: block.body().encoded_2718_transactions(),
1005 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
1006 })
1007 .await
1008 }
1009
1010 pub async fn get_payload_bodies_by_hash_v1_metered(
1012 &self,
1013 hashes: Vec<BlockHash>,
1014 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
1015 let start = Instant::now();
1016 let res = Self::get_payload_bodies_by_hash_v1(self, hashes).await;
1017 self.inner.metrics.latency.get_payload_bodies_by_hash_v1.record(start.elapsed());
1018 res
1019 }
1020
1021 pub async fn get_payload_bodies_by_hash_v2(
1025 &self,
1026 hashes: Vec<BlockHash>,
1027 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
1028 let payload_bodies =
1029 self.get_payload_bodies_by_hash_with(hashes.clone(), |block| ExecutionPayloadBodyV2 {
1030 transactions: block.body().encoded_2718_transactions(),
1031 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
1032 block_access_list: None,
1033 });
1034 let block_access_lists = self.get_block_access_lists_by_hashes(hashes);
1035 let (mut payload_bodies, block_access_lists) =
1036 tokio::try_join!(payload_bodies, block_access_lists)?;
1037
1038 for (payload_body, block_access_list) in payload_bodies.iter_mut().zip(block_access_lists) {
1039 if let Some(payload_body) = payload_body {
1040 payload_body.block_access_list = block_access_list;
1041 }
1042 }
1043
1044 Ok(payload_bodies)
1045 }
1046
1047 pub async fn get_payload_bodies_by_hash_v2_metered(
1049 &self,
1050 hashes: Vec<BlockHash>,
1051 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
1052 let start = Instant::now();
1053 let res = Self::get_payload_bodies_by_hash_v2(self, hashes).await;
1054 self.inner.metrics.latency.get_payload_bodies_by_hash_v2.record(start.elapsed());
1055 res
1056 }
1057
1058 async fn validate_and_execute_forkchoice(
1072 &self,
1073 version: EngineApiMessageVersion,
1074 state: ForkchoiceState,
1075 payload_attrs: Option<EngineT::PayloadAttributes>,
1076 ) -> EngineApiResult<ForkchoiceUpdated> {
1077 if let Some(ref attrs) = payload_attrs {
1078 let attr_validation_res =
1079 self.inner.validator.ensure_well_formed_attributes(version, attrs);
1080
1081 if let Err(err) = attr_validation_res {
1091 let fcu_res = self.inner.beacon_consensus.fork_choice_updated(state, None).await?;
1092 if fcu_res.is_invalid() || fcu_res.payload_status.is_syncing() {
1093 return Ok(fcu_res)
1094 }
1095 return Err(err.into())
1096 }
1097 }
1098
1099 Ok(self.inner.beacon_consensus.fork_choice_updated(state, payload_attrs).await?)
1100 }
1101
1102 pub fn capabilities(&self) -> &EngineCapabilities {
1104 &self.inner.capabilities
1105 }
1106
1107 fn has_blobs(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1108 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1109 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1110 }
1111
1112 self.inner
1113 .tx_pool
1114 .has_blobs_for_versioned_hashes(&versioned_hashes)
1115 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1116 }
1117
1118 pub fn has_blobs_metered(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1120 let start = Instant::now();
1121 let res = Self::has_blobs(self, versioned_hashes);
1122 self.inner.metrics.latency.has_blobs.record(start.elapsed());
1123 res
1124 }
1125
1126 fn get_blobs_v1(
1127 &self,
1128 versioned_hashes: Vec<B256>,
1129 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1130 let current_timestamp =
1132 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1133 if self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1134 return Err(EngineApiError::EngineObjectValidationError(
1135 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1136 ));
1137 }
1138
1139 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1140 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1141 }
1142
1143 self.inner
1144 .tx_pool
1145 .get_blobs_for_versioned_hashes_v1(&versioned_hashes)
1146 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1147 }
1148
1149 pub fn get_blobs_v1_metered(
1151 &self,
1152 versioned_hashes: Vec<B256>,
1153 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1154 let hashes_len = versioned_hashes.len();
1155 let start = Instant::now();
1156 let res = Self::get_blobs_v1(self, versioned_hashes);
1157 self.inner.metrics.latency.get_blobs_v1.record(start.elapsed());
1158
1159 if let Ok(blobs) = &res {
1160 let blobs_found = blobs.iter().flatten().count();
1161 let blobs_missed = hashes_len - blobs_found;
1162
1163 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1164 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1165 }
1166
1167 res
1168 }
1169
1170 fn get_blobs_v2(
1171 &self,
1172 versioned_hashes: Vec<B256>,
1173 ) -> EngineApiResult<Option<Vec<BlobAndProofV2>>> {
1174 let current_timestamp =
1176 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1177 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1178 return Err(EngineApiError::EngineObjectValidationError(
1179 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1180 ));
1181 }
1182
1183 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1184 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1185 }
1186
1187 self.inner
1188 .tx_pool
1189 .get_blobs_for_versioned_hashes_v2(&versioned_hashes)
1190 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1191 }
1192
1193 fn get_blobs_v3(
1194 &self,
1195 versioned_hashes: Vec<B256>,
1196 ) -> EngineApiResult<Option<Vec<Option<BlobAndProofV2>>>> {
1197 let current_timestamp =
1199 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1200 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1201 return Err(EngineApiError::EngineObjectValidationError(
1202 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1203 ));
1204 }
1205
1206 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1207 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1208 }
1209
1210 if (*self.inner.is_syncing)() {
1212 return Ok(None)
1213 }
1214
1215 self.inner
1216 .tx_pool
1217 .get_blobs_for_versioned_hashes_v3(&versioned_hashes)
1218 .map(Some)
1219 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1220 }
1221
1222 fn get_blobs_v4(
1223 &self,
1224 versioned_hashes: Vec<B256>,
1225 indices_bitarray: B128,
1226 ) -> EngineApiResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1227 let cell_mask = BlobCellMask::from_bits(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, cell_mask)
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 for index in [0, 7, 8, 63, 64, 127] {
2436 let wire_mask = B128::from((1u128 << index).to_le_bytes());
2437 let mask = BlobCellMask::from_bits(u128::from_le_bytes(wire_mask.into()));
2438 assert_eq!(mask.selected_indices().collect::<Vec<_>>(), vec![index]);
2439 }
2440 }
2441
2442 #[tokio::test]
2443 async fn fcu_v4_updates_shared_cell_custody_before_forkchoice_result() {
2444 let chain_spec: Arc<ChainSpec> =
2445 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2446 let provider = Arc::new(MockEthProvider::default());
2447 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2448 let (to_engine, mut engine_rx) = unbounded_channel();
2449 let network = NoopNetwork::default();
2450 let cell_custody = network.cell_custody().clone();
2451
2452 let api = EngineApi::new(
2453 provider,
2454 chain_spec.clone(),
2455 ConsensusEngineHandle::new(to_engine),
2456 payload_store.into(),
2457 NoopTransactionPool::default(),
2458 Runtime::test(),
2459 ClientVersionV1 {
2460 code: ClientCode::RH,
2461 name: "Reth".to_string(),
2462 version: "v0.0.0-test".to_string(),
2463 commit: "test".to_string(),
2464 },
2465 EngineCapabilities::default(),
2466 EthereumEngineValidator::new(chain_spec),
2467 false,
2468 network,
2469 );
2470
2471 let state = ForkchoiceState {
2472 head_block_hash: B256::from([0x33; 32]),
2473 safe_block_hash: B256::ZERO,
2474 finalized_block_hash: B256::ZERO,
2475 };
2476 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2477 let expected_custody_columns = B128::from(0b1010u128);
2478
2479 let api_task = tokio::spawn(async move {
2480 api.fork_choice_updated_v4(state, None, Some(custody_columns)).await
2481 });
2482
2483 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2484 .await
2485 .expect("timed out waiting for forkchoiceUpdated request")
2486 .expect("expected forkchoiceUpdated request");
2487 let response_tx = match request {
2488 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2489 assert!(payload_attrs.is_none());
2490 tx
2491 }
2492 other => panic!("unexpected engine message: {other:?}"),
2493 };
2494 assert_eq!(cell_custody.get(), expected_custody_columns);
2495
2496 response_tx
2497 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2498 PayloadStatusEnum::Valid,
2499 ))))
2500 .expect("send valid response");
2501
2502 api_task
2503 .await
2504 .expect("api task should not panic")
2505 .expect("forkchoiceUpdatedV4 should succeed");
2506 assert_eq!(cell_custody.get(), expected_custody_columns);
2507 }
2508
2509 #[tokio::test]
2510 async fn fcu_v4_updates_shared_cell_custody_when_payload_attrs_invalid() {
2511 let chain_spec: Arc<ChainSpec> =
2512 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2513 let provider = Arc::new(MockEthProvider::default());
2514 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2515 let (to_engine, mut engine_rx) = unbounded_channel();
2516 let network = NoopNetwork::default();
2517 let cell_custody = network.cell_custody().clone();
2518
2519 let api = EngineApi::new(
2520 provider,
2521 chain_spec.clone(),
2522 ConsensusEngineHandle::new(to_engine),
2523 payload_store.into(),
2524 NoopTransactionPool::default(),
2525 Runtime::test(),
2526 ClientVersionV1 {
2527 code: ClientCode::RH,
2528 name: "Reth".to_string(),
2529 version: "v0.0.0-test".to_string(),
2530 commit: "test".to_string(),
2531 },
2532 EngineCapabilities::default(),
2533 EthereumEngineValidator::new(chain_spec),
2534 false,
2535 network,
2536 );
2537
2538 let state = ForkchoiceState {
2539 head_block_hash: B256::from([0x44; 32]),
2540 safe_block_hash: B256::ZERO,
2541 finalized_block_hash: B256::ZERO,
2542 };
2543 let payload_attributes = PayloadAttributes {
2544 timestamp: 1,
2545 prev_randao: B256::ZERO,
2546 suggested_fee_recipient: Address::ZERO,
2547 withdrawals: Some(vec![]),
2548 parent_beacon_block_root: None,
2549 slot_number: None,
2550 ..Default::default()
2551 };
2552 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2553 let expected_custody_columns = B128::from(0b1010u128);
2554
2555 let api_task = tokio::spawn(async move {
2556 api.fork_choice_updated_v4(state, Some(payload_attributes), Some(custody_columns)).await
2557 });
2558
2559 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2560 .await
2561 .expect("timed out waiting for forkchoiceUpdated request")
2562 .expect("expected forkchoiceUpdated request");
2563 let response_tx = match request {
2564 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2565 assert!(
2566 payload_attrs.is_none(),
2567 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2568 );
2569 tx
2570 }
2571 other => panic!("unexpected engine message: {other:?}"),
2572 };
2573 assert_eq!(cell_custody.get(), expected_custody_columns);
2574
2575 response_tx
2576 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2577 PayloadStatusEnum::Valid,
2578 ))))
2579 .expect("send valid response");
2580
2581 let response = api_task.await.expect("api task should not panic");
2582 assert_matches!(
2583 response,
2584 Err(EngineApiError::EngineObjectValidationError(
2585 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2586 ))
2587 );
2588 }
2589
2590 #[tokio::test]
2591 async fn fcu_v3_syncing_precedes_invalid_payload_attributes_validation() {
2592 let (mut handle, api) = setup_engine_api();
2593
2594 let state = ForkchoiceState {
2595 head_block_hash: B256::from([0x11; 32]),
2596 safe_block_hash: B256::ZERO,
2597 finalized_block_hash: B256::ZERO,
2598 };
2599 let payload_attributes = PayloadAttributes {
2600 timestamp: 1,
2601 prev_randao: B256::ZERO,
2602 suggested_fee_recipient: Address::ZERO,
2603 withdrawals: Some(vec![]),
2604 parent_beacon_block_root: None,
2606 slot_number: None,
2607 ..Default::default()
2608 };
2609
2610 let api_task = tokio::spawn(async move {
2611 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2612 });
2613
2614 let request =
2615 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2616 .await
2617 .expect("timed out waiting for forkchoiceUpdated request")
2618 .expect("expected forkchoiceUpdated request");
2619 let response_tx = match request {
2620 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2621 assert!(
2622 payload_attrs.is_none(),
2623 "FCU for syncing state should be evaluated before payload attributes"
2624 );
2625 tx
2626 }
2627 other => panic!("unexpected engine message: {other:?}"),
2628 };
2629
2630 response_tx.send(Ok(OnForkChoiceUpdated::syncing())).expect("send syncing response");
2631
2632 let response = api_task
2633 .await
2634 .expect("api task should not panic")
2635 .expect("forkchoiceUpdatedV3 should return a syncing response");
2636 assert!(response.payload_status.is_syncing());
2637 assert!(response.payload_id.is_none());
2638 }
2639
2640 #[tokio::test]
2641 async fn fcu_v3_valid_forkchoice_missing_beacon_root_returns_invalid_attributes() {
2642 let (mut handle, api) = setup_engine_api();
2643
2644 let state = ForkchoiceState {
2645 head_block_hash: B256::from([0x22; 32]),
2646 safe_block_hash: B256::ZERO,
2647 finalized_block_hash: B256::ZERO,
2648 };
2649 let payload_attributes = PayloadAttributes {
2650 timestamp: 1,
2651 prev_randao: B256::ZERO,
2652 suggested_fee_recipient: Address::ZERO,
2653 withdrawals: Some(vec![]),
2654 parent_beacon_block_root: None,
2655 slot_number: None,
2656 ..Default::default()
2657 };
2658
2659 let api_task = tokio::spawn(async move {
2660 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2661 });
2662
2663 let request =
2664 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2665 .await
2666 .expect("timed out waiting for forkchoiceUpdated request")
2667 .expect("expected forkchoiceUpdated request");
2668
2669 let response_tx = match request {
2670 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2671 assert!(
2672 payload_attrs.is_none(),
2673 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2674 );
2675 tx
2676 }
2677 other => panic!("unexpected engine message: {other:?}"),
2678 };
2679
2680 response_tx
2681 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2682 PayloadStatusEnum::Valid,
2683 ))))
2684 .expect("send valid response");
2685
2686 let response = api_task.await.expect("api task should not panic");
2687 assert_matches!(
2688 response,
2689 Err(EngineApiError::EngineObjectValidationError(
2690 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2691 ))
2692 );
2693
2694 match tokio::time::timeout(std::time::Duration::from_millis(100), handle.from_api.recv())
2695 .await
2696 {
2697 Err(_) | Ok(None) => {}
2698 Ok(Some(BeaconEngineMessage::ForkchoiceUpdated { .. })) => {
2699 panic!("no second forkchoiceUpdated call should be sent when attrs are invalid")
2700 }
2701 Ok(Some(other)) => panic!("unexpected engine message: {other:?}"),
2702 }
2703 }
2704
2705 mod get_payload_bodies {
2707 use super::*;
2708 use alloy_rpc_types_engine::ExecutionPayloadBodyV1;
2709 use reth_testing_utils::generators::{self, random_block_range, BlockRangeParams};
2710
2711 #[tokio::test]
2712 async fn invalid_params() {
2713 let (_, api) = setup_engine_api();
2714
2715 let by_range_tests = [
2716 (0, 0),
2718 (0, 1),
2719 (1, 0),
2720 ];
2721
2722 for (start, count) in by_range_tests {
2724 let res = api.get_payload_bodies_by_range_v1(start, count).await;
2725 assert_matches!(res, Err(EngineApiError::InvalidBodiesRange { .. }));
2726 }
2727 }
2728
2729 #[tokio::test]
2730 async fn request_too_large() {
2731 let (_, api) = setup_engine_api();
2732
2733 let request_count = MAX_PAYLOAD_BODIES_LIMIT + 1;
2734 let res = api.get_payload_bodies_by_range_v1(0, request_count).await;
2735 assert_matches!(res, Err(EngineApiError::PayloadRequestTooLarge { .. }));
2736 }
2737
2738 #[tokio::test]
2739 async fn returns_payload_bodies() {
2740 let mut rng = generators::rng();
2741 let (handle, api) = setup_engine_api();
2742
2743 let (start, count) = (1, 10);
2744 let blocks = random_block_range(
2745 &mut rng,
2746 start..=start + count - 1,
2747 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2748 );
2749 handle
2750 .provider
2751 .extend_blocks(blocks.iter().cloned().map(|b| (b.hash(), b.into_block())));
2752
2753 let expected = blocks
2754 .iter()
2755 .cloned()
2756 .map(|b| Some(ExecutionPayloadBodyV1::from_block(b.into_block())))
2757 .collect::<Vec<_>>();
2758
2759 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2760 assert_eq!(res, expected);
2761 }
2762
2763 #[tokio::test]
2764 async fn returns_payload_bodies_with_gaps() {
2765 let mut rng = generators::rng();
2766 let (handle, api) = setup_engine_api();
2767
2768 let (start, count) = (1, 100);
2769 let blocks = random_block_range(
2770 &mut rng,
2771 start..=start + count - 1,
2772 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2773 );
2774
2775 let first_missing_range = 26..=50;
2777 let second_missing_range = 76..=100;
2778 handle.provider.extend_blocks(
2779 blocks
2780 .iter()
2781 .filter(|b| {
2782 !first_missing_range.contains(&b.number) &&
2783 !second_missing_range.contains(&b.number)
2784 })
2785 .map(|b| (b.hash(), b.clone().into_block())),
2786 );
2787
2788 let expected = blocks
2789 .iter()
2790 .filter(|b| !second_missing_range.contains(&b.number))
2793 .cloned()
2794 .map(|b| {
2795 if first_missing_range.contains(&b.number) {
2796 None
2797 } else {
2798 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2799 }
2800 })
2801 .collect::<Vec<_>>();
2802
2803 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2804 assert_eq!(res, expected);
2805
2806 let expected = blocks
2807 .iter()
2808 .cloned()
2809 .map(|b| {
2812 if first_missing_range.contains(&b.number) ||
2813 second_missing_range.contains(&b.number)
2814 {
2815 None
2816 } else {
2817 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2818 }
2819 })
2820 .collect::<Vec<_>>();
2821
2822 let hashes = blocks.iter().map(|b| b.hash()).collect();
2823 let res = api.get_payload_bodies_by_hash_v1(hashes).await.unwrap();
2824 assert_eq!(res, expected);
2825 }
2826 }
2827}