1use 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
64pub 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 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 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 pub async fn engine_api(&self) -> Option<EthEngineApi<Provider, Pool, Validator, ChainSpec>> {
119 self.engine_api.read().await.clone()
120 }
121}
122
123#[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 pub fn new() -> (Self, EngineSszProxyHandle<ChainSpec, Provider, Pool, Validator>) {
134 let handle = EngineSszProxyHandle::new();
135 (Self { handle: handle.clone() }, handle)
136 }
137
138 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#[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
534async 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
594fn decode_blob_hashes_request(body: &[u8]) -> Result<Vec<B256>, &'static str> {
596 Vec::<B256>::from_ssz_bytes(body).map_err(|_| "invalid ssz")
597}
598
599fn 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}