Skip to main content

reth_node_ethereum/
engine_ssz_proxy.rs

1//! HTTP SSZ transport proxy for the authenticated Engine API server.
2//!
3//! Implements the [EIP-8178] SSZ Engine API routes under `/engine/v1`.
4//!
5//! [EIP-8178]: https://eips.ethereum.org/EIPS/eip-8178
6
7use crate::engine_ssz_witness::{
8    EngineSszWitness, EngineSszWitnessError, PayloadStatusWithWitness,
9};
10use alloy_consensus::{Transaction, TxEnvelope};
11use alloy_eips::eip7685::Requests;
12use alloy_primitives::{Bytes, B128, B256};
13use alloy_rpc_types_engine::{
14    ssz_engine_types::{
15        BlobsV1Request, BlobsV1Response, BlobsV2Response, BlobsV3Response, BlobsV4Request,
16        BlobsV4Response, BodiesByHashRequest, BodiesResponse, BuiltPayloadAmsterdam,
17        BuiltPayloadOsaka, BuiltPayloadParis, BuiltPayloadPrague, BuiltPayloadShanghai,
18        ExecutionPayloadBodyAmsterdam, ExecutionPayloadBodyParis, ExecutionPayloadBodyShanghai,
19        ExecutionPayloadEnvelopeAmsterdam, ExecutionPayloadEnvelopeCancun,
20        ExecutionPayloadEnvelopeOsaka, ExecutionPayloadEnvelopeParis,
21        ExecutionPayloadEnvelopePrague, ExecutionPayloadEnvelopeShanghai,
22        ForkchoiceUpdateAmsterdam, ForkchoiceUpdateCancun, ForkchoiceUpdateOsaka,
23        ForkchoiceUpdateParis, ForkchoiceUpdatePrague, ForkchoiceUpdateResponse,
24        ForkchoiceUpdateShanghai, Optional, PayloadStatus as EngineSszPayloadStatus,
25        PayloadStatusKind, MAX_BLOBS_REQUEST, MAX_BODIES_REQUEST,
26    },
27    CancunPayloadFields, ExecutionData, ExecutionPayload, ExecutionPayloadBodyV1,
28    ExecutionPayloadFieldV2, ExecutionPayloadSidecar, ForkchoiceState, PayloadAttributes,
29    PayloadId, PraguePayloadFields,
30};
31use futures::future::{BoxFuture, Either};
32use http_body_util::{BodyExt, LengthLimitError, Limited};
33use jsonrpsee::server::{HttpBody, HttpRequest, HttpResponse};
34use reth_chainspec::{EthereumHardfork, EthereumHardforks};
35use reth_engine_primitives::EngineApiValidator;
36use reth_ethereum_engine_primitives::EthEngineTypes;
37use reth_provider::{BalProvider, BlockReader, HeaderProvider, StateProviderFactory};
38use reth_rpc::EngineApi;
39use reth_rpc_engine_api::EngineApiError;
40use reth_tracing::tracing::debug;
41use reth_transaction_pool::TransactionPool;
42use ssz::Decode;
43use std::{
44    future::Future,
45    sync::Arc,
46    task::{Context, Poll},
47};
48use tokio::sync::RwLock;
49use tower::{BoxError, Layer, Service};
50
51const OCTET_STREAM: &str = "application/octet-stream";
52const APPLICATION_JSON: &str = "application/json";
53const CONTENT_TYPE: &str = "content-type";
54const CACHE_CONTROL: &str = "cache-control";
55const ETH_EXECUTION_VERSION: &str = "eth-execution-version";
56
57const STATUS_OK: u16 = 200;
58const STATUS_BAD_REQUEST: u16 = 400;
59const STATUS_CONFLICT: u16 = 409;
60const STATUS_UNPROCESSABLE_ENTITY: u16 = 422;
61const STATUS_NOT_FOUND: u16 = 404;
62const STATUS_METHOD_NOT_ALLOWED: u16 = 405;
63const STATUS_PAYLOAD_TOO_LARGE: u16 = 413;
64const STATUS_INTERNAL_SERVER_ERROR: u16 = 500;
65const STATUS_SERVICE_UNAVAILABLE: u16 = 503;
66const STATUS_UNSUPPORTED_MEDIA_TYPE: u16 = 415;
67
68const MAX_BLOB_REQUEST_BYTES: usize = 4 + 16 + MAX_BLOBS_REQUEST * 32;
69const MAX_BODIES_REQUEST_BYTES: usize = 4 + MAX_BODIES_REQUEST * 32;
70const MAX_PAYLOAD_BYTES: usize = 64 * 1024 * 1024;
71const PROBLEM_JSON: &str = "application/problem+json";
72
73type EthEngineApi<Provider, Pool, Validator, ChainSpec> =
74    EngineApi<Provider, EthEngineTypes, Pool, Validator, ChainSpec>;
75/// Shared handle used by [`EngineSszProxyLayer`].
76pub struct EngineSszProxyHandle<Api = ()> {
77    state: Arc<RwLock<EngineSszState<Api>>>,
78}
79
80impl<Api> Clone for EngineSszProxyHandle<Api> {
81    fn clone(&self) -> Self {
82        Self { state: self.state.clone() }
83    }
84}
85
86impl<Api> std::fmt::Debug for EngineSszProxyHandle<Api> {
87    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
88        f.debug_struct("EngineSszProxyHandle").finish_non_exhaustive()
89    }
90}
91
92impl<Api> EngineSszProxyHandle<Api> {
93    fn new() -> Self {
94        Self {
95            state: Arc::new(RwLock::new(EngineSszState {
96                engine_api: None,
97                witness_handler: None,
98                witness_enabled: false,
99            })),
100        }
101    }
102
103    fn with_engine_api(engine_api: Api) -> Self {
104        Self {
105            state: Arc::new(RwLock::new(EngineSszState {
106                engine_api: Some(engine_api),
107                witness_handler: None,
108                witness_enabled: false,
109            })),
110        }
111    }
112
113    /// Returns whether both the API and generator support the witness extension.
114    pub async fn witness_enabled(&self) -> bool {
115        self.state.read().await.witness_enabled
116    }
117
118    /// Returns the configured witness generator.
119    pub async fn witness_handler(&self) -> Option<Arc<dyn EngineSszWitness>> {
120        self.state.read().await.witness_handler.clone()
121    }
122}
123
124impl<Api: EngineSszApi> EngineSszProxyHandle<Api> {
125    /// Sets the Engine API implementation used by the proxy.
126    pub async fn set_engine_api(&self, engine_api: Api) {
127        let mut state = self.state.write().await;
128        state.engine_api = Some(engine_api);
129        state.update_witness_support();
130    }
131
132    /// Sets the Engine API implementation during synchronous launch wiring.
133    pub fn set_engine_api_sync(&self, engine_api: Api) {
134        let mut state =
135            self.state.try_write().expect("engine api handle should not be locked during launch");
136        state.engine_api = Some(engine_api);
137        state.update_witness_support();
138    }
139
140    /// Sets the witness generator used by `/payloads/witness`.
141    pub async fn set_witness_handler(&self, witness_handler: Arc<dyn EngineSszWitness>) {
142        let mut state = self.state.write().await;
143        state.witness_handler = Some(witness_handler);
144        state.update_witness_support();
145    }
146
147    /// Sets the witness generator during synchronous launch wiring.
148    pub fn set_witness_handler_sync(&self, witness_handler: Arc<dyn EngineSszWitness>) {
149        let mut state =
150            self.state.try_write().expect("witness handle should not be locked during launch");
151        state.witness_handler = Some(witness_handler);
152        state.update_witness_support();
153    }
154}
155
156impl<Api: Clone> EngineSszProxyHandle<Api> {
157    /// Returns the Engine API implementation used by the proxy.
158    pub async fn engine_api(&self) -> Option<Api> {
159        self.state.read().await.engine_api.clone()
160    }
161}
162
163/// A tower layer that intercepts SSZ Engine API routes under `/engine/v1`.
164#[derive(Clone, Debug)]
165pub struct EngineSszProxyLayer<Api = ()> {
166    handle: EngineSszProxyHandle<Api>,
167}
168
169impl<Api> EngineSszProxyLayer<Api> {
170    /// Creates a new proxy layer and a handle for setting the engine after node launch.
171    pub fn new() -> (Self, EngineSszProxyHandle<Api>) {
172        let handle = EngineSszProxyHandle::new();
173        (Self { handle: handle.clone() }, handle)
174    }
175
176    /// Creates a new proxy layer with an Engine API implementation.
177    pub fn with_engine_api(engine_api: Api) -> (Self, EngineSszProxyHandle<Api>) {
178        let handle = EngineSszProxyHandle::with_engine_api(engine_api);
179        (Self { handle: handle.clone() }, handle)
180    }
181}
182
183impl<S, Api> Layer<S> for EngineSszProxyLayer<Api> {
184    type Service = EngineSszProxyService<S, Api>;
185
186    fn layer(&self, inner: S) -> Self::Service {
187        EngineSszProxyService { inner, handle: self.handle.clone() }
188    }
189}
190
191/// The service produced by [`EngineSszProxyLayer`].
192#[derive(Clone, Debug)]
193pub struct EngineSszProxyService<S, Api = ()> {
194    inner: S,
195    handle: EngineSszProxyHandle<Api>,
196}
197
198impl<S, Api> Service<HttpRequest> for EngineSszProxyService<S, Api>
199where
200    S: Service<HttpRequest, Response = HttpResponse, Error = BoxError> + Send + Clone,
201    S::Future: Send,
202    Api: EngineSszApi,
203{
204    type Response = HttpResponse;
205    type Error = BoxError;
206    type Future = Either<S::Future, BoxFuture<'static, Result<HttpResponse, BoxError>>>;
207
208    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
209        self.inner.poll_ready(cx)
210    }
211
212    fn call(&mut self, request: HttpRequest) -> Self::Future {
213        if !request.uri().path().starts_with("/engine/") {
214            return Either::Left(self.inner.call(request))
215        }
216
217        let handle = self.handle.clone();
218        Either::Right(Box::pin(async move { Ok(handle_engine_ssz_request(handle, request).await) }))
219    }
220}
221
222/// Fork selector used by SSZ Engine API request handling.
223#[derive(Clone, Copy, Debug, Eq, PartialEq)]
224pub enum EngineSszFork {
225    /// Paris fork.
226    Paris,
227    /// Shanghai fork.
228    Shanghai,
229    /// Cancun fork.
230    Cancun,
231    /// Prague fork.
232    Prague,
233    /// Osaka fork.
234    Osaka,
235    /// Amsterdam fork.
236    Amsterdam,
237}
238
239impl EngineSszFork {
240    const fn payloads_version(self) -> u8 {
241        match self {
242            Self::Paris => 1,
243            Self::Shanghai => 2,
244            Self::Cancun => 3,
245            Self::Prague | Self::Osaka => 4,
246            Self::Amsterdam => 5,
247        }
248    }
249
250    const fn forkchoice_version(self) -> u8 {
251        match self {
252            Self::Paris => 1,
253            Self::Shanghai => 2,
254            Self::Cancun | Self::Prague | Self::Osaka => 3,
255            Self::Amsterdam => 4,
256        }
257    }
258}
259
260impl std::str::FromStr for EngineSszFork {
261    type Err = ();
262
263    fn from_str(value: &str) -> Result<Self, Self::Err> {
264        match value {
265            "paris" => Ok(Self::Paris),
266            "shanghai" => Ok(Self::Shanghai),
267            "cancun" => Ok(Self::Cancun),
268            "prague" => Ok(Self::Prague),
269            "osaka" => Ok(Self::Osaka),
270            "amsterdam" => Ok(Self::Amsterdam),
271            _ => Err(()),
272        }
273    }
274}
275
276/// API surface required by the SSZ Engine API proxy.
277///
278/// Custom APIs can opt in with a marker implementation; every route, including
279/// `/capabilities`, then answers 404 so clients fall back to the JSON-RPC Engine API.
280pub trait EngineSszApi: Clone + Send + Sync + 'static {
281    /// Returns the capabilities advertisement.
282    ///
283    /// `witness_enabled` is true when a witness generator is configured and
284    /// [`Self::supports_witness`] holds, so the advertisement can include the extension.
285    fn capabilities(&self, _witness_enabled: bool) -> HttpResponse {
286        problem_response(STATUS_NOT_FOUND, "method-not-found", None)
287    }
288
289    /// Whether the implementation supports the witness extension.
290    fn supports_witness(&self) -> bool {
291        false
292    }
293
294    /// Handles the experimental Amsterdam payload submission with a witness.
295    fn new_payload_with_witness(
296        &self,
297        _body: Bytes,
298        _witness_handler: Arc<dyn EngineSszWitness>,
299    ) -> impl Future<Output = HttpResponse> + Send {
300        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
301    }
302
303    /// Returns the Engine API identity response.
304    fn identity(&self) -> HttpResponse {
305        problem_response(STATUS_NOT_FOUND, "method-not-found", None)
306    }
307
308    /// Handles a new payload request.
309    fn new_payload(
310        &self,
311        _fork: EngineSszFork,
312        _body: Bytes,
313    ) -> impl Future<Output = HttpResponse> + Send {
314        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
315    }
316
317    /// Handles a getPayload request.
318    fn get_payload(
319        &self,
320        _fork: EngineSszFork,
321        _payload_id: PayloadId,
322    ) -> impl Future<Output = HttpResponse> + Send {
323        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
324    }
325
326    /// Handles a forkchoice update request.
327    fn forkchoice_updated(
328        &self,
329        _fork: EngineSszFork,
330        _body: Bytes,
331    ) -> impl Future<Output = HttpResponse> + Send {
332        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
333    }
334
335    /// Handles a getPayloadBodiesByHash request.
336    fn get_payload_bodies_by_hash(
337        &self,
338        _fork: EngineSszFork,
339        _hashes: Vec<B256>,
340    ) -> impl Future<Output = HttpResponse> + Send {
341        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
342    }
343
344    /// Handles a getPayloadBodiesByRange request.
345    fn get_payload_bodies_by_range(
346        &self,
347        _fork: EngineSszFork,
348        _start: u64,
349        _count: u64,
350    ) -> impl Future<Output = HttpResponse> + Send {
351        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
352    }
353
354    /// Handles a getBlobs request.
355    fn get_blobs(&self, _version: u8, _body: Bytes) -> impl Future<Output = HttpResponse> + Send {
356        async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
357    }
358}
359
360impl EngineSszApi for reth_node_builder::rpc::NoopEngineApi {}
361
362impl<Provider, Pool, Validator, ChainSpec> EngineSszApi
363    for EthEngineApi<Provider, Pool, Validator, ChainSpec>
364where
365    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
366    Pool: TransactionPool + 'static,
367    Validator: EngineApiValidator<EthEngineTypes>,
368    ChainSpec: EthereumHardforks + Send + Sync + 'static,
369{
370    fn capabilities(&self, witness_enabled: bool) -> HttpResponse {
371        handle_capabilities(witness_enabled)
372    }
373
374    fn identity(&self) -> HttpResponse {
375        json_response(vec![self.client_version().clone()])
376    }
377
378    async fn new_payload(&self, fork: EngineSszFork, body: Bytes) -> HttpResponse {
379        let payload = match decode_new_payload_request(fork, &body) {
380            Ok(payload) => payload,
381            Err(err) => return err.into_response(),
382        };
383
384        match submit_payload(self, fork, payload).await {
385            Ok(status) => ssz_response(status),
386            Err(response) => response,
387        }
388    }
389
390    fn supports_witness(&self) -> bool {
391        true
392    }
393
394    async fn new_payload_with_witness(
395        &self,
396        body: Bytes,
397        witness_handler: Arc<dyn EngineSszWitness>,
398    ) -> HttpResponse {
399        let payload = match decode_new_payload_request(EngineSszFork::Amsterdam, &body) {
400            Ok(payload) => payload,
401            Err(err) => return err.into_response(),
402        };
403        let status = match submit_payload(self, EngineSszFork::Amsterdam, payload.clone()).await {
404            Ok(status) => status,
405            Err(response) => return response,
406        };
407        let witness = match status.status {
408            PayloadStatusKind::Valid => match witness_handler.generate_witness(payload).await {
409                Ok(witness) => Some(witness),
410                // The block is valid but its parent is only known to the engine tree. The
411                // status stays authoritative; resubmitting once forkchoice has made the parent
412                // canonical yields the witness.
413                Err(EngineSszWitnessError::ParentStateUnavailable { parent, source }) => {
414                    debug!(
415                        target: "engine::ssz",
416                        %parent,
417                        %source,
418                        "witness omitted for valid payload"
419                    );
420                    None
421                }
422                Err(err) => {
423                    return problem_response(
424                        STATUS_INTERNAL_SERVER_ERROR,
425                        "internal",
426                        Some(err.to_string()),
427                    )
428                }
429            },
430            _ => None,
431        };
432        ssz_response(PayloadStatusWithWitness::new(status, witness))
433    }
434
435    async fn get_payload(&self, fork: EngineSszFork, payload_id: PayloadId) -> HttpResponse {
436        match fork {
437            EngineSszFork::Paris => match self.get_payload_v2_metered(payload_id).await {
438                Ok(payload) => {
439                    let block_value = payload.block_value;
440                    match payload.execution_payload {
441                        ExecutionPayloadFieldV2::V1(payload) => {
442                            get_payload_response(BuiltPayloadParis { payload, block_value })
443                        }
444                        ExecutionPayloadFieldV2::V2(_) => {
445                            problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
446                        }
447                    }
448                }
449                Err(err) => engine_error_response(err),
450            },
451            EngineSszFork::Shanghai => match self.get_payload_v2_metered(payload_id).await {
452                Ok(payload) => match BuiltPayloadShanghai::try_from(payload) {
453                    Ok(payload) => get_payload_response(payload),
454                    Err(err) => problem_response(
455                        STATUS_UNPROCESSABLE_ENTITY,
456                        "invalid-body",
457                        Some(err.to_string()),
458                    ),
459                },
460                Err(err) => engine_error_response(err),
461            },
462            EngineSszFork::Cancun => self
463                .get_payload_v3_metered(payload_id)
464                .await
465                .map_or_else(engine_error_response, get_payload_response),
466            EngineSszFork::Prague => self
467                .get_payload_v4_metered(payload_id)
468                .await
469                .map(BuiltPayloadPrague::from)
470                .map_or_else(engine_error_response, get_payload_response),
471            EngineSszFork::Osaka => self
472                .get_payload_v5_metered(payload_id)
473                .await
474                .map(BuiltPayloadOsaka::from)
475                .map_or_else(engine_error_response, get_payload_response),
476            EngineSszFork::Amsterdam => self
477                .get_payload_v6_metered(payload_id)
478                .await
479                .map(BuiltPayloadAmsterdam::from)
480                .map_or_else(engine_error_response, get_payload_response),
481        }
482    }
483
484    async fn forkchoice_updated(&self, fork: EngineSszFork, body: Bytes) -> HttpResponse {
485        let (state, attrs, custody_columns) = match decode_forkchoice_request(fork, &body) {
486            Ok(request) => request,
487            Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
488        };
489
490        if attrs.as_ref().is_some_and(|attrs| {
491            !timestamp_matches_fork(self.chain_spec().as_ref(), fork, attrs.timestamp)
492        }) {
493            return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
494        }
495
496        let response = match fork.forkchoice_version() {
497            1 => self.fork_choice_updated_v1_metered(state, attrs).await,
498            2 => self.fork_choice_updated_v2_metered(state, attrs).await,
499            3 => self.fork_choice_updated_v3_metered(state, attrs).await,
500            4 => self.fork_choice_updated_v4_metered(state, attrs, custody_columns).await,
501            _ => return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None),
502        };
503
504        match response {
505            Ok(updated) => match ForkchoiceUpdateResponse::try_from(updated) {
506                Ok(updated) => ssz_response(updated),
507                Err(err) => problem_response(
508                    STATUS_INTERNAL_SERVER_ERROR,
509                    "internal",
510                    Some(err.to_string()),
511                ),
512            },
513            Err(err) => engine_error_response(err),
514        }
515    }
516
517    async fn get_payload_bodies_by_hash(
518        &self,
519        fork: EngineSszFork,
520        hashes: Vec<B256>,
521    ) -> HttpResponse {
522        handle_get_payload_bodies(self.clone(), fork, PayloadBodiesRequest::Hash(hashes)).await
523    }
524
525    async fn get_payload_bodies_by_range(
526        &self,
527        fork: EngineSszFork,
528        start: u64,
529        count: u64,
530    ) -> HttpResponse {
531        handle_get_payload_bodies(self.clone(), fork, PayloadBodiesRequest::Range { start, count })
532            .await
533    }
534
535    async fn get_blobs(&self, version: u8, body: Bytes) -> HttpResponse {
536        if version == 4 {
537            let request = match BlobsV4Request::from_ssz_bytes(&body) {
538                Ok(request) => request,
539                Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
540            };
541            if request.versioned_hashes.len() > MAX_BLOBS_REQUEST {
542                return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
543            }
544            return blob_response::<BlobsV4Response, _>(
545                self.get_blobs_v4_metered(request.versioned_hashes, request.indices_bitarray),
546            )
547        }
548        let request = match BlobsV1Request::from_ssz_bytes(&body) {
549            Ok(request) => request,
550            Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
551        };
552        if request.versioned_hashes.len() > MAX_BLOBS_REQUEST {
553            return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
554        }
555        let hashes = request.versioned_hashes;
556        match version {
557            1 => blob_response::<BlobsV1Response, _>(self.get_blobs_v1_metered(hashes).map(Some)),
558            2 => blob_response::<BlobsV2Response, _>(self.get_blobs_v2_metered(hashes)),
559            3 => blob_response::<BlobsV3Response, _>(self.get_blobs_v3_metered(hashes)),
560            _ => problem_response(STATUS_NOT_FOUND, "method-not-found", None),
561        }
562    }
563}
564
565struct EngineSszState<Api> {
566    engine_api: Option<Api>,
567    witness_handler: Option<Arc<dyn EngineSszWitness>>,
568    witness_enabled: bool,
569}
570
571impl<Api: EngineSszApi> EngineSszState<Api> {
572    fn update_witness_support(&mut self) {
573        self.witness_enabled = self.witness_handler.is_some() &&
574            self.engine_api.as_ref().is_some_and(EngineSszApi::supports_witness);
575    }
576}
577
578async fn handle_engine_ssz_request<Api>(
579    handle: EngineSszProxyHandle<Api>,
580    request: HttpRequest,
581) -> HttpResponse
582where
583    Api: EngineSszApi,
584{
585    let Some(endpoint) = parse_engine_path(request.uri().path()) else {
586        return problem_response(STATUS_NOT_FOUND, "method-not-found", None)
587    };
588
589    if request.method() != endpoint.method() {
590        return problem_response(STATUS_METHOD_NOT_ALLOWED, "method-not-allowed", None)
591    }
592
593    match endpoint {
594        EngineSszEndpoint::Capabilities => {
595            let Some(engine_api) = handle.engine_api().await else {
596                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
597            };
598            engine_api.capabilities(handle.witness_enabled().await)
599        }
600        EngineSszEndpoint::Identity => {
601            let Some(engine_api) = handle.engine_api().await else {
602                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
603            };
604            engine_api.identity()
605        }
606        EngineSszEndpoint::NewPayload => {
607            let Some(fork) = request_fork(&request) else {
608                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
609            };
610            let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
611                Ok(body) => body,
612                Err(response) => return response,
613            };
614            let Some(engine_api) = handle.engine_api().await else {
615                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
616            };
617            engine_api.new_payload(fork, body).await
618        }
619        EngineSszEndpoint::PayloadsWithWitness => {
620            let Some(fork) = request_fork(&request) else {
621                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
622            };
623            if fork != EngineSszFork::Amsterdam {
624                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
625            }
626            let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
627                Ok(body) => body,
628                Err(response) => return response,
629            };
630            let Some(engine_api) = handle.engine_api().await else {
631                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
632            };
633            let Some(witness_handler) = handle.witness_handler().await else {
634                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
635            };
636            engine_api.new_payload_with_witness(body, witness_handler).await
637        }
638        EngineSszEndpoint::GetPayload(payload_id) => {
639            let Ok(payload_id) = payload_id else {
640                return problem_response(STATUS_BAD_REQUEST, "invalid-request", None)
641            };
642            let Some(fork) = request_fork(&request) else {
643                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
644            };
645            let Some(engine_api) = handle.engine_api().await else {
646                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
647            };
648            engine_api.get_payload(fork, payload_id).await
649        }
650        EngineSszEndpoint::Forkchoice => {
651            let Some(fork) = request_fork(&request) else {
652                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
653            };
654            let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
655                Ok(body) => body,
656                Err(response) => return response,
657            };
658            let Some(engine_api) = handle.engine_api().await else {
659                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
660            };
661            engine_api.forkchoice_updated(fork, body).await
662        }
663        EngineSszEndpoint::PayloadBodiesByHash => {
664            let Some(fork) = request_fork(&request) else {
665                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
666            };
667            let body = match read_ssz_body(request, MAX_BODIES_REQUEST_BYTES).await {
668                Ok(body) => body,
669                Err(response) => return response,
670            };
671            let request = match BodiesByHashRequest::from_ssz_bytes(&body) {
672                Ok(request) if request.block_hashes.len() <= MAX_BODIES_REQUEST => {
673                    request.block_hashes
674                }
675                Ok(_) => {
676                    return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
677                }
678                Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
679            };
680            let Some(engine_api) = handle.engine_api().await else {
681                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
682            };
683            engine_api.get_payload_bodies_by_hash(fork, request).await
684        }
685        EngineSszEndpoint::PayloadBodiesByRange => {
686            let Some(fork) = request_fork(&request) else {
687                return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
688            };
689            let (start, count) =
690                match parse_bodies_range_query(request.uri().query().unwrap_or_default()) {
691                    Ok(range) => range,
692                    Err(response) => return response,
693                };
694            let Some(engine_api) = handle.engine_api().await else {
695                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
696            };
697            engine_api.get_payload_bodies_by_range(fork, start, count).await
698        }
699        EngineSszEndpoint::Blobs(version) => {
700            let body = match read_ssz_body(request, MAX_BLOB_REQUEST_BYTES).await {
701                Ok(body) => body,
702                Err(response) => return response,
703            };
704            let Some(engine_api) = handle.engine_api().await else {
705                return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
706            };
707            engine_api.get_blobs(version, body).await
708        }
709    }
710}
711
712fn parse_method_version(version: &str) -> Option<u8> {
713    version.strip_prefix('v')?.parse().ok().filter(|version| (1..=4).contains(version))
714}
715
716fn request_fork(request: &HttpRequest) -> Option<EngineSszFork> {
717    request.headers().get(ETH_EXECUTION_VERSION)?.to_str().ok()?.parse().ok()
718}
719
720fn parse_engine_path(path: &str) -> Option<EngineSszEndpoint> {
721    let mut segments = path.trim_start_matches('/').split('/');
722    match (segments.next(), segments.next(), segments.next(), segments.next(), segments.next()) {
723        (Some("engine"), Some("v1"), Some("capabilities"), None, None) => {
724            Some(EngineSszEndpoint::Capabilities)
725        }
726        (Some("engine"), Some("v1"), Some("identity"), None, None) => {
727            Some(EngineSszEndpoint::Identity)
728        }
729        (Some("engine"), Some("v1"), Some("payloads"), None, None) => {
730            Some(EngineSszEndpoint::NewPayload)
731        }
732        (Some("engine"), Some("v1"), Some("payloads"), Some("witness"), None) => {
733            Some(EngineSszEndpoint::PayloadsWithWitness)
734        }
735        (Some("engine"), Some("v1"), Some("payloads"), Some(payload_id), None) => {
736            let payload_id = payload_id.parse::<PayloadId>();
737            Some(EngineSszEndpoint::GetPayload(payload_id))
738        }
739        (Some("engine"), Some("v1"), Some("forkchoice"), None, None) => {
740            Some(EngineSszEndpoint::Forkchoice)
741        }
742        (Some("engine"), Some("v1"), Some("blobs"), version, None) => {
743            Some(EngineSszEndpoint::Blobs(parse_method_version(version?)?))
744        }
745        (Some("engine"), Some("v1"), Some("bodies"), Some("hash"), None) => {
746            Some(EngineSszEndpoint::PayloadBodiesByHash)
747        }
748        (Some("engine"), Some("v1"), Some("bodies"), None, None) => {
749            Some(EngineSszEndpoint::PayloadBodiesByRange)
750        }
751        _ => None,
752    }
753}
754
755#[derive(Clone, Copy, Debug, Eq, PartialEq)]
756enum EngineSszEndpoint {
757    Capabilities,
758    Identity,
759    NewPayload,
760    PayloadsWithWitness,
761    GetPayload(Result<PayloadId, <PayloadId as std::str::FromStr>::Err>),
762    Forkchoice,
763    PayloadBodiesByHash,
764    PayloadBodiesByRange,
765    Blobs(u8),
766}
767
768impl EngineSszEndpoint {
769    const fn method(&self) -> &'static str {
770        match self {
771            Self::Capabilities |
772            Self::Identity |
773            Self::GetPayload(_) |
774            Self::PayloadBodiesByRange => "GET",
775            Self::NewPayload |
776            Self::PayloadsWithWitness |
777            Self::Forkchoice |
778            Self::PayloadBodiesByHash |
779            Self::Blobs(_) => "POST",
780        }
781    }
782}
783
784fn handle_capabilities(witness_enabled: bool) -> HttpResponse {
785    let mut fork_scoped_endpoints = vec!["payloads", "forkchoice", "bodies"];
786    if witness_enabled {
787        fork_scoped_endpoints.push("payloads/witness");
788    }
789    json_response(serde_json::json!({
790        "supported_forks": ["paris", "shanghai", "cancun", "prague", "osaka", "amsterdam"],
791        "fork_scoped_endpoints": fork_scoped_endpoints,
792        "independently_versioned": {
793            "blobs": ["v1", "v2", "v3", "v4"],
794        },
795        "unscoped_endpoints": ["capabilities", "identity"],
796        "limits": {
797            "bodies.max_count": MAX_BODIES_REQUEST,
798            "blobs.max_versioned_hashes": MAX_BLOBS_REQUEST,
799            "payload.max_bytes": MAX_PAYLOAD_BYTES,
800        },
801    }))
802}
803
804async fn submit_payload<Provider, Pool, Validator, ChainSpec>(
805    engine_api: &EthEngineApi<Provider, Pool, Validator, ChainSpec>,
806    fork: EngineSszFork,
807    payload: ExecutionData,
808) -> Result<EngineSszPayloadStatus, HttpResponse>
809where
810    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
811    Pool: TransactionPool + 'static,
812    Validator: EngineApiValidator<EthEngineTypes>,
813    ChainSpec: EthereumHardforks + Send + Sync + 'static,
814{
815    let status = match fork.payloads_version() {
816        1 => engine_api.new_payload_v1(payload).await,
817        2 => engine_api.new_payload_v2(payload).await,
818        3 => engine_api.new_payload_v3(payload).await,
819        4 => engine_api.new_payload_v4(payload).await,
820        5 => engine_api.new_payload_v5(payload).await,
821        _ => return Err(problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)),
822    }
823    .map_err(engine_error_response)?;
824    EngineSszPayloadStatus::try_from(status).map_err(|error| {
825        problem_response(STATUS_INTERNAL_SERVER_ERROR, "internal", Some(error.to_string()))
826    })
827}
828
829fn engine_error_response(err: EngineApiError) -> HttpResponse {
830    let detail = err.to_string();
831    let error: jsonrpsee::types::ErrorObjectOwned = err.into();
832    let (status, problem_type) = match error.code() {
833        -32700 => (STATUS_BAD_REQUEST, "parse-error"),
834        -32600 => (STATUS_BAD_REQUEST, "invalid-request"),
835        -32601 => (STATUS_NOT_FOUND, "method-not-found"),
836        -32602 => (STATUS_UNPROCESSABLE_ENTITY, "invalid-body"),
837        -38001 => (STATUS_NOT_FOUND, "unknown-payload"),
838        -38002 => (STATUS_CONFLICT, "invalid-forkchoice"),
839        -38003 => (STATUS_UNPROCESSABLE_ENTITY, "invalid-attributes"),
840        -38004 => (STATUS_PAYLOAD_TOO_LARGE, "request-too-large"),
841        -38005 => (STATUS_BAD_REQUEST, "unsupported-fork"),
842        -38006 => (STATUS_CONFLICT, "reorg-too-deep"),
843        _ => (STATUS_INTERNAL_SERVER_ERROR, "internal"),
844    };
845    problem_response(status, problem_type, Some(detail))
846}
847
848enum PayloadBodiesRequest {
849    Hash(Vec<B256>),
850    Range { start: u64, count: u64 },
851}
852
853async fn handle_get_payload_bodies<Provider, Pool, Validator, ChainSpec>(
854    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
855    fork: EngineSszFork,
856    request: PayloadBodiesRequest,
857) -> HttpResponse
858where
859    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
860    Pool: TransactionPool + 'static,
861    Validator: EngineApiValidator<EthEngineTypes>,
862    ChainSpec: EthereumHardforks + Send + Sync + 'static,
863{
864    let include_bal = fork == EngineSszFork::Amsterdam;
865    let response = match request {
866        PayloadBodiesRequest::Hash(hashes) => {
867            engine_api.get_payload_bodies_by_hash_with_timestamps_metered(hashes, include_bal).await
868        }
869        PayloadBodiesRequest::Range { start, count } => {
870            engine_api
871                .get_payload_bodies_by_range_with_timestamps_metered(start, count, include_bal)
872                .await
873        }
874    };
875    let chain_spec = engine_api.chain_spec().as_ref();
876    match fork {
877        EngineSszFork::Amsterdam => payload_bodies_http_response(
878            response,
879            |body| ExecutionPayloadBodyAmsterdam::try_from(body).ok(),
880            fork,
881            chain_spec,
882        ),
883        EngineSszFork::Paris => payload_bodies_http_response(
884            response,
885            |body| ExecutionPayloadBodyParis::try_from(ExecutionPayloadBodyV1::from(body)).ok(),
886            fork,
887            chain_spec,
888        ),
889        _ => payload_bodies_http_response(
890            response,
891            |body| ExecutionPayloadBodyShanghai::try_from(ExecutionPayloadBodyV1::from(body)).ok(),
892            fork,
893            chain_spec,
894        ),
895    }
896}
897
898fn parse_bodies_range_query(query: &str) -> Result<(u64, u64), HttpResponse> {
899    let invalid = || problem_response(STATUS_BAD_REQUEST, "invalid-request", None);
900    let mut start = None;
901    let mut count = None;
902    for pair in query.split('&') {
903        let (key, value) = pair.split_once('=').ok_or_else(invalid)?;
904        let field = match key {
905            "from" if start.is_none() => &mut start,
906            "count" if count.is_none() => &mut count,
907            _ => return Err(invalid()),
908        };
909        *field = Some(value.parse::<u64>().map_err(|_| invalid())?);
910    }
911    let range = (start.ok_or_else(invalid)?, count.ok_or_else(invalid)?);
912    if range.1 > MAX_BODIES_REQUEST as u64 {
913        return Err(problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None))
914    }
915    Ok(range)
916}
917
918/// Amsterdam bodies with a missing or pruned BAL are unavailable, just like pruned blocks.
919fn payload_bodies_response<LegacyBody, ForkBody>(
920    response: Result<Vec<Option<(u64, LegacyBody)>>, EngineApiError>,
921    convert: impl Fn(LegacyBody) -> Option<ForkBody>,
922    fork: EngineSszFork,
923    chain_spec: &impl EthereumHardforks,
924) -> Result<BodiesResponse<ForkBody>, HttpResponse>
925where
926    ForkBody: Default + ssz::Encode + ssz::Decode,
927{
928    let bodies = response.map_err(engine_error_response)?;
929    let bodies = bodies
930        .into_iter()
931        .map(|body| {
932            body.filter(|(timestamp, _)| timestamp_matches_fork(chain_spec, fork, *timestamp))
933                .map(|(_, body)| body)
934        })
935        .collect();
936    BodiesResponse::from_optional_bodies(bodies, convert).map_err(|error| {
937        problem_response(STATUS_INTERNAL_SERVER_ERROR, "internal", Some(error.to_string()))
938    })
939}
940
941fn payload_bodies_http_response<LegacyBody, ForkBody>(
942    response: Result<Vec<Option<(u64, LegacyBody)>>, EngineApiError>,
943    convert: impl Fn(LegacyBody) -> Option<ForkBody>,
944    fork: EngineSszFork,
945    chain_spec: &impl EthereumHardforks,
946) -> HttpResponse
947where
948    ForkBody: Default + ssz::Encode + ssz::Decode,
949{
950    match payload_bodies_response(response, convert, fork, chain_spec) {
951        Ok(response) => ssz_response(response),
952        Err(response) => response,
953    }
954}
955
956fn timestamp_matches_fork<ChainSpec: EthereumHardforks>(
957    chain_spec: &ChainSpec,
958    fork: EngineSszFork,
959    timestamp: u64,
960) -> bool {
961    let active = |fork| chain_spec.is_ethereum_fork_active_at_timestamp(fork, timestamp);
962    match fork {
963        EngineSszFork::Paris => !active(EthereumHardfork::Shanghai),
964        EngineSszFork::Shanghai => {
965            active(EthereumHardfork::Shanghai) && !active(EthereumHardfork::Cancun)
966        }
967        EngineSszFork::Cancun => {
968            active(EthereumHardfork::Cancun) && !active(EthereumHardfork::Prague)
969        }
970        EngineSszFork::Prague => {
971            active(EthereumHardfork::Prague) && !active(EthereumHardfork::Osaka)
972        }
973        EngineSszFork::Osaka => {
974            active(EthereumHardfork::Osaka) && !active(EthereumHardfork::Amsterdam)
975        }
976        EngineSszFork::Amsterdam => active(EthereumHardfork::Amsterdam),
977    }
978}
979
980async fn read_ssz_body(request: HttpRequest, max_bytes: usize) -> Result<Bytes, HttpResponse> {
981    let content_type = request.headers().get(CONTENT_TYPE).and_then(|value| value.to_str().ok());
982    if content_type != Some(OCTET_STREAM) {
983        return Err(problem_response(STATUS_UNSUPPORTED_MEDIA_TYPE, "unsupported-media-type", None))
984    }
985
986    if let Some(content_length) = request.headers().get("content-length") {
987        let Some(content_length) =
988            content_length.to_str().ok().and_then(|value| value.parse::<usize>().ok())
989        else {
990            return Err(problem_response(STATUS_BAD_REQUEST, "invalid-request", None))
991        };
992        if content_length > max_bytes {
993            return Err(problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None))
994        }
995    }
996
997    match Limited::new(request.into_body(), max_bytes).collect().await {
998        Ok(body) => Ok(body.to_bytes().into()),
999        Err(err) if err.downcast_ref::<LengthLimitError>().is_some() => {
1000            Err(problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None))
1001        }
1002        Err(_) => Err(problem_response(STATUS_BAD_REQUEST, "invalid-request", None)),
1003    }
1004}
1005
1006fn blob_response<Ssz, Legacy>(response: Result<Option<Legacy>, EngineApiError>) -> HttpResponse
1007where
1008    Ssz: TryFrom<Legacy> + ssz::Encode,
1009    Ssz::Error: std::fmt::Display,
1010{
1011    match response {
1012        Ok(Some(response)) => match Ssz::try_from(response) {
1013            Ok(response) => ssz_response(response),
1014            Err(err) => {
1015                problem_response(STATUS_INTERNAL_SERVER_ERROR, "internal", Some(err.to_string()))
1016            }
1017        },
1018        Ok(None) => no_content_response(),
1019        Err(err) => engine_error_response(err),
1020    }
1021}
1022
1023/// Structural SSZ errors and invalid values have distinct REST error codes.
1024#[derive(Debug)]
1025enum PayloadDecodeError {
1026    Ssz(ssz::DecodeError),
1027    InvalidTransaction(alloy_eips::eip2718::Eip2718Error),
1028}
1029
1030impl From<ssz::DecodeError> for PayloadDecodeError {
1031    fn from(error: ssz::DecodeError) -> Self {
1032        Self::Ssz(error)
1033    }
1034}
1035
1036impl PayloadDecodeError {
1037    fn into_response(self) -> HttpResponse {
1038        match self {
1039            Self::Ssz(error) => {
1040                problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", Some(format!("{error:?}")))
1041            }
1042            Self::InvalidTransaction(error) => problem_response(
1043                STATUS_UNPROCESSABLE_ENTITY,
1044                "invalid-body",
1045                Some(error.to_string()),
1046            ),
1047        }
1048    }
1049}
1050
1051fn check_ssz_bound(actual: usize, max: usize, field: &str) -> Result<(), ssz::DecodeError> {
1052    if actual > max {
1053        return Err(ssz::DecodeError::BytesInvalid(format!(
1054            "{field} exceeds SSZ bound: {actual} > {max}"
1055        )))
1056    }
1057    Ok(())
1058}
1059
1060fn decode_new_payload_request(
1061    fork: EngineSszFork,
1062    body: &[u8],
1063) -> Result<ExecutionData, PayloadDecodeError> {
1064    let (payload, parent_root, requests): (ExecutionPayload, _, Option<Requests>) = match fork {
1065        EngineSszFork::Paris => {
1066            let envelope = ExecutionPayloadEnvelopeParis::from_ssz_bytes(body)?;
1067            (envelope.payload.into(), None, None)
1068        }
1069        EngineSszFork::Shanghai => {
1070            let envelope = ExecutionPayloadEnvelopeShanghai::from_ssz_bytes(body)?;
1071            (envelope.payload.into(), None, None)
1072        }
1073        EngineSszFork::Cancun => {
1074            let envelope = ExecutionPayloadEnvelopeCancun::from_ssz_bytes(body)?;
1075            (envelope.payload.into(), Some(envelope.parent_beacon_block_root), None)
1076        }
1077        EngineSszFork::Prague => {
1078            let envelope = ExecutionPayloadEnvelopePrague::from_ssz_bytes(body)?;
1079            (
1080                envelope.payload.into(),
1081                Some(envelope.parent_beacon_block_root),
1082                Some(envelope.execution_requests),
1083            )
1084        }
1085        EngineSszFork::Osaka => {
1086            let envelope = ExecutionPayloadEnvelopeOsaka::from_ssz_bytes(body)?;
1087            (
1088                envelope.payload.into(),
1089                Some(envelope.parent_beacon_block_root),
1090                Some(envelope.execution_requests),
1091            )
1092        }
1093        EngineSszFork::Amsterdam => {
1094            let envelope = ExecutionPayloadEnvelopeAmsterdam::from_ssz_bytes(body)?;
1095            (
1096                ExecutionPayload::V4(envelope.payload),
1097                Some(envelope.parent_beacon_block_root),
1098                Some(envelope.execution_requests),
1099            )
1100        }
1101    };
1102
1103    // Alloy's codecs reproduce the wire layout but do not enforce the SSZ list bounds.
1104    // Check those independently of the smaller, advertised total HTTP request limit.
1105    let inner = payload.as_v1();
1106    check_ssz_bound(inner.extra_data.len(), 32, "extra_data")?;
1107    check_ssz_bound(inner.transactions.len(), 1 << 20, "transactions")?;
1108    for transaction in &inner.transactions {
1109        check_ssz_bound(transaction.len(), 1 << 30, "transaction")?;
1110    }
1111    if let Some(withdrawals) = payload.withdrawals() {
1112        check_ssz_bound(withdrawals.len(), 16, "withdrawals")?;
1113    }
1114    if let Some(bal) = payload.block_access_list() {
1115        check_ssz_bound(bal.len(), 1 << 30, "block_access_list")?;
1116    }
1117    if let Some(requests) = &requests {
1118        check_ssz_bound(requests.len(), 256, "execution_requests")?;
1119        for request in requests.iter() {
1120            check_ssz_bound(request.len(), 1 << 30, "execution_request")?;
1121        }
1122    }
1123
1124    let versioned_hashes = calculate_versioned_hashes(&payload)?;
1125    let sidecar = match parent_root {
1126        Some(parent_beacon_block_root) => {
1127            let cancun = CancunPayloadFields { parent_beacon_block_root, versioned_hashes };
1128            match requests {
1129                Some(requests) => {
1130                    ExecutionPayloadSidecar::v4(cancun, PraguePayloadFields::new(requests))
1131                }
1132                None => ExecutionPayloadSidecar::v3(cancun),
1133            }
1134        }
1135        None => ExecutionPayloadSidecar::none(),
1136    };
1137    Ok(ExecutionData::new(payload, sidecar))
1138}
1139
1140fn calculate_versioned_hashes(payload: &ExecutionPayload) -> Result<Vec<B256>, PayloadDecodeError> {
1141    let mut versioned_hashes = Vec::new();
1142    for transaction in payload.decoded_transactions::<TxEnvelope>() {
1143        let transaction = transaction.map_err(PayloadDecodeError::InvalidTransaction)?;
1144        if let Some(hashes) = transaction.blob_versioned_hashes() {
1145            versioned_hashes.extend_from_slice(hashes);
1146        }
1147    }
1148    Ok(versioned_hashes)
1149}
1150
1151fn decode_forkchoice_request(
1152    fork: EngineSszFork,
1153    body: &[u8],
1154) -> Result<(ForkchoiceState, Option<PayloadAttributes>, Option<B128>), ssz::DecodeError> {
1155    let (state, attrs, custody) = match fork {
1156        EngineSszFork::Paris => {
1157            let ForkchoiceUpdateParis { forkchoice_state, payload_attributes } =
1158                ForkchoiceUpdateParis::from_ssz_bytes(body)?;
1159            (forkchoice_state, optional_attrs(payload_attributes), None)
1160        }
1161        EngineSszFork::Shanghai => {
1162            let ForkchoiceUpdateShanghai { forkchoice_state, payload_attributes } =
1163                ForkchoiceUpdateShanghai::from_ssz_bytes(body)?;
1164            (forkchoice_state, optional_attrs(payload_attributes), None)
1165        }
1166        EngineSszFork::Cancun => {
1167            let ForkchoiceUpdateCancun { forkchoice_state, payload_attributes } =
1168                ForkchoiceUpdateCancun::from_ssz_bytes(body)?;
1169            (forkchoice_state, optional_attrs(payload_attributes), None)
1170        }
1171        EngineSszFork::Prague => {
1172            let ForkchoiceUpdatePrague { forkchoice_state, payload_attributes } =
1173                ForkchoiceUpdatePrague::from_ssz_bytes(body)?;
1174            (forkchoice_state, optional_attrs(payload_attributes), None)
1175        }
1176        EngineSszFork::Osaka => {
1177            let ForkchoiceUpdateOsaka { forkchoice_state, payload_attributes } =
1178                ForkchoiceUpdateOsaka::from_ssz_bytes(body)?;
1179            (forkchoice_state, optional_attrs(payload_attributes), None)
1180        }
1181        EngineSszFork::Amsterdam => {
1182            let ForkchoiceUpdateAmsterdam { forkchoice_state, payload_attributes, custody_columns } =
1183                ForkchoiceUpdateAmsterdam::from_ssz_bytes(body)?;
1184            (forkchoice_state, optional_attrs(payload_attributes), custody_columns.into_option())
1185        }
1186    };
1187    if let Some(withdrawals) =
1188        attrs.as_ref().and_then(|attrs: &PayloadAttributes| attrs.withdrawals.as_ref())
1189    {
1190        check_ssz_bound(withdrawals.len(), 16, "withdrawals")?;
1191    }
1192    Ok((state, attrs, custody))
1193}
1194
1195fn optional_attrs<T>(attrs: Optional<T>) -> Option<PayloadAttributes>
1196where
1197    T: Into<PayloadAttributes>,
1198{
1199    attrs.into_option().map(Into::into)
1200}
1201
1202fn ssz_response<T: ssz::Encode>(value: T) -> HttpResponse {
1203    HttpResponse::builder()
1204        .status(STATUS_OK)
1205        .header(CONTENT_TYPE, OCTET_STREAM)
1206        .body(HttpBody::from(value.as_ssz_bytes()))
1207        .expect("valid response")
1208}
1209
1210fn get_payload_response<T: ssz::Encode>(value: T) -> HttpResponse {
1211    let mut response = ssz_response(value);
1212    response.headers_mut().insert(CACHE_CONTROL, "no-store".parse().expect("valid cache control"));
1213    response
1214}
1215
1216fn json_response<T: serde::Serialize>(value: T) -> HttpResponse {
1217    let Ok(body) = serde_json::to_string(&value) else {
1218        return problem_response(
1219            STATUS_INTERNAL_SERVER_ERROR,
1220            "internal",
1221            Some("failed to encode json".to_string()),
1222        )
1223    };
1224
1225    HttpResponse::builder()
1226        .status(STATUS_OK)
1227        .header(CONTENT_TYPE, APPLICATION_JSON)
1228        .body(HttpBody::from(body))
1229        .expect("valid response")
1230}
1231
1232fn no_content_response() -> HttpResponse {
1233    HttpResponse::builder().status(204).body(HttpBody::empty()).expect("valid response")
1234}
1235
1236fn problem_response(
1237    status: u16,
1238    problem_type: &'static str,
1239    detail: Option<String>,
1240) -> HttpResponse {
1241    let problem_type = format!("/engine-api/errors/{problem_type}");
1242    let body = match detail {
1243        Some(detail) => serde_json::json!({ "type": problem_type, "detail": detail }),
1244        None => serde_json::json!({ "type": problem_type }),
1245    };
1246
1247    HttpResponse::builder()
1248        .status(status)
1249        .header(CONTENT_TYPE, PROBLEM_JSON)
1250        .body(HttpBody::from(body.to_string()))
1251        .expect("valid response")
1252}
1253
1254#[cfg(test)]
1255mod tests {
1256    use super::*;
1257    use alloy_rpc_types_engine::ssz_engine_types::{
1258        PayloadAttributesAmsterdam, PayloadAttributesCancun, PayloadAttributesParis,
1259        PayloadAttributesShanghai,
1260    };
1261    use ssz::Encode;
1262
1263    #[tokio::test]
1264    async fn witness_capabilities_follow_both_wiring_orders() {
1265        #[derive(Clone)]
1266        struct Api(bool);
1267        impl EngineSszApi for Api {
1268            fn supports_witness(&self) -> bool {
1269                self.0
1270            }
1271        }
1272        struct Witness;
1273        impl EngineSszWitness for Witness {
1274            fn generate_witness(
1275                &self,
1276                _: ExecutionData,
1277            ) -> BoxFuture<
1278                'static,
1279                Result<crate::engine_ssz_witness::ExecutionWitnessV1, EngineSszWitnessError>,
1280            > {
1281                Box::pin(async { Ok(Default::default()) })
1282            }
1283        }
1284        let handle = EngineSszProxyHandle::new();
1285        handle.set_witness_handler_sync(Arc::new(Witness));
1286        assert!(!handle.witness_enabled().await);
1287        handle.set_engine_api_sync(Api(true));
1288        assert!(handle.witness_enabled().await);
1289        handle.set_engine_api(Api(false)).await;
1290        assert!(!handle.witness_enabled().await);
1291
1292        let handle = EngineSszProxyHandle::with_engine_api(Api(true));
1293        assert!(!handle.witness_enabled().await);
1294        handle.set_witness_handler(Arc::new(Witness)).await;
1295        assert!(handle.witness_enabled().await);
1296    }
1297
1298    #[tokio::test]
1299    async fn witness_is_only_advertised_when_configured() {
1300        for enabled in [false, true] {
1301            let body = handle_capabilities(enabled).into_body().collect().await.unwrap().to_bytes();
1302            let capabilities: serde_json::Value = serde_json::from_slice(&body).unwrap();
1303            let endpoints = capabilities["fork_scoped_endpoints"].as_array().unwrap();
1304            assert_eq!(endpoints.iter().any(|endpoint| endpoint == "payloads/witness"), enabled);
1305        }
1306        assert_eq!(
1307            parse_engine_path("/engine/v1/payloads/witness"),
1308            Some(EngineSszEndpoint::PayloadsWithWitness)
1309        );
1310    }
1311
1312    #[test]
1313    fn payload_schema_bounds() {
1314        use alloy_rpc_types_engine::{ExecutionPayloadV1, ExecutionPayloadV2, ExecutionPayloadV3};
1315        let mut payload = ExecutionPayloadV3 {
1316            payload_inner: ExecutionPayloadV2 {
1317                payload_inner: ExecutionPayloadV1::from_block_unchecked(
1318                    B256::ZERO,
1319                    &reth_ethereum_primitives::Block::default(),
1320                ),
1321                withdrawals: vec![],
1322            },
1323            blob_gas_used: 0,
1324            excess_blob_gas: 0,
1325        };
1326        payload.payload_inner.payload_inner.extra_data = vec![0; 32].into();
1327        payload.payload_inner.withdrawals = vec![Default::default(); 16];
1328        let mut envelope = ExecutionPayloadEnvelopePrague {
1329            payload,
1330            parent_beacon_block_root: B256::ZERO,
1331            execution_requests: Requests::new(vec![Bytes::new(); 256]),
1332        };
1333        assert!(decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()).is_ok());
1334        envelope.execution_requests = Requests::new(vec![Bytes::new(); 257]);
1335        assert!(matches!(
1336            decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1337            Err(PayloadDecodeError::Ssz(_))
1338        ));
1339        envelope.execution_requests = Requests::default();
1340        envelope.payload.payload_inner.withdrawals.push(Default::default());
1341        assert!(matches!(
1342            decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1343            Err(PayloadDecodeError::Ssz(_))
1344        ));
1345        envelope.payload.payload_inner.withdrawals.clear();
1346        envelope.payload.payload_inner.payload_inner.extra_data = vec![0; 33].into();
1347        assert!(matches!(
1348            decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1349            Err(PayloadDecodeError::Ssz(_))
1350        ));
1351        let shanghai = ExecutionPayloadEnvelopeShanghai {
1352            payload: ExecutionPayloadV2 {
1353                payload_inner: ExecutionPayloadV1::from_block_unchecked(
1354                    B256::ZERO,
1355                    &reth_ethereum_primitives::Block::default(),
1356                ),
1357                withdrawals: vec![Default::default(); 17],
1358            },
1359        };
1360        assert!(matches!(
1361            decode_new_payload_request(EngineSszFork::Shanghai, &shanghai.as_ssz_bytes()),
1362            Err(PayloadDecodeError::Ssz(_))
1363        ));
1364        let paris = ExecutionPayloadEnvelopeParis {
1365            payload: ExecutionPayloadV1 {
1366                transactions: vec![Bytes::new(); (1 << 20) + 1],
1367                ..ExecutionPayloadV1::from_block_unchecked(
1368                    B256::ZERO,
1369                    &reth_ethereum_primitives::Block::default(),
1370                )
1371            },
1372        };
1373        assert!(matches!(
1374            decode_new_payload_request(EngineSszFork::Paris, &paris.as_ssz_bytes()),
1375            Err(PayloadDecodeError::Ssz(_))
1376        ));
1377    }
1378
1379    #[test]
1380    fn forkchoice_withdrawal_bounds() {
1381        for count in [16, 17] {
1382            let withdrawals = vec![Default::default(); count];
1383            let state = ForkchoiceState::default();
1384            let shanghai = ForkchoiceUpdateShanghai {
1385                forkchoice_state: state,
1386                payload_attributes: Optional::some(PayloadAttributesShanghai {
1387                    withdrawals: withdrawals.clone(),
1388                    ..Default::default()
1389                }),
1390            };
1391            let cancun = ForkchoiceUpdateCancun {
1392                forkchoice_state: state,
1393                payload_attributes: Optional::some(PayloadAttributesCancun {
1394                    withdrawals: withdrawals.clone(),
1395                    ..Default::default()
1396                }),
1397            };
1398            let amsterdam = ForkchoiceUpdateAmsterdam {
1399                forkchoice_state: state,
1400                payload_attributes: Optional::some(PayloadAttributesAmsterdam {
1401                    withdrawals,
1402                    ..Default::default()
1403                }),
1404                custody_columns: Optional::none(),
1405            };
1406            for (fork, bytes) in [
1407                (EngineSszFork::Shanghai, shanghai.as_ssz_bytes()),
1408                (EngineSszFork::Cancun, cancun.as_ssz_bytes()),
1409                (EngineSszFork::Prague, cancun.as_ssz_bytes()),
1410                (EngineSszFork::Osaka, cancun.as_ssz_bytes()),
1411                (EngineSszFork::Amsterdam, amsterdam.as_ssz_bytes()),
1412            ] {
1413                assert_eq!(
1414                    decode_forkchoice_request(fork, &bytes).is_ok(),
1415                    count == 16,
1416                    "{fork:?}"
1417                );
1418            }
1419        }
1420    }
1421
1422    #[test]
1423    fn paris_forkchoice_uses_fixed_size_optional_attributes() {
1424        let request = ForkchoiceUpdateParis {
1425            forkchoice_state: ForkchoiceState::default(),
1426            payload_attributes: Optional::some(PayloadAttributesParis {
1427                timestamp: 1_700_000_000,
1428                ..Default::default()
1429            }),
1430        };
1431        let bytes = request.as_ssz_bytes();
1432        assert_eq!(bytes.len(), 160);
1433        assert_eq!(&bytes[96..100], &100u32.to_le_bytes());
1434        assert_eq!(&bytes[100..108], &1_700_000_000u64.to_le_bytes());
1435        assert_eq!(
1436            decode_forkchoice_request(EngineSszFork::Paris, &bytes).unwrap().1.unwrap().timestamp,
1437            1_700_000_000
1438        );
1439    }
1440
1441    #[tokio::test]
1442    async fn post_body_limits_apply_without_content_length() {
1443        for size in [16, 17] {
1444            let request = HttpRequest::builder()
1445                .header(CONTENT_TYPE, OCTET_STREAM)
1446                .body(HttpBody::from(vec![0; size]))
1447                .unwrap();
1448            let response = read_ssz_body(request, 16).await;
1449            if size == 16 {
1450                assert_eq!(response.unwrap().len(), size);
1451            } else {
1452                assert_eq!(response.unwrap_err().status(), STATUS_PAYLOAD_TOO_LARGE);
1453            }
1454        }
1455    }
1456
1457    #[tokio::test]
1458    async fn engine_errors_preserve_validation_semantics() {
1459        use alloy_rpc_types_engine::ForkchoiceUpdateError;
1460        use reth_engine_primitives::BeaconForkChoiceUpdateError;
1461        for (error, status, kind) in [
1462            (EngineApiError::UnknownPayload, 404, "unknown-payload"),
1463            (EngineApiError::BlobRequestTooLarge { len: 129 }, 413, "request-too-large"),
1464            (
1465                EngineApiError::ForkChoiceUpdate(
1466                    BeaconForkChoiceUpdateError::ForkchoiceUpdateError(
1467                        ForkchoiceUpdateError::InvalidState,
1468                    ),
1469                ),
1470                409,
1471                "invalid-forkchoice",
1472            ),
1473            (
1474                EngineApiError::ForkChoiceUpdate(
1475                    BeaconForkChoiceUpdateError::ForkchoiceUpdateError(
1476                        ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes,
1477                    ),
1478                ),
1479                422,
1480                "invalid-attributes",
1481            ),
1482        ] {
1483            let response = engine_error_response(error);
1484            assert_eq!(response.status(), status);
1485            assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1486            let body = response.into_body().collect().await.unwrap().to_bytes();
1487            let problem: serde_json::Value = serde_json::from_slice(&body).unwrap();
1488            assert_eq!(problem["type"], format!("/engine-api/errors/{kind}"));
1489        }
1490    }
1491
1492    #[test]
1493    fn bodies_range_query_validation() {
1494        assert_eq!(parse_bodies_range_query("from=1&count=32").unwrap(), (1, 32));
1495        assert_eq!(parse_bodies_range_query("count=2&from=3").unwrap(), (3, 2));
1496        for query in
1497            ["", "from=1", "from=x&count=1", "from=1&count=1&count=2", "from=1&count=1&other=0"]
1498        {
1499            assert_eq!(parse_bodies_range_query(query).unwrap_err().status(), 400);
1500        }
1501        assert_eq!(parse_bodies_range_query("from=1&count=33").unwrap_err().status(), 413);
1502    }
1503
1504    #[test]
1505    fn payload_bodies_are_filtered_at_fork_boundaries() {
1506        use reth_chainspec::{ChainSpecBuilder, ForkCondition};
1507        let chain_spec = ChainSpecBuilder::default()
1508            .chain(1.into())
1509            .genesis(Default::default())
1510            .with_fork(EthereumHardfork::Shanghai, ForkCondition::Timestamp(10))
1511            .with_fork(EthereumHardfork::Cancun, ForkCondition::Timestamp(20))
1512            .with_fork(EthereumHardfork::Prague, ForkCondition::Timestamp(30))
1513            .with_fork(EthereumHardfork::Osaka, ForkCondition::Timestamp(40))
1514            .with_fork(EthereumHardfork::Amsterdam, ForkCondition::Timestamp(50))
1515            .build();
1516        for (fork, start, end) in [
1517            (EngineSszFork::Shanghai, 10, 20),
1518            (EngineSszFork::Cancun, 20, 30),
1519            (EngineSszFork::Prague, 30, 40),
1520            (EngineSszFork::Osaka, 40, 50),
1521        ] {
1522            let body = ExecutionPayloadBodyV1 {
1523                transactions: vec![Bytes::from_static(&[1, 2, 3])],
1524                withdrawals: Some(vec![]),
1525            };
1526            let response = payload_bodies_response(
1527                Ok(vec![
1528                    Some((start - 1, body.clone())),
1529                    Some((start, body.clone())),
1530                    Some((end - 1, body.clone())),
1531                    Some((end, body)),
1532                    None,
1533                ]),
1534                |body| ExecutionPayloadBodyShanghai::try_from(body).ok(),
1535                fork,
1536                &chain_spec,
1537            )
1538            .unwrap();
1539            assert_eq!(
1540                response.entries.iter().map(|entry| entry.available).collect::<Vec<_>>(),
1541                [false, true, true, false, false]
1542            );
1543            for index in [0, 3, 4] {
1544                assert_eq!(response.entries[index].body, ExecutionPayloadBodyShanghai::default());
1545            }
1546        }
1547        assert!(timestamp_matches_fork(&chain_spec, EngineSszFork::Paris, 9));
1548        assert!(!timestamp_matches_fork(&chain_spec, EngineSszFork::Paris, 10));
1549        assert!(!timestamp_matches_fork(&chain_spec, EngineSszFork::Amsterdam, 49));
1550        assert!(timestamp_matches_fork(&chain_spec, EngineSszFork::Amsterdam, 50));
1551    }
1552
1553    #[tokio::test]
1554    async fn payload_bodies_hash_request_limits_and_media_type() {
1555        for content_length in [false, true] {
1556            for count in [32, 33] {
1557                let bytes =
1558                    BodiesByHashRequest { block_hashes: vec![B256::ZERO; count] }.as_ssz_bytes();
1559                let mut request = HttpRequest::builder().header(CONTENT_TYPE, OCTET_STREAM);
1560                if content_length {
1561                    request = request.header("content-length", bytes.len());
1562                }
1563                let result = read_ssz_body(
1564                    request.body(HttpBody::from(bytes)).unwrap(),
1565                    MAX_BODIES_REQUEST_BYTES,
1566                )
1567                .await;
1568                if count == 32 {
1569                    assert_eq!(
1570                        BodiesByHashRequest::from_ssz_bytes(&result.unwrap())
1571                            .unwrap()
1572                            .block_hashes
1573                            .len(),
1574                        32
1575                    );
1576                } else {
1577                    let response = result.unwrap_err();
1578                    assert_eq!(response.status(), STATUS_PAYLOAD_TOO_LARGE);
1579                    assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1580                }
1581            }
1582        }
1583        let response = read_ssz_body(HttpRequest::new(HttpBody::empty()), MAX_BODIES_REQUEST_BYTES)
1584            .await
1585            .unwrap_err();
1586        assert_eq!(response.status(), STATUS_UNSUPPORTED_MEDIA_TYPE);
1587    }
1588
1589    #[tokio::test]
1590    async fn payload_bodies_errors_and_capabilities_match_rest_spec() {
1591        let response =
1592            engine_error_response(EngineApiError::InvalidBodiesRange { start: 0, count: 1 });
1593        assert_eq!(response.status(), 422);
1594        assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1595        let body = response.into_body().collect().await.unwrap().to_bytes();
1596        let problem: serde_json::Value = serde_json::from_slice(&body).unwrap();
1597        assert_eq!(problem["type"], "/engine-api/errors/invalid-body");
1598        let body = handle_capabilities(false).into_body().collect().await.unwrap().to_bytes();
1599        let capabilities: serde_json::Value = serde_json::from_slice(&body).unwrap();
1600        assert_eq!(capabilities["limits"]["bodies.max_count"], 32);
1601    }
1602
1603    #[test]
1604    fn parses_capabilities_endpoint() {
1605        let endpoint = parse_engine_path("/engine/v1/capabilities").unwrap();
1606        assert_eq!(endpoint, EngineSszEndpoint::Capabilities);
1607    }
1608
1609    #[test]
1610    fn parses_identity_endpoint() {
1611        let endpoint = parse_engine_path("/engine/v1/identity").unwrap();
1612        assert_eq!(endpoint, EngineSszEndpoint::Identity);
1613    }
1614
1615    #[test]
1616    fn parses_fork_scoped_payload_endpoint() {
1617        let endpoint = parse_engine_path("/engine/v1/payloads").unwrap();
1618        assert_eq!(endpoint, EngineSszEndpoint::NewPayload);
1619    }
1620
1621    #[test]
1622    fn parses_fork_scoped_get_payload_endpoint() {
1623        let endpoint = parse_engine_path("/engine/v1/payloads/0x0000000000000001").unwrap();
1624        assert_eq!(
1625            endpoint,
1626            EngineSszEndpoint::GetPayload(Ok(PayloadId::new([0, 0, 0, 0, 0, 0, 0, 1])))
1627        );
1628    }
1629
1630    #[test]
1631    fn matches_malformed_get_payload_endpoint() {
1632        assert!(matches!(
1633            parse_engine_path("/engine/v1/payloads/0x01"),
1634            Some(EngineSszEndpoint::GetPayload(Err(_)))
1635        ));
1636    }
1637
1638    #[test]
1639    fn parses_fork_scoped_forkchoice_endpoint() {
1640        let endpoint = parse_engine_path("/engine/v1/forkchoice").unwrap();
1641        assert_eq!(endpoint, EngineSszEndpoint::Forkchoice);
1642    }
1643
1644    #[test]
1645    fn parses_payload_bodies_endpoints() {
1646        assert_eq!(
1647            parse_engine_path("/engine/v1/bodies/hash").unwrap(),
1648            EngineSszEndpoint::PayloadBodiesByHash
1649        );
1650        assert_eq!(
1651            parse_engine_path("/engine/v1/bodies").unwrap(),
1652            EngineSszEndpoint::PayloadBodiesByRange
1653        );
1654    }
1655
1656    #[test]
1657    fn rejects_legacy_version_scoped_endpoint() {
1658        assert!(parse_engine_path("/engine/v4/payloads").is_none());
1659    }
1660
1661    #[test]
1662    fn decodes_blob_request_container() {
1663        let hashes = vec![B256::ZERO, B256::with_last_byte(1)];
1664        let decoded = BlobsV1Request::from_ssz_bytes(
1665            &BlobsV1Request { versioned_hashes: hashes.clone() }.as_ssz_bytes(),
1666        )
1667        .unwrap()
1668        .versioned_hashes;
1669        assert_eq!(decoded, hashes);
1670    }
1671
1672    #[test]
1673    fn decodes_forkchoice_v4_with_custody_columns() {
1674        let forkchoice_state = ForkchoiceState::same_hash(B256::ZERO);
1675        let encoded = ForkchoiceUpdateAmsterdam {
1676            forkchoice_state,
1677            payload_attributes: Optional::none(),
1678            custody_columns: Optional::some(B128::with_last_byte(1)),
1679        }
1680        .as_ssz_bytes();
1681
1682        let (decoded_state, decoded_attrs, custody_columns) =
1683            decode_forkchoice_request(EngineSszFork::Amsterdam, &encoded).unwrap();
1684        assert_eq!(decoded_state, forkchoice_state);
1685        assert!(decoded_attrs.is_none());
1686        assert_eq!(custody_columns, Some(B128::with_last_byte(1)));
1687    }
1688
1689    #[test]
1690    fn decodes_forkchoice_cancun_payload_attributes() {
1691        let forkchoice_state = ForkchoiceState::same_hash(B256::ZERO);
1692        let attrs = PayloadAttributesCancun {
1693            timestamp: 1,
1694            prev_randao: B256::with_last_byte(2),
1695            suggested_fee_recipient: Default::default(),
1696            withdrawals: Vec::new(),
1697            parent_beacon_block_root: B256::with_last_byte(3),
1698        };
1699        let encoded =
1700            ForkchoiceUpdateCancun { forkchoice_state, payload_attributes: Optional::some(attrs) }
1701                .as_ssz_bytes();
1702
1703        let (decoded_state, decoded_attrs, custody_columns) =
1704            decode_forkchoice_request(EngineSszFork::Cancun, &encoded).unwrap();
1705        assert_eq!(decoded_state, forkchoice_state);
1706        let decoded_attrs = decoded_attrs.unwrap();
1707        assert_eq!(decoded_attrs.timestamp, 1);
1708        assert!(decoded_attrs.withdrawals.as_ref().unwrap().is_empty());
1709        assert_eq!(decoded_attrs.parent_beacon_block_root, Some(B256::with_last_byte(3)));
1710        assert!(custody_columns.is_none());
1711    }
1712}