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_containers::{
8    BuiltPayloadAmsterdam, BuiltPayloadCancun, BuiltPayloadOsaka, BuiltPayloadParis,
9    BuiltPayloadPrague, BuiltPayloadShanghai, ExecutionPayloadEnvelopeAmsterdam,
10    ExecutionPayloadEnvelopeCancun, ExecutionPayloadEnvelopeOsaka, ExecutionPayloadEnvelopeParis,
11    ExecutionPayloadEnvelopePrague, ExecutionPayloadEnvelopeShanghai, ForkchoiceUpdateAmsterdam,
12    ForkchoiceUpdateCancun, ForkchoiceUpdateOsaka, ForkchoiceUpdateParis, ForkchoiceUpdatePrague,
13    ForkchoiceUpdateResponse, ForkchoiceUpdateShanghai, Optional,
14    PayloadStatus as EngineSszPayloadStatus,
15};
16use alloy_consensus::{Transaction, TxEnvelope};
17use alloy_eips::{eip2718::Decodable2718, eip7685::RequestsOrHash};
18use alloy_primitives::{Bytes, B128, B256};
19use alloy_rpc_types_engine::{
20    CancunPayloadFields, ExecutionData, ExecutionPayload, ExecutionPayloadFieldV2,
21    ExecutionPayloadSidecar, ForkchoiceState, PayloadAttributes, PayloadId, PraguePayloadFields,
22};
23use http_body_util::BodyExt;
24use jsonrpsee::server::{HttpBody, HttpRequest, HttpResponse};
25use reth_chainspec::EthereumHardforks;
26use reth_engine_primitives::EngineApiValidator;
27use reth_ethereum_engine_primitives::EthEngineTypes;
28use reth_provider::{BalProvider, BlockReader, HeaderProvider, StateProviderFactory};
29use reth_rpc::EngineApi;
30use reth_rpc_engine_api::EngineApiError;
31use reth_transaction_pool::TransactionPool;
32use ssz::Decode;
33use std::{
34    future::Future,
35    pin::Pin,
36    sync::Arc,
37    task::{Context, Poll},
38};
39use tokio::sync::RwLock;
40use tower::{BoxError, Layer, Service};
41
42const OCTET_STREAM: &str = "application/octet-stream";
43const APPLICATION_JSON: &str = "application/json";
44const TEXT_PLAIN: &str = "text/plain";
45const CONTENT_TYPE: &str = "content-type";
46const CACHE_CONTROL: &str = "cache-control";
47const ETH_EXECUTION_VERSION: &str = "eth-execution-version";
48
49const STATUS_OK: u16 = 200;
50const STATUS_BAD_REQUEST: u16 = 400;
51const STATUS_NOT_FOUND: u16 = 404;
52const STATUS_METHOD_NOT_ALLOWED: u16 = 405;
53const STATUS_INTERNAL_SERVER_ERROR: u16 = 500;
54const STATUS_SERVICE_UNAVAILABLE: u16 = 503;
55
56const MAX_BLOB_LIMIT: usize = 128;
57const MAX_PAYLOAD_BYTES: u64 = 64 * 1024 * 1024;
58
59type EthEngineApi<Provider, Pool, Validator, ChainSpec> =
60    EngineApi<Provider, EthEngineTypes, Pool, Validator, ChainSpec>;
61type SharedEthEngineApi<Provider, Pool, Validator, ChainSpec> =
62    Arc<RwLock<Option<EthEngineApi<Provider, Pool, Validator, ChainSpec>>>>;
63
64/// Shared handle used by [`EngineSszProxyLayer`].
65pub struct EngineSszProxyHandle<ChainSpec, Provider = (), Pool = (), Validator = ()> {
66    engine_api: SharedEthEngineApi<Provider, Pool, Validator, ChainSpec>,
67}
68
69impl<C, Provider, Pool, Validator> Clone for EngineSszProxyHandle<C, Provider, Pool, Validator> {
70    fn clone(&self) -> Self {
71        Self { engine_api: self.engine_api.clone() }
72    }
73}
74
75impl<C, Provider, Pool, Validator> std::fmt::Debug
76    for EngineSszProxyHandle<C, Provider, Pool, Validator>
77{
78    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
79        f.debug_struct("EngineSszProxyHandle").finish_non_exhaustive()
80    }
81}
82
83impl<ChainSpec, Provider, Pool, Validator>
84    EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>
85{
86    fn new() -> Self {
87        Self { engine_api: Default::default() }
88    }
89
90    fn with_engine_api(engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>) -> Self {
91        Self { engine_api: Arc::new(RwLock::new(Some(engine_api))) }
92    }
93
94    /// Sets the Engine API implementation used by the proxy.
95    pub async fn set_engine_api(
96        &self,
97        engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
98    ) {
99        *self.engine_api.write().await = Some(engine_api);
100    }
101
102    /// Sets the Engine API implementation during synchronous launch wiring.
103    pub fn set_engine_api_sync(
104        &self,
105        engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
106    ) {
107        *self
108            .engine_api
109            .try_write()
110            .expect("engine api handle should not be locked during launch") = Some(engine_api);
111    }
112}
113
114impl<ChainSpec, Provider, Pool, Validator>
115    EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>
116{
117    /// Returns the Engine API implementation used by the proxy.
118    pub async fn engine_api(&self) -> Option<EthEngineApi<Provider, Pool, Validator, ChainSpec>> {
119        self.engine_api.read().await.clone()
120    }
121}
122
123/// A tower layer that intercepts SSZ Engine API routes under `/engine/v1`.
124#[derive(Clone, Debug)]
125pub struct EngineSszProxyLayer<ChainSpec, Provider = (), Pool = (), Validator = ()> {
126    handle: EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>,
127}
128
129impl<ChainSpec, Provider, Pool, Validator>
130    EngineSszProxyLayer<ChainSpec, Provider, Pool, Validator>
131{
132    /// Creates a new proxy layer and a handle for setting the engine after node launch.
133    pub fn new() -> (Self, EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>) {
134        let handle = EngineSszProxyHandle::new();
135        (Self { handle: handle.clone() }, handle)
136    }
137
138    /// Creates a new proxy layer with an Engine API implementation.
139    pub fn with_engine_api(
140        engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
141    ) -> (Self, EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>) {
142        let handle = EngineSszProxyHandle::with_engine_api(engine_api);
143        (Self { handle: handle.clone() }, handle)
144    }
145}
146
147impl<S, ChainSpec, Provider, Pool, Validator> Layer<S>
148    for EngineSszProxyLayer<ChainSpec, Provider, Pool, Validator>
149{
150    type Service = EngineSszProxyService<S, ChainSpec, Provider, Pool, Validator>;
151
152    fn layer(&self, inner: S) -> Self::Service {
153        EngineSszProxyService { inner, handle: self.handle.clone() }
154    }
155}
156
157/// The service produced by [`EngineSszProxyLayer`].
158#[derive(Clone, Debug)]
159pub struct EngineSszProxyService<S, ChainSpec, Provider = (), Pool = (), Validator = ()> {
160    inner: S,
161    handle: EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>,
162}
163
164impl<S, ChainSpec, Provider, Pool, Validator> Service<HttpRequest>
165    for EngineSszProxyService<S, ChainSpec, Provider, Pool, Validator>
166where
167    S: Service<HttpRequest, Response = HttpResponse, Error = BoxError> + Send + Clone,
168    S::Future: Send + 'static,
169    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
170    Pool: TransactionPool + 'static,
171    Validator: EngineApiValidator<EthEngineTypes>,
172    ChainSpec: EthereumHardforks + Send + Sync + 'static,
173{
174    type Response = HttpResponse;
175    type Error = BoxError;
176    type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
177
178    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
179        self.inner.poll_ready(cx)
180    }
181
182    fn call(&mut self, request: HttpRequest) -> Self::Future {
183        if !request.uri().path().starts_with("/engine/") {
184            let fut = self.inner.call(request);
185            return Box::pin(fut)
186        }
187
188        let handle = self.handle.clone();
189        Box::pin(async move { Ok(handle_engine_ssz_request(handle, request).await) })
190    }
191}
192
193async fn handle_engine_ssz_request<ChainSpec, Provider, Pool, Validator>(
194    handle: EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>,
195    request: HttpRequest,
196) -> HttpResponse
197where
198    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
199    Pool: TransactionPool + 'static,
200    Validator: EngineApiValidator<EthEngineTypes>,
201    ChainSpec: EthereumHardforks + Send + Sync + 'static,
202{
203    let method = request.method().as_str().to_owned();
204    let path = request.uri().path().to_owned();
205    let Some(endpoint) = parse_engine_path(&path) else {
206        return text_response(STATUS_NOT_FOUND, "unknown engine ssz endpoint")
207    };
208
209    match endpoint {
210        EngineSszEndpoint::Capabilities => {
211            if method != "GET" {
212                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
213            }
214            handle_capabilities()
215        }
216        EngineSszEndpoint::Identity => {
217            if method != "GET" {
218                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
219            }
220            let Some(engine_api) = handle.engine_api().await else {
221                return text_response(STATUS_SERVICE_UNAVAILABLE, "engine api unavailable")
222            };
223            handle_identity(engine_api)
224        }
225        EngineSszEndpoint::NewPayload => {
226            if method != "POST" {
227                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
228            }
229            let Some(fork) = request_fork(&request) else {
230                return text_response(STATUS_BAD_REQUEST, "unsupported fork")
231            };
232            let Ok(body) = request.into_body().collect().await.map(|body| body.to_bytes()) else {
233                return text_response(STATUS_BAD_REQUEST, "failed to read request body")
234            };
235            let Some(engine_api) = handle.engine_api().await else {
236                return text_response(STATUS_SERVICE_UNAVAILABLE, "engine api unavailable")
237            };
238            handle_new_payload(engine_api, fork, &body).await
239        }
240        EngineSszEndpoint::GetPayload(payload_id) => {
241            if method != "GET" {
242                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
243            }
244            let Ok(payload_id) = payload_id else {
245                return text_response(STATUS_BAD_REQUEST, "invalid payload id")
246            };
247            let Some(fork) = request_fork(&request) else {
248                return text_response(STATUS_BAD_REQUEST, "unsupported fork")
249            };
250            let Some(engine_api) = handle.engine_api().await else {
251                return text_response(STATUS_SERVICE_UNAVAILABLE, "engine api unavailable")
252            };
253            handle_get_payload(engine_api, fork, payload_id).await
254        }
255        EngineSszEndpoint::Forkchoice => {
256            if method != "POST" {
257                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
258            }
259            let Some(fork) = request_fork(&request) else {
260                return text_response(STATUS_BAD_REQUEST, "unsupported fork")
261            };
262            let Ok(body) = request.into_body().collect().await.map(|body| body.to_bytes()) else {
263                return text_response(STATUS_BAD_REQUEST, "failed to read request body")
264            };
265            let Some(engine_api) = handle.engine_api().await else {
266                return text_response(STATUS_SERVICE_UNAVAILABLE, "engine api unavailable")
267            };
268            handle_forkchoice_updated(engine_api, fork, &body).await
269        }
270        EngineSszEndpoint::Blobs(version) => {
271            if method != "POST" {
272                return text_response(STATUS_METHOD_NOT_ALLOWED, "method not allowed")
273            }
274            let Ok(body) = request.into_body().collect().await.map(|body| body.to_bytes()) else {
275                return text_response(STATUS_BAD_REQUEST, "failed to read request body")
276            };
277            let Some(engine_api) = handle.engine_api().await else {
278                return text_response(STATUS_SERVICE_UNAVAILABLE, "engine api unavailable")
279            };
280            handle_get_blobs(engine_api, version, &body).await
281        }
282    }
283}
284
285fn request_fork(request: &HttpRequest) -> Option<EngineSszFork> {
286    request.headers().get(ETH_EXECUTION_VERSION)?.to_str().ok()?.parse().ok()
287}
288
289fn parse_engine_path(path: &str) -> Option<EngineSszEndpoint> {
290    let mut segments = path.trim_start_matches('/').split('/');
291    match (segments.next(), segments.next(), segments.next(), segments.next(), segments.next()) {
292        (Some("engine"), Some("v1"), Some("capabilities"), None, None) => {
293            Some(EngineSszEndpoint::Capabilities)
294        }
295        (Some("engine"), Some("v1"), Some("identity"), None, None) => {
296            Some(EngineSszEndpoint::Identity)
297        }
298        (Some("engine"), Some("v1"), Some("payloads"), None, None) => {
299            Some(EngineSszEndpoint::NewPayload)
300        }
301        (Some("engine"), Some("v1"), Some("payloads"), Some(payload_id), None) => {
302            let payload_id = payload_id.parse::<PayloadId>();
303            Some(EngineSszEndpoint::GetPayload(payload_id))
304        }
305        (Some("engine"), Some("v1"), Some("forkchoice"), None, None) => {
306            Some(EngineSszEndpoint::Forkchoice)
307        }
308        (Some("engine"), Some("v1"), Some("blobs"), version, None) => {
309            Some(EngineSszEndpoint::Blobs(parse_method_version(version?)?))
310        }
311        _ => None,
312    }
313}
314
315#[derive(Clone, Copy, Debug, Eq, PartialEq)]
316enum EngineSszEndpoint {
317    Capabilities,
318    Identity,
319    NewPayload,
320    GetPayload(Result<PayloadId, <PayloadId as std::str::FromStr>::Err>),
321    Forkchoice,
322    Blobs(u8),
323}
324
325#[derive(Clone, Copy, Debug, Eq, PartialEq)]
326enum EngineSszFork {
327    Paris,
328    Shanghai,
329    Cancun,
330    Prague,
331    Osaka,
332    Amsterdam,
333}
334
335impl EngineSszFork {
336    const fn payloads_version(self) -> u8 {
337        match self {
338            Self::Paris => 1,
339            Self::Shanghai => 2,
340            Self::Cancun => 3,
341            Self::Prague | Self::Osaka => 4,
342            Self::Amsterdam => 5,
343        }
344    }
345
346    const fn forkchoice_version(self) -> u8 {
347        match self {
348            Self::Paris => 1,
349            Self::Shanghai => 2,
350            Self::Cancun | Self::Prague | Self::Osaka => 3,
351            Self::Amsterdam => 4,
352        }
353    }
354}
355
356impl std::str::FromStr for EngineSszFork {
357    type Err = ();
358
359    fn from_str(value: &str) -> Result<Self, Self::Err> {
360        match value {
361            "paris" => Ok(Self::Paris),
362            "shanghai" => Ok(Self::Shanghai),
363            "cancun" => Ok(Self::Cancun),
364            "prague" => Ok(Self::Prague),
365            "osaka" => Ok(Self::Osaka),
366            "amsterdam" => Ok(Self::Amsterdam),
367            _ => Err(()),
368        }
369    }
370}
371
372fn parse_method_version(version: &str) -> Option<u8> {
373    version.strip_prefix('v')?.parse().ok().filter(|version| (1..=4).contains(version))
374}
375
376fn handle_capabilities() -> HttpResponse {
377    json_response(serde_json::json!({
378        "supported_forks": ["paris", "shanghai", "cancun", "prague", "osaka", "amsterdam"],
379        "fork_scoped_endpoints": ["payloads", "forkchoice", "bodies"],
380        "independently_versioned": {
381            "blobs": ["v1", "v2", "v3", "v4"],
382        },
383        "unscoped_endpoints": ["capabilities", "identity"],
384        "limits": {
385            "bodies.max_count": 128,
386            "blobs.max_versioned_hashes": MAX_BLOB_LIMIT,
387            "payload.max_bytes": MAX_PAYLOAD_BYTES,
388        },
389    }))
390}
391
392fn handle_identity<Provider, Pool, Validator, ChainSpec>(
393    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
394) -> HttpResponse
395where
396    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
397    Pool: TransactionPool + 'static,
398    Validator: EngineApiValidator<EthEngineTypes>,
399    ChainSpec: EthereumHardforks + Send + Sync + 'static,
400{
401    json_response(vec![engine_api.client_version().clone()])
402}
403
404async fn handle_get_payload<Provider, Pool, Validator, ChainSpec>(
405    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
406    fork: EngineSszFork,
407    payload_id: PayloadId,
408) -> HttpResponse
409where
410    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
411    Pool: TransactionPool + 'static,
412    Validator: EngineApiValidator<EthEngineTypes>,
413    ChainSpec: EthereumHardforks + Send + Sync + 'static,
414{
415    match fork {
416        EngineSszFork::Paris => match engine_api.get_payload_v2_metered(payload_id).await {
417            Ok(payload) => {
418                let block_value = payload.block_value;
419                match payload.execution_payload {
420                    ExecutionPayloadFieldV2::V1(payload) => {
421                        get_payload_response(BuiltPayloadParis { payload, block_value })
422                    }
423                    ExecutionPayloadFieldV2::V2(_) => {
424                        text_response(STATUS_BAD_REQUEST, "unsupported fork")
425                    }
426                }
427            }
428            Err(err) => get_payload_error_response(err),
429        },
430        EngineSszFork::Shanghai => match engine_api.get_payload_v2_metered(payload_id).await {
431            Ok(payload) => match BuiltPayloadShanghai::try_from(payload) {
432                Ok(payload) => get_payload_response(payload),
433                Err(err) => text_response(STATUS_BAD_REQUEST, err.to_string()),
434            },
435            Err(err) => get_payload_error_response(err),
436        },
437        EngineSszFork::Cancun => match engine_api.get_payload_v3_metered(payload_id).await {
438            Ok(payload) => get_payload_response(BuiltPayloadCancun::from(payload)),
439            Err(err) => get_payload_error_response(err),
440        },
441        EngineSszFork::Prague => match engine_api.get_payload_v4_metered(payload_id).await {
442            Ok(payload) => get_payload_response(BuiltPayloadPrague::from(payload)),
443            Err(err) => get_payload_error_response(err),
444        },
445        EngineSszFork::Osaka => match engine_api.get_payload_v5_metered(payload_id).await {
446            Ok(payload) => get_payload_response(BuiltPayloadOsaka::from(payload)),
447            Err(err) => get_payload_error_response(err),
448        },
449        EngineSszFork::Amsterdam => match engine_api.get_payload_v6_metered(payload_id).await {
450            Ok(payload) => get_payload_response(BuiltPayloadAmsterdam::from(payload)),
451            Err(err) => get_payload_error_response(err),
452        },
453    }
454}
455
456fn get_payload_error_response(err: EngineApiError) -> HttpResponse {
457    let status = match &err {
458        EngineApiError::UnknownPayload => STATUS_NOT_FOUND,
459        EngineApiError::EngineObjectValidationError(
460            reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
461        ) => STATUS_BAD_REQUEST,
462        _ => STATUS_INTERNAL_SERVER_ERROR,
463    };
464    text_response(status, err.to_string())
465}
466
467async fn handle_new_payload<Provider, Pool, Validator, ChainSpec>(
468    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
469    fork: EngineSszFork,
470    body: &[u8],
471) -> HttpResponse
472where
473    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
474    Pool: TransactionPool + 'static,
475    Validator: EngineApiValidator<EthEngineTypes>,
476    ChainSpec: EthereumHardforks + Send + Sync + 'static,
477{
478    let payload = match decode_new_payload_request(fork, body) {
479        Ok(payload) => payload,
480        Err(err) => return text_response(STATUS_BAD_REQUEST, err),
481    };
482
483    let response = match fork.payloads_version() {
484        1 => engine_api.new_payload_v1(payload).await,
485        2 => engine_api.new_payload_v2(payload).await,
486        3 => engine_api.new_payload_v3(payload).await,
487        4 => engine_api.new_payload_v4(payload).await,
488        5 => engine_api.new_payload_v5(payload).await,
489        _ => return text_response(STATUS_BAD_REQUEST, "unsupported payload endpoint version"),
490    };
491
492    match response {
493        Ok(status) => match EngineSszPayloadStatus::try_from(status) {
494            Ok(status) => ssz_response(status),
495            Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
496        },
497        Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
498    }
499}
500
501async fn handle_forkchoice_updated<Provider, Pool, Validator, ChainSpec>(
502    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
503    fork: EngineSszFork,
504    body: &[u8],
505) -> HttpResponse
506where
507    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
508    Pool: TransactionPool + 'static,
509    Validator: EngineApiValidator<EthEngineTypes>,
510    ChainSpec: EthereumHardforks + Send + Sync + 'static,
511{
512    let (state, attrs, custody_columns) = match decode_forkchoice_request(fork, body) {
513        Ok(request) => request,
514        Err(err) => return text_response(STATUS_BAD_REQUEST, err),
515    };
516
517    let response = match fork.forkchoice_version() {
518        1 => engine_api.fork_choice_updated_v1_metered(state, attrs).await,
519        2 => engine_api.fork_choice_updated_v2_metered(state, attrs).await,
520        3 => engine_api.fork_choice_updated_v3_metered(state, attrs).await,
521        4 => engine_api.fork_choice_updated_v4_metered(state, attrs, custody_columns).await,
522        _ => return text_response(STATUS_BAD_REQUEST, "unsupported forkchoice endpoint version"),
523    };
524
525    match response {
526        Ok(updated) => match ForkchoiceUpdateResponse::try_from(updated) {
527            Ok(updated) => ssz_response(updated),
528            Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
529        },
530        Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
531    }
532}
533
534/// Handles SSZ `engine_getBlobsV*` requests with the node's blob store.
535async fn handle_get_blobs<ChainSpec, Provider, Pool, Validator>(
536    engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
537    version: u8,
538    body: &[u8],
539) -> HttpResponse
540where
541    Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
542    Pool: TransactionPool + 'static,
543    Validator: EngineApiValidator<EthEngineTypes>,
544    ChainSpec: EthereumHardforks + Send + Sync + 'static,
545{
546    match version {
547        1 => {
548            let hashes = match decode_blob_hashes_request(body) {
549                Ok(hashes) => hashes,
550                Err(err) => return text_response(STATUS_BAD_REQUEST, err),
551            };
552            match engine_api.get_blobs_v1_metered(hashes) {
553                Ok(response) => ssz_response(response),
554                Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
555            }
556        }
557        2 => {
558            let hashes = match decode_blob_hashes_request(body) {
559                Ok(hashes) => hashes,
560                Err(err) => return text_response(STATUS_BAD_REQUEST, err),
561            };
562            match engine_api.get_blobs_v2_metered(hashes) {
563                Ok(Some(response)) => ssz_response(response),
564                Ok(None) => no_content_response(),
565                Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
566            }
567        }
568        3 => {
569            let hashes = match decode_blob_hashes_request(body) {
570                Ok(hashes) => hashes,
571                Err(err) => return text_response(STATUS_BAD_REQUEST, err),
572            };
573            match engine_api.get_blobs_v3_metered(hashes) {
574                Ok(Some(response)) => ssz_response(response),
575                Ok(None) => no_content_response(),
576                Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
577            }
578        }
579        4 => {
580            let (hashes, indices_bitarray) = match decode_blob_cells_request(body) {
581                Ok(request) => request,
582                Err(err) => return text_response(STATUS_BAD_REQUEST, err),
583            };
584            match engine_api.get_blobs_v4_metered(hashes, indices_bitarray) {
585                Ok(Some(response)) => ssz_response(response),
586                Ok(None) => no_content_response(),
587                Err(err) => text_response(STATUS_INTERNAL_SERVER_ERROR, err.to_string()),
588            }
589        }
590        _ => text_response(STATUS_NOT_FOUND, "unsupported blobs endpoint version"),
591    }
592}
593
594/// Decodes the common getBlobs request container with only versioned hashes.
595fn decode_blob_hashes_request(body: &[u8]) -> Result<Vec<B256>, &'static str> {
596    Vec::<B256>::from_ssz_bytes(body).map_err(|_| "invalid ssz")
597}
598
599/// Decodes the Amsterdam getBlobs request container with hashes and a cell index mask.
600fn decode_blob_cells_request(body: &[u8]) -> Result<(Vec<B256>, B128), &'static str> {
601    <(Vec<B256>, B128) as ssz::Decode>::from_ssz_bytes(body).map_err(|_| "invalid ssz")
602}
603
604fn decode_new_payload_request(
605    fork: EngineSszFork,
606    body: &[u8],
607) -> Result<ExecutionData, &'static str> {
608    match fork {
609        EngineSszFork::Paris => {
610            let ExecutionPayloadEnvelopeParis { payload: execution_payload } =
611                ExecutionPayloadEnvelopeParis::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
612            Ok(ExecutionData::new(execution_payload.into(), ExecutionPayloadSidecar::none()))
613        }
614        EngineSszFork::Shanghai => {
615            let ExecutionPayloadEnvelopeShanghai { payload: execution_payload } =
616                ExecutionPayloadEnvelopeShanghai::from_ssz_bytes(body)
617                    .map_err(|_| "invalid ssz")?;
618            Ok(ExecutionData::new(execution_payload.into(), ExecutionPayloadSidecar::none()))
619        }
620        EngineSszFork::Cancun => {
621            let ExecutionPayloadEnvelopeCancun {
622                payload: execution_payload,
623                parent_beacon_block_root,
624            } = ExecutionPayloadEnvelopeCancun::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
625            let versioned_hashes = calculate_versioned_hashes(
626                &execution_payload.payload_inner.payload_inner.transactions,
627            )?;
628            let sidecar = ExecutionPayloadSidecar::v3(CancunPayloadFields {
629                parent_beacon_block_root,
630                versioned_hashes,
631            });
632            Ok(ExecutionData::new(execution_payload.into(), sidecar))
633        }
634        EngineSszFork::Prague => {
635            let ExecutionPayloadEnvelopePrague {
636                payload: execution_payload,
637                parent_beacon_block_root,
638                execution_requests,
639            } = ExecutionPayloadEnvelopePrague::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
640            let versioned_hashes = calculate_versioned_hashes(
641                &execution_payload.payload_inner.payload_inner.transactions,
642            )?;
643            let sidecar = ExecutionPayloadSidecar::v4(
644                CancunPayloadFields { parent_beacon_block_root, versioned_hashes },
645                PraguePayloadFields::new(RequestsOrHash::Requests(execution_requests)),
646            );
647            Ok(ExecutionData::new(execution_payload.into(), sidecar))
648        }
649        EngineSszFork::Osaka => {
650            let ExecutionPayloadEnvelopeOsaka {
651                payload: execution_payload,
652                parent_beacon_block_root,
653                execution_requests,
654            } = ExecutionPayloadEnvelopeOsaka::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
655            let versioned_hashes = calculate_versioned_hashes(
656                &execution_payload.payload_inner.payload_inner.transactions,
657            )?;
658            let sidecar = ExecutionPayloadSidecar::v4(
659                CancunPayloadFields { parent_beacon_block_root, versioned_hashes },
660                PraguePayloadFields::new(RequestsOrHash::Requests(execution_requests)),
661            );
662            Ok(ExecutionData::new(execution_payload.into(), sidecar))
663        }
664        EngineSszFork::Amsterdam => {
665            let ExecutionPayloadEnvelopeAmsterdam {
666                payload: execution_payload,
667                parent_beacon_block_root,
668                execution_requests,
669            } = ExecutionPayloadEnvelopeAmsterdam::from_ssz_bytes(body)
670                .map_err(|_| "invalid ssz")?;
671            let versioned_hashes = calculate_versioned_hashes(
672                &execution_payload.payload_inner.payload_inner.payload_inner.transactions,
673            )?;
674            let sidecar = ExecutionPayloadSidecar::v4(
675                CancunPayloadFields { parent_beacon_block_root, versioned_hashes },
676                PraguePayloadFields::new(RequestsOrHash::Requests(execution_requests)),
677            );
678            Ok(ExecutionData::new(ExecutionPayload::V4(execution_payload), sidecar))
679        }
680    }
681}
682
683fn calculate_versioned_hashes(transactions: &[Bytes]) -> Result<Vec<B256>, &'static str> {
684    let mut versioned_hashes = Vec::new();
685    for transaction in transactions {
686        let transaction =
687            TxEnvelope::decode_2718_exact(transaction.as_ref()).map_err(|_| "invalid tx")?;
688        if let Some(hashes) = transaction.blob_versioned_hashes() {
689            versioned_hashes.extend_from_slice(hashes);
690        }
691    }
692
693    Ok(versioned_hashes)
694}
695
696fn decode_forkchoice_request(
697    fork: EngineSszFork,
698    body: &[u8],
699) -> Result<(ForkchoiceState, Option<PayloadAttributes>, Option<B128>), &'static str> {
700    match fork {
701        EngineSszFork::Paris => {
702            let ForkchoiceUpdateParis { forkchoice_state, payload_attributes } =
703                ForkchoiceUpdateParis::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
704            Ok((forkchoice_state, optional_attrs(payload_attributes), None))
705        }
706        EngineSszFork::Shanghai => {
707            let ForkchoiceUpdateShanghai { forkchoice_state, payload_attributes } =
708                ForkchoiceUpdateShanghai::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
709            Ok((forkchoice_state, optional_attrs(payload_attributes), None))
710        }
711        EngineSszFork::Cancun => {
712            let ForkchoiceUpdateCancun { forkchoice_state, payload_attributes } =
713                ForkchoiceUpdateCancun::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
714            Ok((forkchoice_state, optional_attrs(payload_attributes), None))
715        }
716        EngineSszFork::Prague => {
717            let ForkchoiceUpdatePrague { forkchoice_state, payload_attributes } =
718                ForkchoiceUpdatePrague::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
719            Ok((forkchoice_state, optional_attrs(payload_attributes), None))
720        }
721        EngineSszFork::Osaka => {
722            let ForkchoiceUpdateOsaka { forkchoice_state, payload_attributes } =
723                ForkchoiceUpdateOsaka::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
724            Ok((forkchoice_state, optional_attrs(payload_attributes), None))
725        }
726        EngineSszFork::Amsterdam => {
727            let ForkchoiceUpdateAmsterdam { forkchoice_state, payload_attributes, custody_columns } =
728                ForkchoiceUpdateAmsterdam::from_ssz_bytes(body).map_err(|_| "invalid ssz")?;
729            Ok((
730                forkchoice_state,
731                optional_attrs(payload_attributes),
732                custody_columns.into_option(),
733            ))
734        }
735    }
736}
737
738fn optional_attrs<T>(attrs: Optional<T>) -> Option<PayloadAttributes>
739where
740    T: Into<PayloadAttributes>,
741{
742    attrs.into_option().map(Into::into)
743}
744
745fn ssz_response<T: ssz::Encode>(value: T) -> HttpResponse {
746    HttpResponse::builder()
747        .status(STATUS_OK)
748        .header(CONTENT_TYPE, OCTET_STREAM)
749        .body(HttpBody::from(value.as_ssz_bytes()))
750        .expect("valid response")
751}
752
753fn get_payload_response<T: ssz::Encode>(value: T) -> HttpResponse {
754    HttpResponse::builder()
755        .status(STATUS_OK)
756        .header(CONTENT_TYPE, OCTET_STREAM)
757        .header(CACHE_CONTROL, "no-store")
758        .body(HttpBody::from(value.as_ssz_bytes()))
759        .expect("valid response")
760}
761
762fn json_response<T: serde::Serialize>(value: T) -> HttpResponse {
763    let Ok(body) = serde_json::to_string(&value) else {
764        return text_response(STATUS_INTERNAL_SERVER_ERROR, "failed to encode json")
765    };
766
767    HttpResponse::builder()
768        .status(STATUS_OK)
769        .header(CONTENT_TYPE, APPLICATION_JSON)
770        .body(HttpBody::from(body))
771        .expect("valid response")
772}
773
774fn no_content_response() -> HttpResponse {
775    HttpResponse::builder().status(204).body(HttpBody::empty()).expect("valid response")
776}
777
778fn text_response(status: u16, body: impl Into<String>) -> HttpResponse {
779    HttpResponse::builder()
780        .status(status)
781        .header(CONTENT_TYPE, TEXT_PLAIN)
782        .body(HttpBody::from(body.into()))
783        .expect("valid response")
784}
785
786#[cfg(test)]
787mod tests {
788    use super::*;
789    use ssz::Encode;
790
791    #[test]
792    fn parses_capabilities_endpoint() {
793        let endpoint = parse_engine_path("/engine/v1/capabilities").unwrap();
794        assert_eq!(endpoint, EngineSszEndpoint::Capabilities);
795    }
796
797    #[test]
798    fn parses_identity_endpoint() {
799        let endpoint = parse_engine_path("/engine/v1/identity").unwrap();
800        assert_eq!(endpoint, EngineSszEndpoint::Identity);
801    }
802
803    #[test]
804    fn parses_fork_scoped_payload_endpoint() {
805        let endpoint = parse_engine_path("/engine/v1/payloads").unwrap();
806        assert_eq!(endpoint, EngineSszEndpoint::NewPayload);
807    }
808
809    #[test]
810    fn parses_fork_scoped_get_payload_endpoint() {
811        let endpoint = parse_engine_path("/engine/v1/payloads/0x0000000000000001").unwrap();
812        assert_eq!(
813            endpoint,
814            EngineSszEndpoint::GetPayload(Ok(PayloadId::new([0, 0, 0, 0, 0, 0, 0, 1])))
815        );
816    }
817
818    #[test]
819    fn matches_malformed_get_payload_endpoint() {
820        assert!(matches!(
821            parse_engine_path("/engine/v1/payloads/0x01"),
822            Some(EngineSszEndpoint::GetPayload(Err(_)))
823        ));
824    }
825
826    #[test]
827    fn parses_fork_scoped_forkchoice_endpoint() {
828        let endpoint = parse_engine_path("/engine/v1/forkchoice").unwrap();
829        assert_eq!(endpoint, EngineSszEndpoint::Forkchoice);
830    }
831
832    #[test]
833    fn rejects_legacy_version_scoped_endpoint() {
834        assert!(parse_engine_path("/engine/v4/payloads").is_none());
835    }
836
837    #[test]
838    fn decodes_top_level_blob_hashes_request() {
839        let hashes = vec![B256::ZERO, B256::with_last_byte(1)];
840        let decoded = decode_blob_hashes_request(&hashes.as_ssz_bytes()).unwrap();
841        assert_eq!(decoded, hashes);
842    }
843
844    #[test]
845    fn decodes_forkchoice_v4_with_custody_columns() {
846        let forkchoice_state = ForkchoiceState {
847            head_block_hash: B256::ZERO,
848            safe_block_hash: B256::ZERO,
849            finalized_block_hash: B256::ZERO,
850        };
851        let encoded = ForkchoiceUpdateAmsterdam {
852            forkchoice_state,
853            payload_attributes: Optional::none(),
854            custody_columns: Optional::some(B128::with_last_byte(1)),
855        }
856        .as_ssz_bytes();
857
858        let (decoded_state, decoded_attrs, custody_columns) =
859            decode_forkchoice_request(EngineSszFork::Amsterdam, &encoded).unwrap();
860        assert_eq!(decoded_state, forkchoice_state);
861        assert!(decoded_attrs.is_none());
862        assert_eq!(custody_columns, Some(B128::with_last_byte(1)));
863    }
864
865    #[test]
866    fn decodes_forkchoice_cancun_payload_attributes() {
867        let forkchoice_state = ForkchoiceState {
868            head_block_hash: B256::ZERO,
869            safe_block_hash: B256::ZERO,
870            finalized_block_hash: B256::ZERO,
871        };
872        let attrs = crate::engine_ssz_containers::PayloadAttributesCancun {
873            timestamp: 1,
874            prev_randao: B256::with_last_byte(2),
875            suggested_fee_recipient: Default::default(),
876            withdrawals: Vec::new(),
877            parent_beacon_block_root: B256::with_last_byte(3),
878        };
879        let encoded =
880            ForkchoiceUpdateCancun { forkchoice_state, payload_attributes: Optional::some(attrs) }
881                .as_ssz_bytes();
882
883        let (decoded_state, decoded_attrs, custody_columns) =
884            decode_forkchoice_request(EngineSszFork::Cancun, &encoded).unwrap();
885        assert_eq!(decoded_state, forkchoice_state);
886        let decoded_attrs = decoded_attrs.unwrap();
887        assert_eq!(decoded_attrs.timestamp, 1);
888        assert!(decoded_attrs.withdrawals.as_ref().unwrap().is_empty());
889        assert_eq!(decoded_attrs.parent_beacon_block_root, Some(B256::with_last_byte(3)));
890        assert!(custody_columns.is_none());
891    }
892}