Skip to main content

reth_rpc_engine_api/
engine_api.rs

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
42/// The Engine API response sender.
43pub type EngineApiSender<Ok> = oneshot::Sender<EngineApiResult<Ok>>;
44
45/// The upper limit for payload bodies request.
46const MAX_PAYLOAD_BODIES_LIMIT: u64 = 1024;
47
48/// The upper limit for blobs in `engine_getBlobsVx`.
49const MAX_BLOB_LIMIT: usize = 128;
50
51/// The Engine API implementation that grants the Consensus layer access to data and
52/// functions in the Execution layer that are crucial for the consensus process.
53///
54/// This type is generic over [`EngineTypes`] and intended to be used as the entrypoint for engine
55/// API processing. It can be reused by other non L1 engine APIs that deviate from the L1 spec but
56/// are still follow the engine API model.
57///
58/// ## Implementers
59///
60/// Implementing support for an engine API jsonrpsee RPC handler is done by defining the engine API
61/// server trait and implementing it on a type that can either wrap this [`EngineApi`] type or
62/// use a custom [`EngineTypes`] implementation if it mirrors ethereum's versioned engine API
63/// endpoints (e.g. opstack).
64/// See also [`EngineApiServer`] implementation for this type which is the
65/// L1 implementation.
66pub 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    /// Returns the configured chainspec.
74    pub fn chain_spec(&self) -> &Arc<ChainSpec> {
75        &self.inner.chain_spec
76    }
77
78    /// Returns the configured client version.
79    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    /// Create new instance of [`EngineApi`].
94    #[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    /// Fetches the client version.
129    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    /// Fetches the timestamp of the payload with the given id.
137    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    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_newpayloadv1>
147    /// Caution: This should not accept the `withdrawals` field
148    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    /// Metered version of `new_payload_v1`.
166    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    /// See also <https://github.com/ethereum/execution-apis/blob/584905270d8ad665718058060267061ecfd79ca5/src/engine/shanghai.md#engine_newpayloadv2>
178    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    /// Metered version of `new_payload_v2`.
194    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    /// See also <https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/cancun.md#engine_newpayloadv3>
206    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    /// Metrics version of `new_payload_v3`
223    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    /// See also <https://github.com/ethereum/execution-apis/blob/7907424db935b93c2fe6a3c0faab943adebe8557/src/engine/prague.md#engine_newpayloadv4>
236    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    /// Metrics version of `new_payload_v4`
253    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    /// Handler for `engine_newPayloadV5`
266    ///
267    /// Post-Amsterdam payload handler.
268    ///
269    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_newpayloadv5>
270    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    /// Metrics version of `new_payload_v5`
286    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    /// Handler for `engine_newPayloadV6`.
298    ///
299    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/bogota.md#engine_newpayloadv6>
300    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    /// Metrics version of `new_payload_v6`.
317    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    /// Returns whether the engine accepts execution requests hash.
328    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    /// Sends a message to the beacon consensus engine to update the fork choice _without_
343    /// withdrawals.
344    ///
345    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_forkchoiceUpdatedV1>
346    ///
347    /// Caution: This should not accept the `withdrawals` field
348    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    /// Metrics version of `fork_choice_updated_v1`
358    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    /// Sends a message to the beacon consensus engine to update the fork choice _with_ withdrawals,
370    /// but only _after_ shanghai.
371    ///
372    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_forkchoiceupdatedv2>
373    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    /// Metrics version of `fork_choice_updated_v2`
383    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    /// Sends a message to the beacon consensus engine to update the fork choice _with_ withdrawals,
395    /// but only _after_ cancun.
396    ///
397    /// See also  <https://github.com/ethereum/execution-apis/blob/main/src/engine/cancun.md#engine_forkchoiceupdatedv3>
398    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    /// Metrics version of `fork_choice_updated_v3`
408    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    /// Sends a message to the beacon consensus engine to update the fork choice _with_ slot number,
420    /// but only _after_ amsterdam.
421    ///
422    /// See also  <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_forkchoiceupdatedv4>
423    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    /// Metrics version of `fork_choice_updated_v4`
437    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    /// Handler for `engine_forkchoiceUpdatedV5`.
450    ///
451    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/bogota.md#engine_forkchoiceupdatedv5>
452    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        // Todo: Validate IL and populate `inclusion_list_satisfied` properly and test.
463        Ok(self
464            .validate_and_execute_forkchoice(EngineApiMessageVersion::V5, state, payload_attrs)
465            .await?
466            .into())
467    }
468
469    /// Metrics version of `fork_choice_updated_v5`
470    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    /// Builds an EIP-7805 inclusion list from the local transaction pool.
483    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    /// Metrics version of `get_inclusion_list_v1`.
502    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    /// Helper function for retrieving the build payload by id.
510    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    /// Helper function for validating the payload timestamp and retrieving & converting the payload
523    /// into desired envelope.
524    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        // Validate timestamp according to engine rules
533        // Enforces Osaka restrictions on `getPayloadV4`.
534        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        // Now resolve the payload
543        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    /// Returns the most recent version of the payload that is available in the corresponding
550    /// payload build process at the time of receiving this call.
551    ///
552    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_getPayloadV1>
553    ///
554    /// Caution: This should not return the `withdrawals` field
555    ///
556    /// Note:
557    /// > Provider software MAY stop the corresponding build process after serving this call.
558    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    /// Metrics version of `get_payload_v1`
569    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    /// Returns the most recent version of the payload that is available in the corresponding
580    /// payload build process at the time of receiving this call.
581    ///
582    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_getpayloadv2>
583    ///
584    /// Note:
585    /// > Provider software MAY stop the corresponding build process after serving this call.
586    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    /// Metrics version of `get_payload_v2`
594    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    /// Returns the most recent version of the payload that is available in the corresponding
605    /// payload build process at the time of receiving this call.
606    ///
607    /// See also <https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/cancun.md#engine_getpayloadv3>
608    ///
609    /// Note:
610    /// > Provider software MAY stop the corresponding build process after serving this call.
611    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    /// Metrics version of `get_payload_v3`
619    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    /// Returns the most recent version of the payload that is available in the corresponding
630    /// payload build process at the time of receiving this call.
631    ///
632    /// See also <https://github.com/ethereum/execution-apis/blob/7907424db935b93c2fe6a3c0faab943adebe8557/src/engine/prague.md#engine_getpayloadv4>
633    ///
634    /// Note:
635    /// > Provider software MAY stop the corresponding build process after serving this call.
636    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    /// Metrics version of `get_payload_v4`
644    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    /// Handler for `engine_getPayloadV5`
655    ///
656    /// Returns the most recent version of the payload that is available in the corresponding
657    /// payload build process at the time of receiving this call.
658    ///
659    /// See also <https://github.com/ethereum/execution-apis/blob/15399c2e2f16a5f800bf3f285640357e2c245ad9/src/engine/osaka.md#engine_getpayloadv5>
660    ///
661    /// Note:
662    /// > Provider software MAY stop the corresponding build process after serving this call.
663    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    /// Metrics version of `get_payload_v5`
671    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    /// Handler for `engine_getPayloadV6`
682    ///
683    /// Post-Amsterdam payload handler that includes Block Access Lists (BAL).
684    ///
685    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_getpayloadv6>
686    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    /// Metrics version of `get_payload_v6`
694    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    /// Fetches all the blocks for the provided range starting at `start`, containing `count`
705    /// blocks and returns the mapped payload bodies.
706    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            // -1 so range is inclusive
733            let mut end = start.saturating_add(count - 1);
734
735            // > Client software MUST NOT return trailing null values if the request extends past the current latest known block.
736            // truncate the end if it's greater than the last block
737            if let Ok(best_block) = inner.provider.best_block_number()
738                && end > best_block {
739                    end = best_block;
740                }
741
742            // Check if the requested range starts before the earliest available block due to pruning/expiry
743            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    /// Returns payload bodies and their timestamps from the same block read.
771    ///
772    /// SSZ transports use the timestamp to select the block's fork schema. Reading it
773    /// separately by block number could pair a body with another block after a reorg.
774    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    /// Records SSZ body requests in the corresponding JSON-RPC V1/V2 latency histogram.
787    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    /// Returns the execution payload bodies by the range starting at `start`, containing `count`
806    /// blocks.
807    ///
808    /// WARNING: This method is associated with the `BeaconBlocksByRange` message in the consensus
809    /// layer p2p specification, meaning the input should be treated as untrusted or potentially
810    /// adversarial.
811    ///
812    /// Implementers should take care when acting on the input to this method, specifically
813    /// ensuring that the range is limited properly, and that the range boundaries are computed
814    /// correctly and without panics.
815    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    /// Metrics version of `get_payload_bodies_by_range_v1`
828    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    /// Returns the execution payload bodies by the range (V2).
840    ///
841    /// V2 includes the `block_access_list` field for EIP-7928 BAL support.
842    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    /// Metrics version of `get_payload_bodies_by_range_v2`
856    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    /// Called to retrieve execution payload bodies by hashes.
868    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    /// Returns payload bodies and their timestamps from the same block read.
910    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    /// Records SSZ body requests in the corresponding JSON-RPC V1/V2 latency histogram.
921    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    /// Called to retrieve execution payload bodies by hashes.
999    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    /// Metrics version of `get_payload_bodies_by_hash_v1`
1011    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    /// Called to retrieve execution payload bodies by hashes (V2).
1022    ///
1023    /// V2 includes the `block_access_list` field for EIP-7928 BAL support.
1024    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    /// Metrics version of `get_payload_bodies_by_hash_v2`
1048    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    /// Validates the `engine_forkchoiceUpdated` payload attributes and executes the forkchoice
1059    /// update.
1060    ///
1061    /// The payload attributes will be validated according to the engine API rules for the given
1062    /// message version:
1063    /// * If the version is [`EngineApiMessageVersion::V1`], then the payload attributes will be
1064    ///   validated according to the Paris rules.
1065    /// * If the version is [`EngineApiMessageVersion::V2`], then the payload attributes will be
1066    ///   validated according to the Shanghai rules, as well as the validity changes from cancun:
1067    ///   <https://github.com/ethereum/execution-apis/blob/584905270d8ad665718058060267061ecfd79ca5/src/engine/cancun.md#update-the-methods-of-previous-forks>
1068    ///
1069    /// * If the version above [`EngineApiMessageVersion::V3`], then the payload attributes will be
1070    ///   validated according to the Cancun rules.
1071    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            // From the engine API spec:
1082            //
1083            // Client software MUST ensure that payloadAttributes.timestamp is greater than
1084            // timestamp of a block referenced by forkchoiceState.headBlockHash. If this condition
1085            // isn't held client software MUST respond with -38003: Invalid payload attributes and
1086            // MUST NOT begin a payload build process. In such an event, the forkchoiceState
1087            // update MUST NOT be rolled back.
1088            //
1089            // NOTE: This also applies to cancun/shanghai-specific payload attributes.
1090            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    /// Returns reference to supported capabilities.
1103    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    /// Metered version of `has_blobs`.
1119    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        // Only allow this method before Osaka fork
1131        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    /// Metered version of `get_blobs_v1`.
1150    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        // Check if Osaka fork is active
1175        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        // Check if Osaka fork is active
1198        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        // Spec requires returning `null` if syncing.
1211        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        // Engine API bitvectors encode the lowest cell indices in the first byte.
1228        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        // Spec requires returning `null` if syncing.
1242        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    /// Metered version of `get_blobs_v2`.
1254    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    /// Metered version of `get_blobs_v3`.
1290    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    /// Metered version of `get_blobs_v4`.
1311    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// This is the concrete ethereum engine API implementation.
1334#[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    /// Handler for `engine_newPayloadV1`
1345    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_newpayloadv1>
1346    /// Caution: This should not accept the `withdrawals` field
1347    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    /// Handler for `engine_newPayloadV2`
1355    /// See also <https://github.com/ethereum/execution-apis/blob/584905270d8ad665718058060267061ecfd79ca5/src/engine/shanghai.md#engine_newpayloadv2>
1356    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    /// Handler for `engine_newPayloadV3`
1367    /// See also <https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/cancun.md#engine_newpayloadv3>
1368    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    /// Handler for `engine_newPayloadV4`
1387    /// See also <https://github.com/ethereum/execution-apis/blob/03911ffc053b8b806123f1fc237184b0092a485a/src/engine/prague.md#engine_newpayloadv4>
1388    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        // Accept requests as a hash only if it is explicitly allowed
1398        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    /// Handler for `engine_newPayloadV5`
1414    ///
1415    /// Post Amsterdam payload handler.
1416    ///
1417    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_newpayloadv5>
1418    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        // Accept requests as a hash only if it is explicitly allowed.
1427        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    /// Handler for `engine_newPayloadV6`.
1443    ///
1444    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/bogota.md#engine_newpayloadv6>
1445    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        // TODO: perform structural validation of the inclusion list transactions and populate
1468        // `inclusion_list_satisfied` for VALID payloads
1469        Ok(self.new_payload_v6_metered(payload).await?)
1470    }
1471
1472    /// Handler for `engine_forkchoiceUpdatedV1`
1473    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_forkchoiceupdatedv1>
1474    ///
1475    /// Caution: This should not accept the `withdrawals` field
1476    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    /// Handler for `engine_forkchoiceUpdatedV2`
1486    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_forkchoiceupdatedv2>
1487    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    /// Handler for `engine_forkchoiceUpdatedV3`
1497    ///
1498    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/cancun.md#engine_forkchoiceupdatedv3>
1499    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    /// Handler for `engine_forkchoiceUpdatedV4`
1509    ///
1510    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_forkchoiceupdatedv4>
1511    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    /// Handler for `engine_forkchoiceUpdatedV5`.
1524    ///
1525    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/bogota.md#engine_forkchoiceupdatedv5>
1526    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    /// Handler for `engine_getPayloadV1`
1539    ///
1540    /// Returns the most recent version of the payload that is available in the corresponding
1541    /// payload build process at the time of receiving this call.
1542    ///
1543    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/paris.md#engine_getPayloadV1>
1544    ///
1545    /// Caution: This should not return the `withdrawals` field
1546    ///
1547    /// Note:
1548    /// > Provider software MAY stop the corresponding build process after serving this call.
1549    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    /// Handler for `engine_getPayloadV2`
1558    ///
1559    /// Returns the most recent version of the payload that is available in the corresponding
1560    /// payload build process at the time of receiving this call.
1561    ///
1562    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_getpayloadv2>
1563    ///
1564    /// Note:
1565    /// > Provider software MAY stop the corresponding build process after serving this call.
1566    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    /// Handler for `engine_getPayloadV3`
1575    ///
1576    /// Returns the most recent version of the payload that is available in the corresponding
1577    /// payload build process at the time of receiving this call.
1578    ///
1579    /// See also <https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/cancun.md#engine_getpayloadv3>
1580    ///
1581    /// Note:
1582    /// > Provider software MAY stop the corresponding build process after serving this call.
1583    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    /// Handler for `engine_getPayloadV4`
1592    ///
1593    /// Returns the most recent version of the payload that is available in the corresponding
1594    /// payload build process at the time of receiving this call.
1595    ///
1596    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/prague.md#engine_getpayloadv4>
1597    ///
1598    /// Note:
1599    /// > Provider software MAY stop the corresponding build process after serving this call.
1600    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    /// Handler for `engine_getPayloadV5`
1609    ///
1610    /// Returns the most recent version of the payload that is available in the corresponding
1611    /// payload build process at the time of receiving this call.
1612    ///
1613    /// See also <https://github.com/ethereum/execution-apis/blob/15399c2e2f16a5f800bf3f285640357e2c245ad9/src/engine/osaka.md#engine_getpayloadv5>
1614    ///
1615    /// Note:
1616    /// > Provider software MAY stop the corresponding build process after serving this call.
1617    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    /// Handler for `engine_getPayloadV6`
1626    ///
1627    /// Post-Amsterdam payload handler that includes Block Access Lists (BAL).
1628    ///
1629    /// See also <https://github.com/ethereum/execution-apis/blob/main/src/engine/amsterdam.md#engine_getpayloadv6>
1630    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    /// Handler for `engine_getInclusionListV1`.
1639    ///
1640    /// See also <https://github.com/ethereum/execution-apis/pull/609>.
1641    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    /// Handler for `engine_getPayloadBodiesByHashV1`
1647    /// See also <https://github.com/ethereum/execution-apis/blob/6452a6b194d7db269bf1dbd087a267251d3cc7f8/src/engine/shanghai.md#engine_getpayloadbodiesbyhashv1>
1648    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    /// Handler for `engine_getPayloadBodiesByHashV2`
1657    ///
1658    /// V2 includes the `block_access_list` field for EIP-7928 BAL support.
1659    ///
1660    /// See also <https://eips.ethereum.org/EIPS/eip-7928>
1661    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    /// Handler for `engine_getPayloadBodiesByRangeV1`
1670    ///
1671    /// See also <https://github.com/ethereum/execution-apis/blob/6452a6b194d7db269bf1dbd087a267251d3cc7f8/src/engine/shanghai.md#engine_getpayloadbodiesbyrangev1>
1672    ///
1673    /// Returns the execution payload bodies by the range starting at `start`, containing `count`
1674    /// blocks.
1675    ///
1676    /// WARNING: This method is associated with the `BeaconBlocksByRange` message in the consensus
1677    /// layer p2p specification, meaning the input should be treated as untrusted or potentially
1678    /// adversarial.
1679    ///
1680    /// Implementers should take care when acting on the input to this method, specifically
1681    /// ensuring that the range is limited properly, and that the range boundaries are computed
1682    /// correctly and without panics.
1683    ///
1684    /// Note: If a block is pre shanghai, `withdrawals` field will be `null`.
1685    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    /// Handler for `engine_getPayloadBodiesByRangeV2`
1695    ///
1696    /// V2 includes the `block_access_list` field for EIP-7928 BAL support.
1697    ///
1698    /// See also <https://eips.ethereum.org/EIPS/eip-7928>
1699    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    /// Handler for `engine_getClientVersionV1`
1709    ///
1710    /// See also <https://github.com/ethereum/execution-apis/blob/03911ffc053b8b806123f1fc237184b0092a485a/src/engine/identification.md>
1711    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    /// Handler for `engine_exchangeCapabilitiesV1`
1720    /// See also <https://github.com/ethereum/execution-apis/blob/6452a6b194d7db269bf1dbd087a267251d3cc7f8/src/engine/common.md#capabilities>
1721    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
1800/// The container type for the engine API internals.
1801struct EngineApiInner<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
1802    /// The provider to interact with the chain.
1803    provider: Provider,
1804    /// Consensus configuration
1805    chain_spec: Arc<ChainSpec>,
1806    /// The channel to send messages to the beacon consensus engine.
1807    beacon_consensus: ConsensusEngineHandle<PayloadT>,
1808    /// The type that can communicate with the payload service to retrieve payloads.
1809    payload_store: PayloadStore<PayloadT>,
1810    /// For spawning and executing async tasks
1811    task_spawner: Runtime,
1812    /// The latency and response type metrics for engine api calls
1813    metrics: EngineApiMetrics,
1814    /// Identification of the execution client used by the consensus client
1815    client: ClientVersionV1,
1816    /// The list of all supported Engine capabilities available over the engine endpoint.
1817    capabilities: EngineCapabilities,
1818    /// Transaction pool.
1819    tx_pool: Pool,
1820    /// Engine validator.
1821    validator: Validator,
1822    accept_execution_requests_hash: bool,
1823    /// Shared blob cell custody bitmap.
1824    cell_custody: CellCustody,
1825    /// Returns `true` if the node is currently syncing.
1826    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    /// Signs candidate transactions because signature lengths affect the final encoded size.
1997    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            // Invalid for V3/Cancun, but should be ignored if forkchoice is SYNCING.
2605            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    // tests covering `engine_getPayloadBodiesByRange` and `engine_getPayloadBodiesByHash`
2706    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                // (start, count)
2717                (0, 0),
2718                (0, 1),
2719                (1, 0),
2720            ];
2721
2722            // test [EngineApiMessage::GetPayloadBodiesByRange]
2723            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            // Insert only blocks in ranges 1-25 and 50-75
2776            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 anything after the second missing range to ensure we don't expect trailing
2791                // `None`s
2792                .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                // ensure we still return trailing `None`s here because by-hash will not be aware
2810                // of the missing block's number, and cannot compare it to the current best block
2811                .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}