1use 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>;
75pub 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 pub async fn witness_enabled(&self) -> bool {
115 self.state.read().await.witness_enabled
116 }
117
118 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 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 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 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 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 pub async fn engine_api(&self) -> Option<Api> {
159 self.state.read().await.engine_api.clone()
160 }
161}
162
163#[derive(Clone, Debug)]
165pub struct EngineSszProxyLayer<Api = ()> {
166 handle: EngineSszProxyHandle<Api>,
167}
168
169impl<Api> EngineSszProxyLayer<Api> {
170 pub fn new() -> (Self, EngineSszProxyHandle<Api>) {
172 let handle = EngineSszProxyHandle::new();
173 (Self { handle: handle.clone() }, handle)
174 }
175
176 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#[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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
224pub enum EngineSszFork {
225 Paris,
227 Shanghai,
229 Cancun,
231 Prague,
233 Osaka,
235 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
276pub trait EngineSszApi: Clone + Send + Sync + 'static {
281 fn capabilities(&self, _witness_enabled: bool) -> HttpResponse {
286 problem_response(STATUS_NOT_FOUND, "method-not-found", None)
287 }
288
289 fn supports_witness(&self) -> bool {
291 false
292 }
293
294 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 fn identity(&self) -> HttpResponse {
305 problem_response(STATUS_NOT_FOUND, "method-not-found", None)
306 }
307
308 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 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 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 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 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 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 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
918fn 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#[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 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}