1use crate::{
8 engine_ssz_containers::{
9 BlobsV1Request, BlobsV1Response, BlobsV2Response, BlobsV3Response, BlobsV4Request,
10 BlobsV4Response, BodiesByHashRequest, BodiesResponse, BuiltPayloadAmsterdam,
11 BuiltPayloadOsaka, BuiltPayloadParis, BuiltPayloadPrague, BuiltPayloadShanghai,
12 ExecutionPayloadBodyAmsterdam, ExecutionPayloadBodyParis, ExecutionPayloadBodyShanghai,
13 ExecutionPayloadEnvelopeAmsterdam, ExecutionPayloadEnvelopeCancun,
14 ExecutionPayloadEnvelopeOsaka, ExecutionPayloadEnvelopeParis,
15 ExecutionPayloadEnvelopePrague, ExecutionPayloadEnvelopeShanghai,
16 ForkchoiceUpdateAmsterdam, ForkchoiceUpdateCancun, ForkchoiceUpdateOsaka,
17 ForkchoiceUpdateParis, ForkchoiceUpdatePrague, ForkchoiceUpdateResponse,
18 ForkchoiceUpdateShanghai, Optional, PayloadStatus as EngineSszPayloadStatus,
19 PayloadStatusWithWitness,
20 },
21 engine_ssz_witness::{EngineSszWitness, EngineSszWitnessError},
22};
23use alloy_consensus::{Transaction, TxEnvelope};
24use alloy_eips::{eip2718::Decodable2718, eip7685::Requests};
25use alloy_primitives::{Bytes, B128, B256};
26use alloy_rpc_types_engine::{
27 CancunPayloadFields, ExecutionData, ExecutionPayload, ExecutionPayloadBodyV1,
28 ExecutionPayloadFieldV2, ExecutionPayloadSidecar, ForkchoiceState, PayloadAttributes,
29 PayloadId, PayloadStatusEnum, 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_LIMIT: usize = crate::engine_ssz_containers::MAX_BLOBS_REQUEST;
69const MAX_BLOB_REQUEST_BYTES: usize = 4 + 16 + MAX_BLOB_LIMIT * 32;
70const MAX_BODIES_REQUEST: usize = crate::engine_ssz_containers::MAX_BODIES_REQUEST;
71const MAX_BODIES_REQUEST_BYTES: usize = 4 + MAX_BODIES_REQUEST * 32;
72const MAX_PAYLOAD_BYTES: usize = 64 * 1024 * 1024;
73const PROBLEM_JSON: &str = "application/problem+json";
74
75type EthEngineApi<Provider, Pool, Validator, ChainSpec> =
76 EngineApi<Provider, EthEngineTypes, Pool, Validator, ChainSpec>;
77pub struct EngineSszProxyHandle<Api = ()> {
79 state: Arc<RwLock<EngineSszState<Api>>>,
80}
81
82impl<Api> Clone for EngineSszProxyHandle<Api> {
83 fn clone(&self) -> Self {
84 Self { state: self.state.clone() }
85 }
86}
87
88impl<Api> std::fmt::Debug for EngineSszProxyHandle<Api> {
89 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90 f.debug_struct("EngineSszProxyHandle").finish_non_exhaustive()
91 }
92}
93
94impl<Api> EngineSszProxyHandle<Api> {
95 fn new() -> Self {
96 Self {
97 state: Arc::new(RwLock::new(EngineSszState {
98 engine_api: None,
99 witness_handler: None,
100 witness_enabled: false,
101 })),
102 }
103 }
104
105 fn with_engine_api(engine_api: Api) -> Self {
106 Self {
107 state: Arc::new(RwLock::new(EngineSszState {
108 engine_api: Some(engine_api),
109 witness_handler: None,
110 witness_enabled: false,
111 })),
112 }
113 }
114
115 pub async fn witness_enabled(&self) -> bool {
117 self.state.read().await.witness_enabled
118 }
119
120 pub async fn witness_handler(&self) -> Option<Arc<dyn EngineSszWitness>> {
122 self.state.read().await.witness_handler.clone()
123 }
124}
125
126impl<Api: EngineSszApi> EngineSszProxyHandle<Api> {
127 pub async fn set_engine_api(&self, engine_api: Api) {
129 let mut state = self.state.write().await;
130 state.engine_api = Some(engine_api);
131 state.update_witness_support();
132 }
133
134 pub fn set_engine_api_sync(&self, engine_api: Api) {
136 let mut state =
137 self.state.try_write().expect("engine api handle should not be locked during launch");
138 state.engine_api = Some(engine_api);
139 state.update_witness_support();
140 }
141
142 pub async fn set_witness_handler(&self, witness_handler: Arc<dyn EngineSszWitness>) {
144 let mut state = self.state.write().await;
145 state.witness_handler = Some(witness_handler);
146 state.update_witness_support();
147 }
148
149 pub fn set_witness_handler_sync(&self, witness_handler: Arc<dyn EngineSszWitness>) {
151 let mut state =
152 self.state.try_write().expect("witness handle should not be locked during launch");
153 state.witness_handler = Some(witness_handler);
154 state.update_witness_support();
155 }
156}
157
158impl<Api: Clone> EngineSszProxyHandle<Api> {
159 pub async fn engine_api(&self) -> Option<Api> {
161 self.state.read().await.engine_api.clone()
162 }
163}
164
165#[derive(Clone, Debug)]
167pub struct EngineSszProxyLayer<Api = ()> {
168 handle: EngineSszProxyHandle<Api>,
169}
170
171impl<Api> EngineSszProxyLayer<Api> {
172 pub fn new() -> (Self, EngineSszProxyHandle<Api>) {
174 let handle = EngineSszProxyHandle::new();
175 (Self { handle: handle.clone() }, handle)
176 }
177
178 pub fn with_engine_api(engine_api: Api) -> (Self, EngineSszProxyHandle<Api>) {
180 let handle = EngineSszProxyHandle::with_engine_api(engine_api);
181 (Self { handle: handle.clone() }, handle)
182 }
183}
184
185impl<S, Api> Layer<S> for EngineSszProxyLayer<Api> {
186 type Service = EngineSszProxyService<S, Api>;
187
188 fn layer(&self, inner: S) -> Self::Service {
189 EngineSszProxyService { inner, handle: self.handle.clone() }
190 }
191}
192
193#[derive(Clone, Debug)]
195pub struct EngineSszProxyService<S, Api = ()> {
196 inner: S,
197 handle: EngineSszProxyHandle<Api>,
198}
199
200impl<S, Api> Service<HttpRequest> for EngineSszProxyService<S, Api>
201where
202 S: Service<HttpRequest, Response = HttpResponse, Error = BoxError> + Send + Clone,
203 S::Future: Send,
204 Api: EngineSszApi,
205{
206 type Response = HttpResponse;
207 type Error = BoxError;
208 type Future = Either<S::Future, BoxFuture<'static, Result<HttpResponse, BoxError>>>;
209
210 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
211 self.inner.poll_ready(cx)
212 }
213
214 fn call(&mut self, request: HttpRequest) -> Self::Future {
215 if !request.uri().path().starts_with("/engine/") {
216 return Either::Left(self.inner.call(request))
217 }
218
219 let handle = self.handle.clone();
220 Either::Right(Box::pin(async move { Ok(handle_engine_ssz_request(handle, request).await) }))
221 }
222}
223
224#[derive(Clone, Copy, Debug, Eq, PartialEq)]
226pub enum EngineSszFork {
227 Paris,
229 Shanghai,
231 Cancun,
233 Prague,
235 Osaka,
237 Amsterdam,
239}
240
241impl EngineSszFork {
242 const fn payloads_version(self) -> u8 {
243 match self {
244 Self::Paris => 1,
245 Self::Shanghai => 2,
246 Self::Cancun => 3,
247 Self::Prague | Self::Osaka => 4,
248 Self::Amsterdam => 5,
249 }
250 }
251
252 const fn forkchoice_version(self) -> u8 {
253 match self {
254 Self::Paris => 1,
255 Self::Shanghai => 2,
256 Self::Cancun | Self::Prague | Self::Osaka => 3,
257 Self::Amsterdam => 4,
258 }
259 }
260}
261
262impl std::str::FromStr for EngineSszFork {
263 type Err = ();
264
265 fn from_str(value: &str) -> Result<Self, Self::Err> {
266 match value {
267 "paris" => Ok(Self::Paris),
268 "shanghai" => Ok(Self::Shanghai),
269 "cancun" => Ok(Self::Cancun),
270 "prague" => Ok(Self::Prague),
271 "osaka" => Ok(Self::Osaka),
272 "amsterdam" => Ok(Self::Amsterdam),
273 _ => Err(()),
274 }
275 }
276}
277
278pub trait EngineSszApi: Clone + Send + Sync + 'static {
283 fn capabilities(&self, _witness_enabled: bool) -> HttpResponse {
288 problem_response(STATUS_NOT_FOUND, "method-not-found", None)
289 }
290
291 fn supports_witness(&self) -> bool {
293 false
294 }
295
296 fn new_payload_with_witness(
298 &self,
299 _body: Bytes,
300 _witness_handler: Arc<dyn EngineSszWitness>,
301 ) -> impl Future<Output = HttpResponse> + Send {
302 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
303 }
304
305 fn identity(&self) -> HttpResponse {
307 problem_response(STATUS_NOT_FOUND, "method-not-found", None)
308 }
309
310 fn new_payload(
312 &self,
313 _fork: EngineSszFork,
314 _body: Bytes,
315 ) -> impl Future<Output = HttpResponse> + Send {
316 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
317 }
318
319 fn get_payload(
321 &self,
322 _fork: EngineSszFork,
323 _payload_id: PayloadId,
324 ) -> impl Future<Output = HttpResponse> + Send {
325 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
326 }
327
328 fn forkchoice_updated(
330 &self,
331 _fork: EngineSszFork,
332 _body: Bytes,
333 ) -> impl Future<Output = HttpResponse> + Send {
334 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
335 }
336
337 fn get_payload_bodies_by_hash(
339 &self,
340 _fork: EngineSszFork,
341 _hashes: Vec<B256>,
342 ) -> impl Future<Output = HttpResponse> + Send {
343 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
344 }
345
346 fn get_payload_bodies_by_range(
348 &self,
349 _fork: EngineSszFork,
350 _start: u64,
351 _count: u64,
352 ) -> impl Future<Output = HttpResponse> + Send {
353 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
354 }
355
356 fn get_blobs(&self, _version: u8, _body: Bytes) -> impl Future<Output = HttpResponse> + Send {
358 async { problem_response(STATUS_NOT_FOUND, "method-not-found", None) }
359 }
360}
361
362impl EngineSszApi for reth_node_builder::rpc::NoopEngineApi {}
363
364impl<Provider, Pool, Validator, ChainSpec> EngineSszApi
365 for EthEngineApi<Provider, Pool, Validator, ChainSpec>
366where
367 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
368 Pool: TransactionPool + 'static,
369 Validator: EngineApiValidator<EthEngineTypes>,
370 ChainSpec: EthereumHardforks + Send + Sync + 'static,
371{
372 fn capabilities(&self, witness_enabled: bool) -> HttpResponse {
373 handle_capabilities(witness_enabled)
374 }
375
376 fn identity(&self) -> HttpResponse {
377 json_response(vec![self.client_version().clone()])
378 }
379
380 async fn new_payload(&self, fork: EngineSszFork, body: Bytes) -> HttpResponse {
381 let payload = match decode_new_payload_request(fork, &body) {
382 Ok(payload) => payload,
383 Err(err) => return err.into_response(),
384 };
385
386 match submit_payload(self, fork, payload).await {
387 Ok(status) => ssz_response(status),
388 Err(response) => response,
389 }
390 }
391
392 fn supports_witness(&self) -> bool {
393 true
394 }
395
396 async fn new_payload_with_witness(
397 &self,
398 body: Bytes,
399 witness_handler: Arc<dyn EngineSszWitness>,
400 ) -> HttpResponse {
401 let payload = match decode_new_payload_request(EngineSszFork::Amsterdam, &body) {
402 Ok(payload) => payload,
403 Err(err) => return err.into_response(),
404 };
405 let status = match submit_payload(self, EngineSszFork::Amsterdam, payload.clone()).await {
406 Ok(status) => status,
407 Err(response) => return response,
408 };
409 let witness = match status.status {
410 PayloadStatusEnum::Valid => match witness_handler.generate_witness(payload).await {
411 Ok(witness) => Some(witness),
412 Err(EngineSszWitnessError::ParentStateUnavailable { parent, source }) => {
416 debug!(
417 target: "engine::ssz",
418 %parent,
419 %source,
420 "witness omitted for valid payload"
421 );
422 None
423 }
424 Err(err) => {
425 return problem_response(
426 STATUS_INTERNAL_SERVER_ERROR,
427 "internal",
428 Some(err.to_string()),
429 )
430 }
431 },
432 _ => None,
433 };
434 ssz_response(PayloadStatusWithWitness::new(status, witness))
435 }
436
437 async fn get_payload(&self, fork: EngineSszFork, payload_id: PayloadId) -> HttpResponse {
438 match fork {
439 EngineSszFork::Paris => match self.get_payload_v2_metered(payload_id).await {
440 Ok(payload) => {
441 let block_value = payload.block_value;
442 match payload.execution_payload {
443 ExecutionPayloadFieldV2::V1(payload) => {
444 get_payload_response(BuiltPayloadParis { payload, block_value })
445 }
446 ExecutionPayloadFieldV2::V2(_) => {
447 problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
448 }
449 }
450 }
451 Err(err) => engine_error_response(err),
452 },
453 EngineSszFork::Shanghai => match self.get_payload_v2_metered(payload_id).await {
454 Ok(payload) => match BuiltPayloadShanghai::try_from(payload) {
455 Ok(payload) => get_payload_response(payload),
456 Err(err) => problem_response(
457 STATUS_UNPROCESSABLE_ENTITY,
458 "invalid-body",
459 Some(err.to_string()),
460 ),
461 },
462 Err(err) => engine_error_response(err),
463 },
464 EngineSszFork::Cancun => self
465 .get_payload_v3_metered(payload_id)
466 .await
467 .map_or_else(engine_error_response, get_payload_response),
468 EngineSszFork::Prague => self
469 .get_payload_v4_metered(payload_id)
470 .await
471 .map(BuiltPayloadPrague::from)
472 .map_or_else(engine_error_response, get_payload_response),
473 EngineSszFork::Osaka => self
474 .get_payload_v5_metered(payload_id)
475 .await
476 .map(BuiltPayloadOsaka::from)
477 .map_or_else(engine_error_response, get_payload_response),
478 EngineSszFork::Amsterdam => self
479 .get_payload_v6_metered(payload_id)
480 .await
481 .map(BuiltPayloadAmsterdam::from)
482 .map_or_else(engine_error_response, get_payload_response),
483 }
484 }
485
486 async fn forkchoice_updated(&self, fork: EngineSszFork, body: Bytes) -> HttpResponse {
487 let (state, attrs, custody_columns) = match decode_forkchoice_request(fork, &body) {
488 Ok(request) => request,
489 Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
490 };
491
492 if attrs.as_ref().is_some_and(|attrs| {
493 !timestamp_matches_fork(self.chain_spec().as_ref(), fork, attrs.timestamp)
494 }) {
495 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
496 }
497
498 let response = match fork.forkchoice_version() {
499 1 => self.fork_choice_updated_v1_metered(state, attrs).await,
500 2 => self.fork_choice_updated_v2_metered(state, attrs).await,
501 3 => self.fork_choice_updated_v3_metered(state, attrs).await,
502 4 => self.fork_choice_updated_v4_metered(state, attrs, custody_columns).await,
503 _ => return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None),
504 };
505
506 match response {
507 Ok(updated) => match ForkchoiceUpdateResponse::try_from(updated) {
508 Ok(updated) => ssz_response(updated),
509 Err(err) => problem_response(
510 STATUS_INTERNAL_SERVER_ERROR,
511 "internal",
512 Some(err.to_string()),
513 ),
514 },
515 Err(err) => engine_error_response(err),
516 }
517 }
518
519 async fn get_payload_bodies_by_hash(
520 &self,
521 fork: EngineSszFork,
522 hashes: Vec<B256>,
523 ) -> HttpResponse {
524 handle_get_payload_bodies(self.clone(), fork, PayloadBodiesRequest::Hash(hashes)).await
525 }
526
527 async fn get_payload_bodies_by_range(
528 &self,
529 fork: EngineSszFork,
530 start: u64,
531 count: u64,
532 ) -> HttpResponse {
533 handle_get_payload_bodies(self.clone(), fork, PayloadBodiesRequest::Range { start, count })
534 .await
535 }
536
537 async fn get_blobs(&self, version: u8, body: Bytes) -> HttpResponse {
538 if version == 4 {
539 let request = match BlobsV4Request::from_ssz_bytes(&body) {
540 Ok(request) => request,
541 Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
542 };
543 if request.versioned_hashes.len() > MAX_BLOB_LIMIT {
544 return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
545 }
546 return blob_response::<BlobsV4Response, _>(
547 self.get_blobs_v4_metered(request.versioned_hashes, request.indices_bitarray),
548 )
549 }
550 let request = match BlobsV1Request::from_ssz_bytes(&body) {
551 Ok(request) => request,
552 Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
553 };
554 if request.versioned_hashes.len() > MAX_BLOB_LIMIT {
555 return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
556 }
557 let hashes = request.versioned_hashes;
558 match version {
559 1 => blob_response::<BlobsV1Response, _>(self.get_blobs_v1_metered(hashes).map(Some)),
560 2 => blob_response::<BlobsV2Response, _>(self.get_blobs_v2_metered(hashes)),
561 3 => blob_response::<BlobsV3Response, _>(self.get_blobs_v3_metered(hashes)),
562 _ => problem_response(STATUS_NOT_FOUND, "method-not-found", None),
563 }
564 }
565}
566
567struct EngineSszState<Api> {
568 engine_api: Option<Api>,
569 witness_handler: Option<Arc<dyn EngineSszWitness>>,
570 witness_enabled: bool,
571}
572
573impl<Api: EngineSszApi> EngineSszState<Api> {
574 fn update_witness_support(&mut self) {
575 self.witness_enabled = self.witness_handler.is_some() &&
576 self.engine_api.as_ref().is_some_and(EngineSszApi::supports_witness);
577 }
578}
579
580async fn handle_engine_ssz_request<Api>(
581 handle: EngineSszProxyHandle<Api>,
582 request: HttpRequest,
583) -> HttpResponse
584where
585 Api: EngineSszApi,
586{
587 let Some(endpoint) = parse_engine_path(request.uri().path()) else {
588 return problem_response(STATUS_NOT_FOUND, "method-not-found", None)
589 };
590
591 if request.method() != endpoint.method() {
592 return problem_response(STATUS_METHOD_NOT_ALLOWED, "method-not-allowed", None)
593 }
594
595 match endpoint {
596 EngineSszEndpoint::Capabilities => {
597 let Some(engine_api) = handle.engine_api().await else {
598 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
599 };
600 engine_api.capabilities(handle.witness_enabled().await)
601 }
602 EngineSszEndpoint::Identity => {
603 let Some(engine_api) = handle.engine_api().await else {
604 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
605 };
606 engine_api.identity()
607 }
608 EngineSszEndpoint::NewPayload => {
609 let Some(fork) = request_fork(&request) else {
610 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
611 };
612 let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
613 Ok(body) => body,
614 Err(response) => return response,
615 };
616 let Some(engine_api) = handle.engine_api().await else {
617 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
618 };
619 engine_api.new_payload(fork, body).await
620 }
621 EngineSszEndpoint::PayloadsWithWitness => {
622 let Some(fork) = request_fork(&request) else {
623 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
624 };
625 if fork != EngineSszFork::Amsterdam {
626 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
627 }
628 let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
629 Ok(body) => body,
630 Err(response) => return response,
631 };
632 let Some(engine_api) = handle.engine_api().await else {
633 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
634 };
635 let Some(witness_handler) = handle.witness_handler().await else {
636 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
637 };
638 engine_api.new_payload_with_witness(body, witness_handler).await
639 }
640 EngineSszEndpoint::GetPayload(payload_id) => {
641 let Ok(payload_id) = payload_id else {
642 return problem_response(STATUS_BAD_REQUEST, "invalid-request", None)
643 };
644 let Some(fork) = request_fork(&request) else {
645 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
646 };
647 let Some(engine_api) = handle.engine_api().await else {
648 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
649 };
650 engine_api.get_payload(fork, payload_id).await
651 }
652 EngineSszEndpoint::Forkchoice => {
653 let Some(fork) = request_fork(&request) else {
654 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
655 };
656 let body = match read_ssz_body(request, MAX_PAYLOAD_BYTES).await {
657 Ok(body) => body,
658 Err(response) => return response,
659 };
660 let Some(engine_api) = handle.engine_api().await else {
661 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
662 };
663 engine_api.forkchoice_updated(fork, body).await
664 }
665 EngineSszEndpoint::PayloadBodiesByHash => {
666 let Some(fork) = request_fork(&request) else {
667 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
668 };
669 let body = match read_ssz_body(request, MAX_BODIES_REQUEST_BYTES).await {
670 Ok(body) => body,
671 Err(response) => return response,
672 };
673 let request = match BodiesByHashRequest::from_ssz_bytes(&body) {
674 Ok(request) if request.block_hashes.len() <= MAX_BODIES_REQUEST => {
675 request.block_hashes
676 }
677 Ok(_) => {
678 return problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None)
679 }
680 Err(_) => return problem_response(STATUS_BAD_REQUEST, "ssz-decode-error", None),
681 };
682 let Some(engine_api) = handle.engine_api().await else {
683 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
684 };
685 engine_api.get_payload_bodies_by_hash(fork, request).await
686 }
687 EngineSszEndpoint::PayloadBodiesByRange => {
688 let Some(fork) = request_fork(&request) else {
689 return problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)
690 };
691 let (start, count) =
692 match parse_bodies_range_query(request.uri().query().unwrap_or_default()) {
693 Ok(range) => range,
694 Err(response) => return response,
695 };
696 let Some(engine_api) = handle.engine_api().await else {
697 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
698 };
699 engine_api.get_payload_bodies_by_range(fork, start, count).await
700 }
701 EngineSszEndpoint::Blobs(version) => {
702 let body = match read_ssz_body(request, MAX_BLOB_REQUEST_BYTES).await {
703 Ok(body) => body,
704 Err(response) => return response,
705 };
706 let Some(engine_api) = handle.engine_api().await else {
707 return problem_response(STATUS_SERVICE_UNAVAILABLE, "service-unavailable", None)
708 };
709 engine_api.get_blobs(version, body).await
710 }
711 }
712}
713
714fn parse_method_version(version: &str) -> Option<u8> {
715 version.strip_prefix('v')?.parse().ok().filter(|version| (1..=4).contains(version))
716}
717
718fn request_fork(request: &HttpRequest) -> Option<EngineSszFork> {
719 request.headers().get(ETH_EXECUTION_VERSION)?.to_str().ok()?.parse().ok()
720}
721
722fn parse_engine_path(path: &str) -> Option<EngineSszEndpoint> {
723 let mut segments = path.trim_start_matches('/').split('/');
724 match (segments.next(), segments.next(), segments.next(), segments.next(), segments.next()) {
725 (Some("engine"), Some("v1"), Some("capabilities"), None, None) => {
726 Some(EngineSszEndpoint::Capabilities)
727 }
728 (Some("engine"), Some("v1"), Some("identity"), None, None) => {
729 Some(EngineSszEndpoint::Identity)
730 }
731 (Some("engine"), Some("v1"), Some("payloads"), None, None) => {
732 Some(EngineSszEndpoint::NewPayload)
733 }
734 (Some("engine"), Some("v1"), Some("payloads"), Some("witness"), None) => {
735 Some(EngineSszEndpoint::PayloadsWithWitness)
736 }
737 (Some("engine"), Some("v1"), Some("payloads"), Some(payload_id), None) => {
738 let payload_id = payload_id.parse::<PayloadId>();
739 Some(EngineSszEndpoint::GetPayload(payload_id))
740 }
741 (Some("engine"), Some("v1"), Some("forkchoice"), None, None) => {
742 Some(EngineSszEndpoint::Forkchoice)
743 }
744 (Some("engine"), Some("v1"), Some("blobs"), version, None) => {
745 Some(EngineSszEndpoint::Blobs(parse_method_version(version?)?))
746 }
747 (Some("engine"), Some("v1"), Some("bodies"), Some("hash"), None) => {
748 Some(EngineSszEndpoint::PayloadBodiesByHash)
749 }
750 (Some("engine"), Some("v1"), Some("bodies"), None, None) => {
751 Some(EngineSszEndpoint::PayloadBodiesByRange)
752 }
753 _ => None,
754 }
755}
756
757#[derive(Clone, Copy, Debug, Eq, PartialEq)]
758enum EngineSszEndpoint {
759 Capabilities,
760 Identity,
761 NewPayload,
762 PayloadsWithWitness,
763 GetPayload(Result<PayloadId, <PayloadId as std::str::FromStr>::Err>),
764 Forkchoice,
765 PayloadBodiesByHash,
766 PayloadBodiesByRange,
767 Blobs(u8),
768}
769
770impl EngineSszEndpoint {
771 const fn method(&self) -> &'static str {
772 match self {
773 Self::Capabilities |
774 Self::Identity |
775 Self::GetPayload(_) |
776 Self::PayloadBodiesByRange => "GET",
777 Self::NewPayload |
778 Self::PayloadsWithWitness |
779 Self::Forkchoice |
780 Self::PayloadBodiesByHash |
781 Self::Blobs(_) => "POST",
782 }
783 }
784}
785
786fn handle_capabilities(witness_enabled: bool) -> HttpResponse {
787 let mut fork_scoped_endpoints = vec!["payloads", "forkchoice", "bodies"];
788 if witness_enabled {
789 fork_scoped_endpoints.push("payloads/witness");
790 }
791 json_response(serde_json::json!({
792 "supported_forks": ["paris", "shanghai", "cancun", "prague", "osaka", "amsterdam"],
793 "fork_scoped_endpoints": fork_scoped_endpoints,
794 "independently_versioned": {
795 "blobs": ["v1", "v2", "v3", "v4"],
796 },
797 "unscoped_endpoints": ["capabilities", "identity"],
798 "limits": {
799 "bodies.max_count": MAX_BODIES_REQUEST,
800 "blobs.max_versioned_hashes": MAX_BLOB_LIMIT,
801 "payload.max_bytes": MAX_PAYLOAD_BYTES,
802 },
803 }))
804}
805
806async fn submit_payload<Provider, Pool, Validator, ChainSpec>(
807 engine_api: &EthEngineApi<Provider, Pool, Validator, ChainSpec>,
808 fork: EngineSszFork,
809 payload: ExecutionData,
810) -> Result<EngineSszPayloadStatus, HttpResponse>
811where
812 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
813 Pool: TransactionPool + 'static,
814 Validator: EngineApiValidator<EthEngineTypes>,
815 ChainSpec: EthereumHardforks + Send + Sync + 'static,
816{
817 let status = match fork.payloads_version() {
818 1 => engine_api.new_payload_v1(payload).await,
819 2 => engine_api.new_payload_v2(payload).await,
820 3 => engine_api.new_payload_v3(payload).await,
821 4 => engine_api.new_payload_v4(payload).await,
822 5 => engine_api.new_payload_v5(payload).await,
823 _ => return Err(problem_response(STATUS_BAD_REQUEST, "unsupported-fork", None)),
824 }
825 .map_err(engine_error_response)?;
826 EngineSszPayloadStatus::try_from(status).map_err(|error| {
827 problem_response(STATUS_INTERNAL_SERVER_ERROR, "internal", Some(error.to_string()))
828 })
829}
830
831fn engine_error_response(err: EngineApiError) -> HttpResponse {
832 let detail = err.to_string();
833 let error: jsonrpsee::types::ErrorObjectOwned = err.into();
834 let (status, problem_type) = match error.code() {
835 -32700 => (STATUS_BAD_REQUEST, "parse-error"),
836 -32600 => (STATUS_BAD_REQUEST, "invalid-request"),
837 -32601 => (STATUS_NOT_FOUND, "method-not-found"),
838 -32602 => (STATUS_UNPROCESSABLE_ENTITY, "invalid-body"),
839 -38001 => (STATUS_NOT_FOUND, "unknown-payload"),
840 -38002 => (STATUS_CONFLICT, "invalid-forkchoice"),
841 -38003 => (STATUS_UNPROCESSABLE_ENTITY, "invalid-attributes"),
842 -38004 => (STATUS_PAYLOAD_TOO_LARGE, "request-too-large"),
843 -38005 => (STATUS_BAD_REQUEST, "unsupported-fork"),
844 -38006 => (STATUS_CONFLICT, "reorg-too-deep"),
845 _ => (STATUS_INTERNAL_SERVER_ERROR, "internal"),
846 };
847 problem_response(status, problem_type, Some(detail))
848}
849
850enum PayloadBodiesRequest {
851 Hash(Vec<B256>),
852 Range { start: u64, count: u64 },
853}
854
855async fn handle_get_payload_bodies<Provider, Pool, Validator, ChainSpec>(
856 engine_api: EthEngineApi<Provider, Pool, Validator, ChainSpec>,
857 fork: EngineSszFork,
858 request: PayloadBodiesRequest,
859) -> HttpResponse
860where
861 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
862 Pool: TransactionPool + 'static,
863 Validator: EngineApiValidator<EthEngineTypes>,
864 ChainSpec: EthereumHardforks + Send + Sync + 'static,
865{
866 let include_bal = fork == EngineSszFork::Amsterdam;
867 let response = match request {
868 PayloadBodiesRequest::Hash(hashes) => {
869 engine_api.get_payload_bodies_by_hash_with_timestamps_metered(hashes, include_bal).await
870 }
871 PayloadBodiesRequest::Range { start, count } => {
872 engine_api
873 .get_payload_bodies_by_range_with_timestamps_metered(start, count, include_bal)
874 .await
875 }
876 };
877 let chain_spec = engine_api.chain_spec().as_ref();
878 match fork {
879 EngineSszFork::Amsterdam => payload_bodies_http_response(
880 response,
881 |body| ExecutionPayloadBodyAmsterdam::try_from(body).ok(),
882 fork,
883 chain_spec,
884 ),
885 EngineSszFork::Paris => payload_bodies_http_response(
886 response,
887 |body| ExecutionPayloadBodyParis::try_from(ExecutionPayloadBodyV1::from(body)).ok(),
888 fork,
889 chain_spec,
890 ),
891 _ => payload_bodies_http_response(
892 response,
893 |body| ExecutionPayloadBodyShanghai::try_from(ExecutionPayloadBodyV1::from(body)).ok(),
894 fork,
895 chain_spec,
896 ),
897 }
898}
899
900fn parse_bodies_range_query(query: &str) -> Result<(u64, u64), HttpResponse> {
901 let invalid = || problem_response(STATUS_BAD_REQUEST, "invalid-request", None);
902 let mut start = None;
903 let mut count = None;
904 for pair in query.split('&') {
905 let (key, value) = pair.split_once('=').ok_or_else(invalid)?;
906 let field = match key {
907 "from" if start.is_none() => &mut start,
908 "count" if count.is_none() => &mut count,
909 _ => return Err(invalid()),
910 };
911 *field = Some(value.parse::<u64>().map_err(|_| invalid())?);
912 }
913 let range = (start.ok_or_else(invalid)?, count.ok_or_else(invalid)?);
914 if range.1 > MAX_BODIES_REQUEST as u64 {
915 return Err(problem_response(STATUS_PAYLOAD_TOO_LARGE, "request-too-large", None))
916 }
917 Ok(range)
918}
919
920fn payload_bodies_response<LegacyBody, ForkBody>(
922 response: Result<Vec<Option<(u64, LegacyBody)>>, EngineApiError>,
923 convert: impl Fn(LegacyBody) -> Option<ForkBody>,
924 fork: EngineSszFork,
925 chain_spec: &impl EthereumHardforks,
926) -> Result<BodiesResponse<ForkBody>, EngineApiError>
927where
928 ForkBody: Default + ssz::Encode + ssz::Decode,
929{
930 let bodies = response?;
931 let bodies = bodies
932 .into_iter()
933 .map(|body| {
934 body.filter(|(timestamp, _)| timestamp_matches_fork(chain_spec, fork, *timestamp))
935 .map(|(_, body)| body)
936 })
937 .collect();
938 Ok(BodiesResponse::from_optional_bodies(bodies, convert))
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(err) => engine_error_response(err),
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(&inner.transactions)?;
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(transactions: &[Bytes]) -> Result<Vec<B256>, PayloadDecodeError> {
1141 let mut versioned_hashes = Vec::new();
1142 for transaction in transactions {
1143 let transaction = TxEnvelope::decode_2718_exact(transaction.as_ref())
1144 .map_err(PayloadDecodeError::InvalidTransaction)?;
1145 if let Some(hashes) = transaction.blob_versioned_hashes() {
1146 versioned_hashes.extend_from_slice(hashes);
1147 }
1148 }
1149 Ok(versioned_hashes)
1150}
1151
1152fn decode_forkchoice_request(
1153 fork: EngineSszFork,
1154 body: &[u8],
1155) -> Result<(ForkchoiceState, Option<PayloadAttributes>, Option<B128>), ssz::DecodeError> {
1156 let (state, attrs, custody) = match fork {
1157 EngineSszFork::Paris => {
1158 let ForkchoiceUpdateParis { forkchoice_state, payload_attributes } =
1159 ForkchoiceUpdateParis::from_ssz_bytes(body)?;
1160 (forkchoice_state, optional_attrs(payload_attributes), None)
1161 }
1162 EngineSszFork::Shanghai => {
1163 let ForkchoiceUpdateShanghai { forkchoice_state, payload_attributes } =
1164 ForkchoiceUpdateShanghai::from_ssz_bytes(body)?;
1165 (forkchoice_state, optional_attrs(payload_attributes), None)
1166 }
1167 EngineSszFork::Cancun => {
1168 let ForkchoiceUpdateCancun { forkchoice_state, payload_attributes } =
1169 ForkchoiceUpdateCancun::from_ssz_bytes(body)?;
1170 (forkchoice_state, optional_attrs(payload_attributes), None)
1171 }
1172 EngineSszFork::Prague => {
1173 let ForkchoiceUpdatePrague { forkchoice_state, payload_attributes } =
1174 ForkchoiceUpdatePrague::from_ssz_bytes(body)?;
1175 (forkchoice_state, optional_attrs(payload_attributes), None)
1176 }
1177 EngineSszFork::Osaka => {
1178 let ForkchoiceUpdateOsaka { forkchoice_state, payload_attributes } =
1179 ForkchoiceUpdateOsaka::from_ssz_bytes(body)?;
1180 (forkchoice_state, optional_attrs(payload_attributes), None)
1181 }
1182 EngineSszFork::Amsterdam => {
1183 let ForkchoiceUpdateAmsterdam { forkchoice_state, payload_attributes, custody_columns } =
1184 ForkchoiceUpdateAmsterdam::from_ssz_bytes(body)?;
1185 (forkchoice_state, optional_attrs(payload_attributes), custody_columns.into_option())
1186 }
1187 };
1188 if let Some(withdrawals) =
1189 attrs.as_ref().and_then(|attrs: &PayloadAttributes| attrs.withdrawals.as_ref())
1190 {
1191 check_ssz_bound(withdrawals.len(), 16, "withdrawals")?;
1192 }
1193 Ok((state, attrs, custody))
1194}
1195
1196fn optional_attrs<T>(attrs: Optional<T>) -> Option<PayloadAttributes>
1197where
1198 T: Into<PayloadAttributes>,
1199{
1200 attrs.into_option().map(Into::into)
1201}
1202
1203fn ssz_response<T: ssz::Encode>(value: T) -> HttpResponse {
1204 HttpResponse::builder()
1205 .status(STATUS_OK)
1206 .header(CONTENT_TYPE, OCTET_STREAM)
1207 .body(HttpBody::from(value.as_ssz_bytes()))
1208 .expect("valid response")
1209}
1210
1211fn get_payload_response<T: ssz::Encode>(value: T) -> HttpResponse {
1212 let mut response = ssz_response(value);
1213 response.headers_mut().insert(CACHE_CONTROL, "no-store".parse().expect("valid cache control"));
1214 response
1215}
1216
1217fn json_response<T: serde::Serialize>(value: T) -> HttpResponse {
1218 let Ok(body) = serde_json::to_string(&value) else {
1219 return problem_response(
1220 STATUS_INTERNAL_SERVER_ERROR,
1221 "internal",
1222 Some("failed to encode json".to_string()),
1223 )
1224 };
1225
1226 HttpResponse::builder()
1227 .status(STATUS_OK)
1228 .header(CONTENT_TYPE, APPLICATION_JSON)
1229 .body(HttpBody::from(body))
1230 .expect("valid response")
1231}
1232
1233fn no_content_response() -> HttpResponse {
1234 HttpResponse::builder().status(204).body(HttpBody::empty()).expect("valid response")
1235}
1236
1237fn problem_response(
1238 status: u16,
1239 problem_type: &'static str,
1240 detail: Option<String>,
1241) -> HttpResponse {
1242 let problem_type = format!("/engine-api/errors/{problem_type}");
1243 let body = match detail {
1244 Some(detail) => serde_json::json!({ "type": problem_type, "detail": detail }),
1245 None => serde_json::json!({ "type": problem_type }),
1246 };
1247
1248 HttpResponse::builder()
1249 .status(status)
1250 .header(CONTENT_TYPE, PROBLEM_JSON)
1251 .body(HttpBody::from(body.to_string()))
1252 .expect("valid response")
1253}
1254
1255#[cfg(test)]
1256mod tests {
1257 use super::*;
1258 use ssz::Encode;
1259
1260 #[tokio::test]
1261 async fn witness_capabilities_follow_both_wiring_orders() {
1262 #[derive(Clone)]
1263 struct Api(bool);
1264 impl EngineSszApi for Api {
1265 fn supports_witness(&self) -> bool {
1266 self.0
1267 }
1268 }
1269 struct Witness;
1270 impl EngineSszWitness for Witness {
1271 fn generate_witness(
1272 &self,
1273 _: ExecutionData,
1274 ) -> BoxFuture<
1275 'static,
1276 Result<crate::engine_ssz_containers::ExecutionWitnessV1, EngineSszWitnessError>,
1277 > {
1278 Box::pin(async { Ok(Default::default()) })
1279 }
1280 }
1281 let handle = EngineSszProxyHandle::new();
1282 handle.set_witness_handler_sync(Arc::new(Witness));
1283 assert!(!handle.witness_enabled().await);
1284 handle.set_engine_api_sync(Api(true));
1285 assert!(handle.witness_enabled().await);
1286 handle.set_engine_api(Api(false)).await;
1287 assert!(!handle.witness_enabled().await);
1288
1289 let handle = EngineSszProxyHandle::with_engine_api(Api(true));
1290 assert!(!handle.witness_enabled().await);
1291 handle.set_witness_handler(Arc::new(Witness)).await;
1292 assert!(handle.witness_enabled().await);
1293 }
1294
1295 #[tokio::test]
1296 async fn witness_is_only_advertised_when_configured() {
1297 for enabled in [false, true] {
1298 let body = handle_capabilities(enabled).into_body().collect().await.unwrap().to_bytes();
1299 let capabilities: serde_json::Value = serde_json::from_slice(&body).unwrap();
1300 let endpoints = capabilities["fork_scoped_endpoints"].as_array().unwrap();
1301 assert_eq!(endpoints.iter().any(|endpoint| endpoint == "payloads/witness"), enabled);
1302 }
1303 assert_eq!(
1304 parse_engine_path("/engine/v1/payloads/witness"),
1305 Some(EngineSszEndpoint::PayloadsWithWitness)
1306 );
1307 }
1308
1309 #[test]
1310 fn payload_schema_bounds() {
1311 use alloy_rpc_types_engine::{ExecutionPayloadV1, ExecutionPayloadV2, ExecutionPayloadV3};
1312 let mut payload = ExecutionPayloadV3 {
1313 payload_inner: ExecutionPayloadV2 {
1314 payload_inner: ExecutionPayloadV1::from_block_unchecked(
1315 B256::ZERO,
1316 &reth_ethereum_primitives::Block::default(),
1317 ),
1318 withdrawals: vec![],
1319 },
1320 blob_gas_used: 0,
1321 excess_blob_gas: 0,
1322 };
1323 payload.payload_inner.payload_inner.extra_data = vec![0; 32].into();
1324 payload.payload_inner.withdrawals = vec![Default::default(); 16];
1325 let mut envelope = ExecutionPayloadEnvelopePrague {
1326 payload,
1327 parent_beacon_block_root: B256::ZERO,
1328 execution_requests: Requests::new(vec![Bytes::new(); 256]),
1329 };
1330 assert!(decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()).is_ok());
1331 envelope.execution_requests = Requests::new(vec![Bytes::new(); 257]);
1332 assert!(matches!(
1333 decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1334 Err(PayloadDecodeError::Ssz(_))
1335 ));
1336 envelope.execution_requests = Requests::default();
1337 envelope.payload.payload_inner.withdrawals.push(Default::default());
1338 assert!(matches!(
1339 decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1340 Err(PayloadDecodeError::Ssz(_))
1341 ));
1342 envelope.payload.payload_inner.withdrawals.clear();
1343 envelope.payload.payload_inner.payload_inner.extra_data = vec![0; 33].into();
1344 assert!(matches!(
1345 decode_new_payload_request(EngineSszFork::Prague, &envelope.as_ssz_bytes()),
1346 Err(PayloadDecodeError::Ssz(_))
1347 ));
1348 let shanghai = ExecutionPayloadEnvelopeShanghai {
1349 payload: ExecutionPayloadV2 {
1350 payload_inner: ExecutionPayloadV1::from_block_unchecked(
1351 B256::ZERO,
1352 &reth_ethereum_primitives::Block::default(),
1353 ),
1354 withdrawals: vec![Default::default(); 17],
1355 },
1356 };
1357 assert!(matches!(
1358 decode_new_payload_request(EngineSszFork::Shanghai, &shanghai.as_ssz_bytes()),
1359 Err(PayloadDecodeError::Ssz(_))
1360 ));
1361 let paris = ExecutionPayloadEnvelopeParis {
1362 payload: ExecutionPayloadV1 {
1363 transactions: vec![Bytes::new(); (1 << 20) + 1],
1364 ..ExecutionPayloadV1::from_block_unchecked(
1365 B256::ZERO,
1366 &reth_ethereum_primitives::Block::default(),
1367 )
1368 },
1369 };
1370 assert!(matches!(
1371 decode_new_payload_request(EngineSszFork::Paris, &paris.as_ssz_bytes()),
1372 Err(PayloadDecodeError::Ssz(_))
1373 ));
1374 }
1375
1376 #[test]
1377 fn forkchoice_withdrawal_bounds() {
1378 use crate::engine_ssz_containers::{
1379 PayloadAttributesAmsterdam, PayloadAttributesCancun, PayloadAttributesShanghai,
1380 };
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 use crate::engine_ssz_containers::PayloadAttributesParis;
1425 let request = ForkchoiceUpdateParis {
1426 forkchoice_state: ForkchoiceState::default(),
1427 payload_attributes: Optional::some(PayloadAttributesParis {
1428 timestamp: 1_700_000_000,
1429 ..Default::default()
1430 }),
1431 };
1432 let bytes = request.as_ssz_bytes();
1433 assert_eq!(bytes.len(), 160);
1434 assert_eq!(&bytes[96..100], &100u32.to_le_bytes());
1435 assert_eq!(&bytes[100..108], &1_700_000_000u64.to_le_bytes());
1436 assert_eq!(
1437 decode_forkchoice_request(EngineSszFork::Paris, &bytes).unwrap().1.unwrap().timestamp,
1438 1_700_000_000
1439 );
1440 }
1441
1442 #[tokio::test]
1443 async fn post_body_limits_apply_without_content_length() {
1444 for size in [16, 17] {
1445 let request = HttpRequest::builder()
1446 .header(CONTENT_TYPE, OCTET_STREAM)
1447 .body(HttpBody::from(vec![0; size]))
1448 .unwrap();
1449 let response = read_ssz_body(request, 16).await;
1450 if size == 16 {
1451 assert_eq!(response.unwrap().len(), size);
1452 } else {
1453 assert_eq!(response.unwrap_err().status(), STATUS_PAYLOAD_TOO_LARGE);
1454 }
1455 }
1456 }
1457
1458 #[tokio::test]
1459 async fn engine_errors_preserve_validation_semantics() {
1460 use alloy_rpc_types_engine::ForkchoiceUpdateError;
1461 use reth_engine_primitives::BeaconForkChoiceUpdateError;
1462 for (error, status, kind) in [
1463 (EngineApiError::UnknownPayload, 404, "unknown-payload"),
1464 (EngineApiError::BlobRequestTooLarge { len: 129 }, 413, "request-too-large"),
1465 (
1466 EngineApiError::ForkChoiceUpdate(
1467 BeaconForkChoiceUpdateError::ForkchoiceUpdateError(
1468 ForkchoiceUpdateError::InvalidState,
1469 ),
1470 ),
1471 409,
1472 "invalid-forkchoice",
1473 ),
1474 (
1475 EngineApiError::ForkChoiceUpdate(
1476 BeaconForkChoiceUpdateError::ForkchoiceUpdateError(
1477 ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes,
1478 ),
1479 ),
1480 422,
1481 "invalid-attributes",
1482 ),
1483 ] {
1484 let response = engine_error_response(error);
1485 assert_eq!(response.status(), status);
1486 assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1487 let body = response.into_body().collect().await.unwrap().to_bytes();
1488 let problem: serde_json::Value = serde_json::from_slice(&body).unwrap();
1489 assert_eq!(problem["type"], format!("/engine-api/errors/{kind}"));
1490 }
1491 }
1492
1493 #[test]
1494 fn bodies_range_query_validation() {
1495 assert_eq!(parse_bodies_range_query("from=1&count=32").unwrap(), (1, 32));
1496 assert_eq!(parse_bodies_range_query("count=2&from=3").unwrap(), (3, 2));
1497 for query in
1498 ["", "from=1", "from=x&count=1", "from=1&count=1&count=2", "from=1&count=1&other=0"]
1499 {
1500 assert_eq!(parse_bodies_range_query(query).unwrap_err().status(), 400);
1501 }
1502 assert_eq!(parse_bodies_range_query("from=1&count=33").unwrap_err().status(), 413);
1503 }
1504
1505 #[test]
1506 fn payload_bodies_are_filtered_at_fork_boundaries() {
1507 use reth_chainspec::{ChainSpecBuilder, ForkCondition};
1508 let chain_spec = ChainSpecBuilder::default()
1509 .chain(1.into())
1510 .genesis(Default::default())
1511 .with_fork(EthereumHardfork::Shanghai, ForkCondition::Timestamp(10))
1512 .with_fork(EthereumHardfork::Cancun, ForkCondition::Timestamp(20))
1513 .with_fork(EthereumHardfork::Prague, ForkCondition::Timestamp(30))
1514 .with_fork(EthereumHardfork::Osaka, ForkCondition::Timestamp(40))
1515 .with_fork(EthereumHardfork::Amsterdam, ForkCondition::Timestamp(50))
1516 .build();
1517 for (fork, start, end) in [
1518 (EngineSszFork::Shanghai, 10, 20),
1519 (EngineSszFork::Cancun, 20, 30),
1520 (EngineSszFork::Prague, 30, 40),
1521 (EngineSszFork::Osaka, 40, 50),
1522 ] {
1523 let body = ExecutionPayloadBodyV1 {
1524 transactions: vec![Bytes::from_static(&[1, 2, 3])],
1525 withdrawals: Some(vec![]),
1526 };
1527 let response = payload_bodies_response(
1528 Ok(vec![
1529 Some((start - 1, body.clone())),
1530 Some((start, body.clone())),
1531 Some((end - 1, body.clone())),
1532 Some((end, body)),
1533 None,
1534 ]),
1535 |body| ExecutionPayloadBodyShanghai::try_from(body).ok(),
1536 fork,
1537 &chain_spec,
1538 )
1539 .unwrap();
1540 assert_eq!(
1541 response.entries.iter().map(|entry| entry.available).collect::<Vec<_>>(),
1542 [false, true, true, false, false]
1543 );
1544 for index in [0, 3, 4] {
1545 assert_eq!(response.entries[index].body, ExecutionPayloadBodyShanghai::default());
1546 }
1547 }
1548 assert!(timestamp_matches_fork(&chain_spec, EngineSszFork::Paris, 9));
1549 assert!(!timestamp_matches_fork(&chain_spec, EngineSszFork::Paris, 10));
1550 assert!(!timestamp_matches_fork(&chain_spec, EngineSszFork::Amsterdam, 49));
1551 assert!(timestamp_matches_fork(&chain_spec, EngineSszFork::Amsterdam, 50));
1552 }
1553
1554 #[tokio::test]
1555 async fn payload_bodies_hash_request_limits_and_media_type() {
1556 for content_length in [false, true] {
1557 for count in [32, 33] {
1558 let bytes =
1559 BodiesByHashRequest { block_hashes: vec![B256::ZERO; count] }.as_ssz_bytes();
1560 let mut request = HttpRequest::builder().header(CONTENT_TYPE, OCTET_STREAM);
1561 if content_length {
1562 request = request.header("content-length", bytes.len());
1563 }
1564 let result = read_ssz_body(
1565 request.body(HttpBody::from(bytes)).unwrap(),
1566 MAX_BODIES_REQUEST_BYTES,
1567 )
1568 .await;
1569 if count == 32 {
1570 assert_eq!(
1571 BodiesByHashRequest::from_ssz_bytes(&result.unwrap())
1572 .unwrap()
1573 .block_hashes
1574 .len(),
1575 32
1576 );
1577 } else {
1578 let response = result.unwrap_err();
1579 assert_eq!(response.status(), STATUS_PAYLOAD_TOO_LARGE);
1580 assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1581 }
1582 }
1583 }
1584 let response = read_ssz_body(HttpRequest::new(HttpBody::empty()), MAX_BODIES_REQUEST_BYTES)
1585 .await
1586 .unwrap_err();
1587 assert_eq!(response.status(), STATUS_UNSUPPORTED_MEDIA_TYPE);
1588 }
1589
1590 #[tokio::test]
1591 async fn payload_bodies_errors_and_capabilities_match_rest_spec() {
1592 let response =
1593 engine_error_response(EngineApiError::InvalidBodiesRange { start: 0, count: 1 });
1594 assert_eq!(response.status(), 422);
1595 assert_eq!(response.headers()[CONTENT_TYPE], PROBLEM_JSON);
1596 let body = response.into_body().collect().await.unwrap().to_bytes();
1597 let problem: serde_json::Value = serde_json::from_slice(&body).unwrap();
1598 assert_eq!(problem["type"], "/engine-api/errors/invalid-body");
1599 let body = handle_capabilities(false).into_body().collect().await.unwrap().to_bytes();
1600 let capabilities: serde_json::Value = serde_json::from_slice(&body).unwrap();
1601 assert_eq!(capabilities["limits"]["bodies.max_count"], 32);
1602 }
1603
1604 #[test]
1605 fn parses_capabilities_endpoint() {
1606 let endpoint = parse_engine_path("/engine/v1/capabilities").unwrap();
1607 assert_eq!(endpoint, EngineSszEndpoint::Capabilities);
1608 }
1609
1610 #[test]
1611 fn parses_identity_endpoint() {
1612 let endpoint = parse_engine_path("/engine/v1/identity").unwrap();
1613 assert_eq!(endpoint, EngineSszEndpoint::Identity);
1614 }
1615
1616 #[test]
1617 fn parses_fork_scoped_payload_endpoint() {
1618 let endpoint = parse_engine_path("/engine/v1/payloads").unwrap();
1619 assert_eq!(endpoint, EngineSszEndpoint::NewPayload);
1620 }
1621
1622 #[test]
1623 fn parses_fork_scoped_get_payload_endpoint() {
1624 let endpoint = parse_engine_path("/engine/v1/payloads/0x0000000000000001").unwrap();
1625 assert_eq!(
1626 endpoint,
1627 EngineSszEndpoint::GetPayload(Ok(PayloadId::new([0, 0, 0, 0, 0, 0, 0, 1])))
1628 );
1629 }
1630
1631 #[test]
1632 fn matches_malformed_get_payload_endpoint() {
1633 assert!(matches!(
1634 parse_engine_path("/engine/v1/payloads/0x01"),
1635 Some(EngineSszEndpoint::GetPayload(Err(_)))
1636 ));
1637 }
1638
1639 #[test]
1640 fn parses_fork_scoped_forkchoice_endpoint() {
1641 let endpoint = parse_engine_path("/engine/v1/forkchoice").unwrap();
1642 assert_eq!(endpoint, EngineSszEndpoint::Forkchoice);
1643 }
1644
1645 #[test]
1646 fn parses_payload_bodies_endpoints() {
1647 assert_eq!(
1648 parse_engine_path("/engine/v1/bodies/hash").unwrap(),
1649 EngineSszEndpoint::PayloadBodiesByHash
1650 );
1651 assert_eq!(
1652 parse_engine_path("/engine/v1/bodies").unwrap(),
1653 EngineSszEndpoint::PayloadBodiesByRange
1654 );
1655 }
1656
1657 #[test]
1658 fn rejects_legacy_version_scoped_endpoint() {
1659 assert!(parse_engine_path("/engine/v4/payloads").is_none());
1660 }
1661
1662 #[test]
1663 fn decodes_blob_request_container() {
1664 let hashes = vec![B256::ZERO, B256::with_last_byte(1)];
1665 let decoded = BlobsV1Request::from_ssz_bytes(
1666 &BlobsV1Request { versioned_hashes: hashes.clone() }.as_ssz_bytes(),
1667 )
1668 .unwrap()
1669 .versioned_hashes;
1670 assert_eq!(decoded, hashes);
1671 }
1672
1673 #[test]
1674 fn decodes_forkchoice_v4_with_custody_columns() {
1675 let forkchoice_state = ForkchoiceState {
1676 head_block_hash: B256::ZERO,
1677 safe_block_hash: B256::ZERO,
1678 finalized_block_hash: B256::ZERO,
1679 };
1680 let encoded = ForkchoiceUpdateAmsterdam {
1681 forkchoice_state,
1682 payload_attributes: Optional::none(),
1683 custody_columns: Optional::some(B128::with_last_byte(1)),
1684 }
1685 .as_ssz_bytes();
1686
1687 let (decoded_state, decoded_attrs, custody_columns) =
1688 decode_forkchoice_request(EngineSszFork::Amsterdam, &encoded).unwrap();
1689 assert_eq!(decoded_state, forkchoice_state);
1690 assert!(decoded_attrs.is_none());
1691 assert_eq!(custody_columns, Some(B128::with_last_byte(1)));
1692 }
1693
1694 #[test]
1695 fn decodes_forkchoice_cancun_payload_attributes() {
1696 let forkchoice_state = ForkchoiceState {
1697 head_block_hash: B256::ZERO,
1698 safe_block_hash: B256::ZERO,
1699 finalized_block_hash: B256::ZERO,
1700 };
1701 let attrs = crate::engine_ssz_containers::PayloadAttributesCancun {
1702 timestamp: 1,
1703 prev_randao: B256::with_last_byte(2),
1704 suggested_fee_recipient: Default::default(),
1705 withdrawals: Vec::new(),
1706 parent_beacon_block_root: B256::with_last_byte(3),
1707 };
1708 let encoded =
1709 ForkchoiceUpdateCancun { forkchoice_state, payload_attributes: Optional::some(attrs) }
1710 .as_ssz_bytes();
1711
1712 let (decoded_state, decoded_attrs, custody_columns) =
1713 decode_forkchoice_request(EngineSszFork::Cancun, &encoded).unwrap();
1714 assert_eq!(decoded_state, forkchoice_state);
1715 let decoded_attrs = decoded_attrs.unwrap();
1716 assert_eq!(decoded_attrs.timestamp, 1);
1717 assert!(decoded_attrs.withdrawals.as_ref().unwrap().is_empty());
1718 assert_eq!(decoded_attrs.parent_beacon_block_root, Some(B256::with_last_byte(3)));
1719 assert!(custody_columns.is_none());
1720 }
1721}