1use crate::{
2 capabilities::EngineCapabilities, metrics::EngineApiMetrics, EngineApiError, EngineApiResult,
3};
4use alloy_eips::{
5 eip1898::BlockHashOrNumber,
6 eip4844::{BlobAndProofV1, BlobAndProofV2, BlobCellsAndProofsV1},
7 eip4895::Withdrawals,
8 eip7685::RequestsOrHash,
9 Encodable2718,
10};
11use alloy_primitives::{BlockHash, BlockNumber, Bytes, Sealable, B128, B256, U64};
12use alloy_rpc_types_engine::{
13 BogotaPayloadFields, CancunPayloadFields, ClientVersionV1, ExecutionData,
14 ExecutionPayloadBodiesV1, ExecutionPayloadBodiesV2, ExecutionPayloadBodyV1,
15 ExecutionPayloadBodyV2, ExecutionPayloadInputV2, ExecutionPayloadSidecar, ExecutionPayloadV1,
16 ExecutionPayloadV3, ExecutionPayloadV4, ForkchoiceState, ForkchoiceUpdated,
17 ForkchoiceUpdatedResponseV2, PayloadId, PayloadStatus, PayloadStatusV2, PraguePayloadFields,
18 MAX_BYTES_PER_INCLUSION_LIST,
19};
20use async_trait::async_trait;
21use jsonrpsee_core::{server::RpcModule, RpcResult};
22use reth_chainspec::EthereumHardforks;
23use reth_engine_primitives::{ConsensusEngineHandle, EngineApiValidator, EngineTypes};
24use reth_network_api::{CellCustody, NetworkInfo};
25use reth_payload_builder::PayloadStore;
26use reth_payload_primitives::{
27 validate_payload_timestamp, EngineApiMessageVersion, MessageValidationKind,
28 PayloadOrAttributes, PayloadTypes,
29};
30use reth_primitives_traits::{Block, BlockBody};
31use reth_rpc_api::{EngineApiServer, IntoEngineApiRpcModule};
32use reth_storage_api::{BalProvider, BlockReader, HeaderProvider, StateProviderFactory};
33use reth_tasks::Runtime;
34use reth_transaction_pool::{BestTransactions, PoolTransaction, TransactionPool};
35use std::{
36 sync::Arc,
37 time::{Instant, SystemTime},
38};
39use tokio::sync::oneshot;
40use tracing::{debug, trace, warn};
41
42pub type EngineApiSender<Ok> = oneshot::Sender<EngineApiResult<Ok>>;
44
45const MAX_PAYLOAD_BODIES_LIMIT: u64 = 1024;
47
48const MAX_BLOB_LIMIT: usize = 128;
50
51pub struct EngineApi<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
67 inner: Arc<EngineApiInner<Provider, PayloadT, Pool, Validator, ChainSpec>>,
68}
69
70impl<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec>
71 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
72{
73 pub fn chain_spec(&self) -> &Arc<ChainSpec> {
75 &self.inner.chain_spec
76 }
77
78 pub fn client_version(&self) -> &ClientVersionV1 {
80 &self.inner.client
81 }
82}
83
84impl<Provider, PayloadT, Pool, Validator, ChainSpec>
85 EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
86where
87 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
88 PayloadT: PayloadTypes,
89 Pool: TransactionPool + 'static,
90 Validator: EngineApiValidator<PayloadT>,
91 ChainSpec: EthereumHardforks + Send + Sync + 'static,
92{
93 #[expect(clippy::too_many_arguments)]
95 pub fn new(
96 provider: Provider,
97 chain_spec: Arc<ChainSpec>,
98 beacon_consensus: ConsensusEngineHandle<PayloadT>,
99 payload_store: PayloadStore<PayloadT>,
100 tx_pool: Pool,
101 task_spawner: Runtime,
102 client: ClientVersionV1,
103 capabilities: EngineCapabilities,
104 validator: Validator,
105 accept_execution_requests_hash: bool,
106 network: impl NetworkInfo + 'static,
107 ) -> Self {
108 let cell_custody = network.cell_custody().clone();
109 let is_syncing = Arc::new(move || network.is_syncing());
110 let inner = Arc::new(EngineApiInner {
111 provider,
112 chain_spec,
113 beacon_consensus,
114 payload_store,
115 task_spawner,
116 metrics: EngineApiMetrics::default(),
117 client,
118 capabilities,
119 tx_pool,
120 validator,
121 accept_execution_requests_hash,
122 cell_custody,
123 is_syncing,
124 });
125 Self { inner }
126 }
127
128 pub fn get_client_version_v1(
130 &self,
131 _client: ClientVersionV1,
132 ) -> EngineApiResult<Vec<ClientVersionV1>> {
133 Ok(vec![self.inner.client.clone()])
134 }
135
136 async fn get_payload_timestamp(&self, payload_id: PayloadId) -> EngineApiResult<u64> {
138 Ok(self
139 .inner
140 .payload_store
141 .payload_timestamp(payload_id)
142 .await
143 .ok_or(EngineApiError::UnknownPayload)??)
144 }
145
146 pub async fn new_payload_v1(
149 &self,
150 payload: PayloadT::ExecutionData,
151 ) -> EngineApiResult<PayloadStatus> {
152 let payload_or_attrs = PayloadOrAttributes::<
153 '_,
154 PayloadT::ExecutionData,
155 PayloadT::PayloadAttributes,
156 >::from_execution_payload(&payload);
157
158 self.inner
159 .validator
160 .validate_version_specific_fields(EngineApiMessageVersion::V1, payload_or_attrs)?;
161
162 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
163 }
164
165 pub async fn new_payload_v1_metered(
167 &self,
168 payload: PayloadT::ExecutionData,
169 ) -> EngineApiResult<PayloadStatus> {
170 let start = Instant::now();
171 let res = Self::new_payload_v1(self, payload).await;
172 let elapsed = start.elapsed();
173 self.inner.metrics.latency.new_payload_v1.record(elapsed);
174 res
175 }
176
177 pub async fn new_payload_v2(
179 &self,
180 payload: PayloadT::ExecutionData,
181 ) -> EngineApiResult<PayloadStatus> {
182 let payload_or_attrs = PayloadOrAttributes::<
183 '_,
184 PayloadT::ExecutionData,
185 PayloadT::PayloadAttributes,
186 >::from_execution_payload(&payload);
187 self.inner
188 .validator
189 .validate_version_specific_fields(EngineApiMessageVersion::V2, payload_or_attrs)?;
190 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
191 }
192
193 pub async fn new_payload_v2_metered(
195 &self,
196 payload: PayloadT::ExecutionData,
197 ) -> EngineApiResult<PayloadStatus> {
198 let start = Instant::now();
199 let res = Self::new_payload_v2(self, payload).await;
200 let elapsed = start.elapsed();
201 self.inner.metrics.latency.new_payload_v2.record(elapsed);
202 res
203 }
204
205 pub async fn new_payload_v3(
207 &self,
208 payload: PayloadT::ExecutionData,
209 ) -> EngineApiResult<PayloadStatus> {
210 let payload_or_attrs = PayloadOrAttributes::<
211 '_,
212 PayloadT::ExecutionData,
213 PayloadT::PayloadAttributes,
214 >::from_execution_payload(&payload);
215 self.inner
216 .validator
217 .validate_version_specific_fields(EngineApiMessageVersion::V3, payload_or_attrs)?;
218
219 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
220 }
221
222 pub async fn new_payload_v3_metered(
224 &self,
225 payload: PayloadT::ExecutionData,
226 ) -> RpcResult<PayloadStatus> {
227 let start = Instant::now();
228
229 let res = Self::new_payload_v3(self, payload).await;
230 let elapsed = start.elapsed();
231 self.inner.metrics.latency.new_payload_v3.record(elapsed);
232 Ok(res?)
233 }
234
235 pub async fn new_payload_v4(
237 &self,
238 payload: PayloadT::ExecutionData,
239 ) -> EngineApiResult<PayloadStatus> {
240 let payload_or_attrs = PayloadOrAttributes::<
241 '_,
242 PayloadT::ExecutionData,
243 PayloadT::PayloadAttributes,
244 >::from_execution_payload(&payload);
245 self.inner
246 .validator
247 .validate_version_specific_fields(EngineApiMessageVersion::V4, payload_or_attrs)?;
248
249 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
250 }
251
252 pub async fn new_payload_v4_metered(
254 &self,
255 payload: PayloadT::ExecutionData,
256 ) -> RpcResult<PayloadStatus> {
257 let start = Instant::now();
258 let res = Self::new_payload_v4(self, payload).await;
259
260 let elapsed = start.elapsed();
261 self.inner.metrics.latency.new_payload_v4.record(elapsed);
262 Ok(res?)
263 }
264
265 pub async fn new_payload_v5(
271 &self,
272 payload: PayloadT::ExecutionData,
273 ) -> EngineApiResult<PayloadStatus> {
274 let payload_or_attrs = PayloadOrAttributes::<
275 '_,
276 PayloadT::ExecutionData,
277 PayloadT::PayloadAttributes,
278 >::from_execution_payload(&payload);
279 self.inner
280 .validator
281 .validate_version_specific_fields(EngineApiMessageVersion::V5, payload_or_attrs)?;
282 Ok(self.inner.beacon_consensus.new_payload(payload).await?)
283 }
284
285 pub async fn new_payload_v5_metered(
287 &self,
288 payload: PayloadT::ExecutionData,
289 ) -> RpcResult<PayloadStatus> {
290 let start = Instant::now();
291 let res = Self::new_payload_v5(self, payload).await;
292 let elapsed = start.elapsed();
293 self.inner.metrics.latency.new_payload_v5.record(elapsed);
294 Ok(res?)
295 }
296
297 pub async fn new_payload_v6(
301 &self,
302 payload: PayloadT::ExecutionData,
303 ) -> EngineApiResult<PayloadStatusV2> {
304 let payload_or_attrs = PayloadOrAttributes::<
305 '_,
306 PayloadT::ExecutionData,
307 PayloadT::PayloadAttributes,
308 >::from_execution_payload(&payload);
309 self.inner
310 .validator
311 .validate_version_specific_fields(EngineApiMessageVersion::V6, payload_or_attrs)?;
312
313 Ok(self.inner.beacon_consensus.new_payload(payload).await?.into())
314 }
315
316 pub async fn new_payload_v6_metered(
318 &self,
319 payload: PayloadT::ExecutionData,
320 ) -> EngineApiResult<PayloadStatusV2> {
321 let start = Instant::now();
322 let result = Self::new_payload_v6(self, payload).await;
323 self.inner.metrics.latency.new_payload_v6.record(start.elapsed());
324 result
325 }
326
327 pub fn accept_execution_requests_hash(&self) -> bool {
329 self.inner.accept_execution_requests_hash
330 }
331}
332
333impl<Provider, EngineT, Pool, Validator, ChainSpec>
334 EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
335where
336 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
337 EngineT: EngineTypes,
338 Pool: TransactionPool + 'static,
339 Validator: EngineApiValidator<EngineT>,
340 ChainSpec: EthereumHardforks + Send + Sync + 'static,
341{
342 pub async fn fork_choice_updated_v1(
349 &self,
350 state: ForkchoiceState,
351 payload_attrs: Option<EngineT::PayloadAttributes>,
352 ) -> EngineApiResult<ForkchoiceUpdated> {
353 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V1, state, payload_attrs)
354 .await
355 }
356
357 pub async fn fork_choice_updated_v1_metered(
359 &self,
360 state: ForkchoiceState,
361 payload_attrs: Option<EngineT::PayloadAttributes>,
362 ) -> EngineApiResult<ForkchoiceUpdated> {
363 let start = Instant::now();
364 let res = Self::fork_choice_updated_v1(self, state, payload_attrs).await;
365 self.inner.metrics.latency.fork_choice_updated_v1.record(start.elapsed());
366 res
367 }
368
369 pub async fn fork_choice_updated_v2(
374 &self,
375 state: ForkchoiceState,
376 payload_attrs: Option<EngineT::PayloadAttributes>,
377 ) -> EngineApiResult<ForkchoiceUpdated> {
378 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V2, state, payload_attrs)
379 .await
380 }
381
382 pub async fn fork_choice_updated_v2_metered(
384 &self,
385 state: ForkchoiceState,
386 payload_attrs: Option<EngineT::PayloadAttributes>,
387 ) -> EngineApiResult<ForkchoiceUpdated> {
388 let start = Instant::now();
389 let res = Self::fork_choice_updated_v2(self, state, payload_attrs).await;
390 self.inner.metrics.latency.fork_choice_updated_v2.record(start.elapsed());
391 res
392 }
393
394 pub async fn fork_choice_updated_v3(
399 &self,
400 state: ForkchoiceState,
401 payload_attrs: Option<EngineT::PayloadAttributes>,
402 ) -> EngineApiResult<ForkchoiceUpdated> {
403 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V3, state, payload_attrs)
404 .await
405 }
406
407 pub async fn fork_choice_updated_v3_metered(
409 &self,
410 state: ForkchoiceState,
411 payload_attrs: Option<EngineT::PayloadAttributes>,
412 ) -> EngineApiResult<ForkchoiceUpdated> {
413 let start = Instant::now();
414 let res = Self::fork_choice_updated_v3(self, state, payload_attrs).await;
415 self.inner.metrics.latency.fork_choice_updated_v3.record(start.elapsed());
416 res
417 }
418
419 pub async fn fork_choice_updated_v4(
424 &self,
425 state: ForkchoiceState,
426 payload_attrs: Option<EngineT::PayloadAttributes>,
427 custody_columns: Option<B128>,
428 ) -> EngineApiResult<ForkchoiceUpdated> {
429 if let Some(custody_columns) = custody_columns {
430 self.inner.cell_custody.set_from_engine_api(custody_columns);
431 }
432 self.validate_and_execute_forkchoice(EngineApiMessageVersion::V4, state, payload_attrs)
433 .await
434 }
435
436 pub async fn fork_choice_updated_v4_metered(
438 &self,
439 state: ForkchoiceState,
440 payload_attrs: Option<EngineT::PayloadAttributes>,
441 custody_columns: Option<B128>,
442 ) -> EngineApiResult<ForkchoiceUpdated> {
443 let start = Instant::now();
444 let res = Self::fork_choice_updated_v4(self, state, payload_attrs, custody_columns).await;
445 self.inner.metrics.latency.fork_choice_updated_v4.record(start.elapsed());
446 res
447 }
448
449 pub async fn fork_choice_updated_v5(
453 &self,
454 state: ForkchoiceState,
455 payload_attrs: Option<EngineT::PayloadAttributes>,
456 custody_columns: Option<B128>,
457 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
458 if let Some(custody_columns) = custody_columns {
459 self.inner.cell_custody.set_from_engine_api(custody_columns);
460 }
461
462 Ok(self
464 .validate_and_execute_forkchoice(EngineApiMessageVersion::V5, state, payload_attrs)
465 .await?
466 .into())
467 }
468
469 pub async fn fork_choice_updated_v5_metered(
471 &self,
472 state: ForkchoiceState,
473 payload_attrs: Option<EngineT::PayloadAttributes>,
474 custody_columns: Option<B128>,
475 ) -> EngineApiResult<ForkchoiceUpdatedResponseV2> {
476 let start = Instant::now();
477 let res = Self::fork_choice_updated_v5(self, state, payload_attrs, custody_columns).await;
478 self.inner.metrics.latency.fork_choice_updated_v5.record(start.elapsed());
479 res
480 }
481
482 pub fn get_inclusion_list_v1(&self) -> EngineApiResult<Vec<Bytes>> {
484 let mut total_size = 0;
485 let mut inclusion_list = Vec::new();
486
487 for pool_tx in self.inner.tx_pool.best_transactions().without_blobs().without_updates() {
488 let encoded: Bytes = pool_tx.transaction.consensus_ref().encoded_2718().into();
489 let new_size = total_size + alloy_rlp::Encodable::length(&encoded);
490 if new_size + alloy_rlp::length_of_length(new_size) >
491 MAX_BYTES_PER_INCLUSION_LIST as usize
492 {
493 break
494 }
495
496 total_size = new_size;
497 inclusion_list.push(encoded);
498 }
499
500 Ok(inclusion_list)
501 }
502
503 pub fn get_inclusion_list_v1_metered(&self) -> EngineApiResult<Vec<Bytes>> {
505 let start = Instant::now();
506 let result = Self::get_inclusion_list_v1(self);
507 self.inner.metrics.latency.get_inclusion_list_v1.record(start.elapsed());
508 result
509 }
510
511 async fn get_built_payload(
513 &self,
514 payload_id: PayloadId,
515 ) -> EngineApiResult<EngineT::BuiltPayload> {
516 self.inner
517 .payload_store
518 .resolve(payload_id)
519 .await
520 .ok_or(EngineApiError::UnknownPayload)?
521 .map_err(|_| EngineApiError::UnknownPayload)
522 }
523
524 async fn get_payload_inner<R>(
527 &self,
528 payload_id: PayloadId,
529 version: EngineApiMessageVersion,
530 ) -> EngineApiResult<R>
531 where
532 EngineT::BuiltPayload: TryInto<R>,
533 {
534 let timestamp = self.get_payload_timestamp(payload_id).await?;
537 validate_payload_timestamp(
538 &self.inner.chain_spec,
539 version,
540 timestamp,
541 MessageValidationKind::GetPayload,
542 )?;
543
544 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
546 warn!(?version, "could not transform built payload");
547 EngineApiError::UnknownPayload
548 })
549 }
550
551 pub async fn get_payload_v1(
561 &self,
562 payload_id: PayloadId,
563 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
564 self.get_built_payload(payload_id).await?.try_into().map_err(|_| {
565 warn!(version = ?EngineApiMessageVersion::V1, "could not transform built payload");
566 EngineApiError::UnknownPayload
567 })
568 }
569
570 pub async fn get_payload_v1_metered(
572 &self,
573 payload_id: PayloadId,
574 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV1> {
575 let start = Instant::now();
576 let res = Self::get_payload_v1(self, payload_id).await;
577 self.inner.metrics.latency.get_payload_v1.record(start.elapsed());
578 res
579 }
580
581 pub async fn get_payload_v2(
589 &self,
590 payload_id: PayloadId,
591 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
592 self.get_payload_inner(payload_id, EngineApiMessageVersion::V2).await
593 }
594
595 pub async fn get_payload_v2_metered(
597 &self,
598 payload_id: PayloadId,
599 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV2> {
600 let start = Instant::now();
601 let res = Self::get_payload_v2(self, payload_id).await;
602 self.inner.metrics.latency.get_payload_v2.record(start.elapsed());
603 res
604 }
605
606 pub async fn get_payload_v3(
614 &self,
615 payload_id: PayloadId,
616 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
617 self.get_payload_inner(payload_id, EngineApiMessageVersion::V3).await
618 }
619
620 pub async fn get_payload_v3_metered(
622 &self,
623 payload_id: PayloadId,
624 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV3> {
625 let start = Instant::now();
626 let res = Self::get_payload_v3(self, payload_id).await;
627 self.inner.metrics.latency.get_payload_v3.record(start.elapsed());
628 res
629 }
630
631 pub async fn get_payload_v4(
639 &self,
640 payload_id: PayloadId,
641 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
642 self.get_payload_inner(payload_id, EngineApiMessageVersion::V4).await
643 }
644
645 pub async fn get_payload_v4_metered(
647 &self,
648 payload_id: PayloadId,
649 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV4> {
650 let start = Instant::now();
651 let res = Self::get_payload_v4(self, payload_id).await;
652 self.inner.metrics.latency.get_payload_v4.record(start.elapsed());
653 res
654 }
655
656 pub async fn get_payload_v5(
666 &self,
667 payload_id: PayloadId,
668 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
669 self.get_payload_inner(payload_id, EngineApiMessageVersion::V5).await
670 }
671
672 pub async fn get_payload_v5_metered(
674 &self,
675 payload_id: PayloadId,
676 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV5> {
677 let start = Instant::now();
678 let res = Self::get_payload_v5(self, payload_id).await;
679 self.inner.metrics.latency.get_payload_v5.record(start.elapsed());
680 res
681 }
682
683 pub async fn get_payload_v6(
689 &self,
690 payload_id: PayloadId,
691 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
692 self.get_payload_inner(payload_id, EngineApiMessageVersion::V6).await
693 }
694
695 pub async fn get_payload_v6_metered(
697 &self,
698 payload_id: PayloadId,
699 ) -> EngineApiResult<EngineT::ExecutionPayloadEnvelopeV6> {
700 let start = Instant::now();
701 let res = Self::get_payload_v6(self, payload_id).await;
702 self.inner.metrics.latency.get_payload_v6.record(start.elapsed());
703 res
704 }
705
706 pub async fn get_payload_bodies_by_range_with<F, R>(
709 &self,
710 start: BlockNumber,
711 count: u64,
712 f: F,
713 ) -> EngineApiResult<Vec<Option<R>>>
714 where
715 F: Fn(Provider::Block) -> R + Send + 'static,
716 R: Send + 'static,
717 {
718 let (tx, rx) = oneshot::channel();
719 let inner = self.inner.clone();
720
721 self.inner.task_spawner.spawn_blocking_task(async move {
722 if count > MAX_PAYLOAD_BODIES_LIMIT {
723 tx.send(Err(EngineApiError::PayloadRequestTooLarge { len: count })).ok();
724 return;
725 }
726
727 if start == 0 || count == 0 {
728 tx.send(Err(EngineApiError::InvalidBodiesRange { start, count })).ok();
729 return;
730 }
731
732 let mut result = Vec::with_capacity(count as usize);
733
734 let mut end = start.saturating_add(count - 1);
736
737 if let Ok(best_block) = inner.provider.best_block_number()
740 && end > best_block {
741 end = best_block;
742 }
743
744 let earliest_block = inner.provider.earliest_block_number().unwrap_or(0);
746 for num in start..=end {
747 if tx.is_closed() {
748 return;
749 }
750
751 if num < earliest_block {
752 result.push(None);
753 continue;
754 }
755 let block_result = inner.provider.block(BlockHashOrNumber::Number(num));
756 match block_result {
757 Ok(block) => {
758 result.push(block.map(&f));
759 }
760 Err(err) => {
761 tx.send(Err(EngineApiError::Internal(Box::new(err)))).ok();
762 return;
763 }
764 };
765 }
766 tx.send(Ok(result)).ok();
767 });
768
769 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
770 }
771
772 pub async fn get_payload_bodies_by_range_v1(
783 &self,
784 start: BlockNumber,
785 count: u64,
786 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
787 self.get_payload_bodies_by_range_with(start, count, |block| ExecutionPayloadBodyV1 {
788 transactions: block.body().encoded_2718_transactions(),
789 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
790 })
791 .await
792 }
793
794 pub async fn get_payload_bodies_by_range_v1_metered(
796 &self,
797 start: BlockNumber,
798 count: u64,
799 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
800 let start_time = Instant::now();
801 let res = Self::get_payload_bodies_by_range_v1(self, start, count).await;
802 self.inner.metrics.latency.get_payload_bodies_by_range_v1.record(start_time.elapsed());
803 res
804 }
805
806 pub async fn get_payload_bodies_by_range_v2(
810 &self,
811 start: BlockNumber,
812 count: u64,
813 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
814 let mut payload_bodies = self
815 .get_payload_bodies_by_range_with(start, count, |block| {
816 let block_hash = block.header().hash_slow();
817 (
818 block_hash,
819 ExecutionPayloadBodyV2 {
820 transactions: block.body().encoded_2718_transactions(),
821 withdrawals: block
822 .body()
823 .withdrawals()
824 .cloned()
825 .map(Withdrawals::into_inner),
826 block_access_list: None,
827 },
828 )
829 })
830 .await?;
831
832 let block_hashes = payload_bodies
833 .iter()
834 .filter_map(|payload_body| payload_body.as_ref().map(|(block_hash, _)| *block_hash))
835 .collect::<Vec<_>>();
836 let block_access_lists = self.get_block_access_lists_by_hashes(block_hashes).await?;
837
838 for (payload_body, block_access_list) in
839 payload_bodies.iter_mut().filter_map(Option::as_mut).zip(block_access_lists)
840 {
841 payload_body.1.block_access_list = block_access_list;
842 }
843
844 Ok(payload_bodies
845 .into_iter()
846 .map(|payload_body| payload_body.map(|(_, payload_body)| payload_body))
847 .collect())
848 }
849
850 pub async fn get_payload_bodies_by_range_v2_metered(
852 &self,
853 start: BlockNumber,
854 count: u64,
855 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
856 let start_time = Instant::now();
857 let res = Self::get_payload_bodies_by_range_v2(self, start, count).await;
858 self.inner.metrics.latency.get_payload_bodies_by_range_v2.record(start_time.elapsed());
859 res
860 }
861
862 pub async fn get_payload_bodies_by_hash_with<F, R>(
864 &self,
865 hashes: Vec<BlockHash>,
866 f: F,
867 ) -> EngineApiResult<Vec<Option<R>>>
868 where
869 F: Fn(Provider::Block) -> R + Send + 'static,
870 R: Send + 'static,
871 {
872 let len = hashes.len() as u64;
873 if len > MAX_PAYLOAD_BODIES_LIMIT {
874 return Err(EngineApiError::PayloadRequestTooLarge { len });
875 }
876
877 let (tx, rx) = oneshot::channel();
878 let inner = self.inner.clone();
879
880 self.inner.task_spawner.spawn_blocking_task(async move {
881 let mut result = Vec::with_capacity(hashes.len());
882 for hash in hashes {
883 if tx.is_closed() {
884 return;
885 }
886
887 let block_result = inner.provider.block(BlockHashOrNumber::Hash(hash));
888 match block_result {
889 Ok(block) => {
890 result.push(block.map(&f));
891 }
892 Err(err) => {
893 let _ = tx.send(Err(EngineApiError::Internal(Box::new(err))));
894 return;
895 }
896 }
897 }
898 tx.send(Ok(result)).ok();
899 });
900
901 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
902 }
903
904 async fn get_block_access_lists_by_hashes(
905 &self,
906 hashes: Vec<BlockHash>,
907 ) -> EngineApiResult<Vec<Option<Bytes>>> {
908 let len = hashes.len() as u64;
909 if len > MAX_PAYLOAD_BODIES_LIMIT {
910 return Err(EngineApiError::PayloadRequestTooLarge { len });
911 }
912
913 let (tx, rx) = oneshot::channel();
914 let inner = self.inner.clone();
915
916 self.inner.task_spawner.spawn_blocking_task(async move {
917 if tx.is_closed() {
918 return;
919 }
920
921 tx.send(
922 inner
923 .provider
924 .get_bals_by_hashes(&hashes)
925 .map_err(|err| EngineApiError::Internal(Box::new(err))),
926 )
927 .ok();
928 });
929
930 rx.await.map_err(|err| EngineApiError::Internal(Box::new(err)))?
931 }
932
933 pub async fn get_payload_bodies_by_hash_v1(
935 &self,
936 hashes: Vec<BlockHash>,
937 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
938 self.get_payload_bodies_by_hash_with(hashes, |block| ExecutionPayloadBodyV1 {
939 transactions: block.body().encoded_2718_transactions(),
940 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
941 })
942 .await
943 }
944
945 pub async fn get_payload_bodies_by_hash_v1_metered(
947 &self,
948 hashes: Vec<BlockHash>,
949 ) -> EngineApiResult<ExecutionPayloadBodiesV1> {
950 let start = Instant::now();
951 let res = Self::get_payload_bodies_by_hash_v1(self, hashes).await;
952 self.inner.metrics.latency.get_payload_bodies_by_hash_v1.record(start.elapsed());
953 res
954 }
955
956 pub async fn get_payload_bodies_by_hash_v2(
960 &self,
961 hashes: Vec<BlockHash>,
962 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
963 let payload_bodies =
964 self.get_payload_bodies_by_hash_with(hashes.clone(), |block| ExecutionPayloadBodyV2 {
965 transactions: block.body().encoded_2718_transactions(),
966 withdrawals: block.body().withdrawals().cloned().map(Withdrawals::into_inner),
967 block_access_list: None,
968 });
969 let block_access_lists = self.get_block_access_lists_by_hashes(hashes);
970 let (mut payload_bodies, block_access_lists) =
971 tokio::try_join!(payload_bodies, block_access_lists)?;
972
973 for (payload_body, block_access_list) in payload_bodies.iter_mut().zip(block_access_lists) {
974 if let Some(payload_body) = payload_body {
975 payload_body.block_access_list = block_access_list;
976 }
977 }
978
979 Ok(payload_bodies)
980 }
981
982 pub async fn get_payload_bodies_by_hash_v2_metered(
984 &self,
985 hashes: Vec<BlockHash>,
986 ) -> EngineApiResult<ExecutionPayloadBodiesV2> {
987 let start = Instant::now();
988 let res = Self::get_payload_bodies_by_hash_v2(self, hashes).await;
989 self.inner.metrics.latency.get_payload_bodies_by_hash_v2.record(start.elapsed());
990 res
991 }
992
993 async fn validate_and_execute_forkchoice(
1007 &self,
1008 version: EngineApiMessageVersion,
1009 state: ForkchoiceState,
1010 payload_attrs: Option<EngineT::PayloadAttributes>,
1011 ) -> EngineApiResult<ForkchoiceUpdated> {
1012 if let Some(ref attrs) = payload_attrs {
1013 let attr_validation_res =
1014 self.inner.validator.ensure_well_formed_attributes(version, attrs);
1015
1016 if let Err(err) = attr_validation_res {
1026 let fcu_res = self.inner.beacon_consensus.fork_choice_updated(state, None).await?;
1027 if fcu_res.is_invalid() || fcu_res.payload_status.is_syncing() {
1028 return Ok(fcu_res)
1029 }
1030 return Err(err.into())
1031 }
1032 }
1033
1034 Ok(self.inner.beacon_consensus.fork_choice_updated(state, payload_attrs).await?)
1035 }
1036
1037 pub fn capabilities(&self) -> &EngineCapabilities {
1039 &self.inner.capabilities
1040 }
1041
1042 fn has_blobs(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1043 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1044 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1045 }
1046
1047 self.inner
1048 .tx_pool
1049 .has_blobs_for_versioned_hashes(&versioned_hashes)
1050 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1051 }
1052
1053 pub fn has_blobs_metered(&self, versioned_hashes: Vec<B256>) -> EngineApiResult<Vec<bool>> {
1055 let start = Instant::now();
1056 let res = Self::has_blobs(self, versioned_hashes);
1057 self.inner.metrics.latency.has_blobs.record(start.elapsed());
1058 res
1059 }
1060
1061 fn get_blobs_v1(
1062 &self,
1063 versioned_hashes: Vec<B256>,
1064 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1065 let current_timestamp =
1067 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1068 if self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1069 return Err(EngineApiError::EngineObjectValidationError(
1070 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1071 ));
1072 }
1073
1074 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1075 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1076 }
1077
1078 self.inner
1079 .tx_pool
1080 .get_blobs_for_versioned_hashes_v1(&versioned_hashes)
1081 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1082 }
1083
1084 pub fn get_blobs_v1_metered(
1086 &self,
1087 versioned_hashes: Vec<B256>,
1088 ) -> EngineApiResult<Vec<Option<BlobAndProofV1>>> {
1089 let hashes_len = versioned_hashes.len();
1090 let start = Instant::now();
1091 let res = Self::get_blobs_v1(self, versioned_hashes);
1092 self.inner.metrics.latency.get_blobs_v1.record(start.elapsed());
1093
1094 if let Ok(blobs) = &res {
1095 let blobs_found = blobs.iter().flatten().count();
1096 let blobs_missed = hashes_len - blobs_found;
1097
1098 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1099 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1100 }
1101
1102 res
1103 }
1104
1105 fn get_blobs_v2(
1106 &self,
1107 versioned_hashes: Vec<B256>,
1108 ) -> EngineApiResult<Option<Vec<BlobAndProofV2>>> {
1109 let current_timestamp =
1111 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1112 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1113 return Err(EngineApiError::EngineObjectValidationError(
1114 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1115 ));
1116 }
1117
1118 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1119 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1120 }
1121
1122 self.inner
1123 .tx_pool
1124 .get_blobs_for_versioned_hashes_v2(&versioned_hashes)
1125 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1126 }
1127
1128 fn get_blobs_v3(
1129 &self,
1130 versioned_hashes: Vec<B256>,
1131 ) -> EngineApiResult<Option<Vec<Option<BlobAndProofV2>>>> {
1132 let current_timestamp =
1134 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1135 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1136 return Err(EngineApiError::EngineObjectValidationError(
1137 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1138 ));
1139 }
1140
1141 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1142 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1143 }
1144
1145 if (*self.inner.is_syncing)() {
1147 return Ok(None)
1148 }
1149
1150 self.inner
1151 .tx_pool
1152 .get_blobs_for_versioned_hashes_v3(&versioned_hashes)
1153 .map(Some)
1154 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1155 }
1156
1157 fn get_blobs_v4(
1158 &self,
1159 versioned_hashes: Vec<B256>,
1160 indices_bitarray: B128,
1161 ) -> EngineApiResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1162 let indices_bitarray = B128::from(u128::from_le_bytes(indices_bitarray.into()));
1165 let current_timestamp =
1166 SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap_or_default().as_secs();
1167 if !self.inner.chain_spec.is_osaka_active_at_timestamp(current_timestamp) {
1168 return Err(EngineApiError::EngineObjectValidationError(
1169 reth_payload_primitives::EngineObjectValidationError::UnsupportedFork,
1170 ));
1171 }
1172
1173 if versioned_hashes.len() > MAX_BLOB_LIMIT {
1174 return Err(EngineApiError::BlobRequestTooLarge { len: versioned_hashes.len() })
1175 }
1176
1177 if (*self.inner.is_syncing)() {
1179 return Ok(None)
1180 }
1181
1182 self.inner
1183 .tx_pool
1184 .get_blobs_for_versioned_hashes_v4(&versioned_hashes, indices_bitarray)
1185 .map(Some)
1186 .map_err(|err| EngineApiError::Internal(Box::new(err)))
1187 }
1188
1189 pub fn get_blobs_v2_metered(
1191 &self,
1192 versioned_hashes: Vec<B256>,
1193 ) -> EngineApiResult<Option<Vec<BlobAndProofV2>>> {
1194 let hashes_len = versioned_hashes.len();
1195 let start = Instant::now();
1196 let res = Self::get_blobs_v2(self, versioned_hashes);
1197 self.inner.metrics.latency.get_blobs_v2.record(start.elapsed());
1198
1199 if let Ok(blobs) = &res {
1200 let blobs_found = blobs.iter().flatten().count();
1201
1202 self.inner
1203 .metrics
1204 .blob_metrics
1205 .get_blobs_requests_blobs_total
1206 .increment(hashes_len as u64);
1207 self.inner
1208 .metrics
1209 .blob_metrics
1210 .get_blobs_requests_blobs_in_blobpool_total
1211 .increment(blobs_found as u64);
1212
1213 if blobs_found == hashes_len {
1214 self.inner.metrics.blob_metrics.get_blobs_requests_success_total.increment(1);
1215 } else {
1216 self.inner.metrics.blob_metrics.get_blobs_requests_failure_total.increment(1);
1217 }
1218 } else {
1219 self.inner.metrics.blob_metrics.get_blobs_requests_failure_total.increment(1);
1220 }
1221
1222 res
1223 }
1224
1225 pub fn get_blobs_v3_metered(
1227 &self,
1228 versioned_hashes: Vec<B256>,
1229 ) -> EngineApiResult<Option<Vec<Option<BlobAndProofV2>>>> {
1230 let hashes_len = versioned_hashes.len();
1231 let start = Instant::now();
1232 let res = Self::get_blobs_v3(self, versioned_hashes);
1233 self.inner.metrics.latency.get_blobs_v3.record(start.elapsed());
1234
1235 if let Ok(Some(blobs)) = &res {
1236 let blobs_found = blobs.iter().flatten().count();
1237 let blobs_missed = hashes_len - blobs_found;
1238
1239 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1240 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1241 }
1242
1243 res
1244 }
1245
1246 pub fn get_blobs_v4_metered(
1248 &self,
1249 versioned_hashes: Vec<B256>,
1250 indices_bitarray: B128,
1251 ) -> EngineApiResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1252 let hashes_len = versioned_hashes.len();
1253 let start = Instant::now();
1254 let res = Self::get_blobs_v4(self, versioned_hashes, indices_bitarray);
1255 self.inner.metrics.latency.get_blobs_v4.record(start.elapsed());
1256
1257 if let Ok(Some(blobs)) = &res {
1258 let blobs_found = blobs.iter().flatten().count();
1259 let blobs_missed = hashes_len - blobs_found;
1260
1261 self.inner.metrics.blob_metrics.blob_count.increment(blobs_found as u64);
1262 self.inner.metrics.blob_metrics.blob_misses.increment(blobs_missed as u64);
1263 }
1264
1265 res
1266 }
1267}
1268
1269#[async_trait]
1271impl<Provider, EngineT, Pool, Validator, ChainSpec> EngineApiServer<EngineT>
1272 for EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
1273where
1274 Provider: HeaderProvider + BlockReader + StateProviderFactory + BalProvider + 'static,
1275 EngineT: EngineTypes<ExecutionData = ExecutionData>,
1276 Pool: TransactionPool + 'static,
1277 Validator: EngineApiValidator<EngineT>,
1278 ChainSpec: EthereumHardforks + Send + Sync + 'static,
1279{
1280 async fn new_payload_v1(&self, payload: ExecutionPayloadV1) -> RpcResult<PayloadStatus> {
1284 trace!(target: "rpc::engine", "Serving engine_newPayloadV1");
1285 let payload =
1286 ExecutionData { payload: payload.into(), sidecar: ExecutionPayloadSidecar::none() };
1287 Ok(self.new_payload_v1_metered(payload).await?)
1288 }
1289
1290 async fn new_payload_v2(&self, payload: ExecutionPayloadInputV2) -> RpcResult<PayloadStatus> {
1293 trace!(target: "rpc::engine", "Serving engine_newPayloadV2");
1294 let payload = ExecutionData {
1295 payload: payload.into_payload(),
1296 sidecar: ExecutionPayloadSidecar::none(),
1297 };
1298
1299 Ok(self.new_payload_v2_metered(payload).await?)
1300 }
1301
1302 async fn new_payload_v3(
1305 &self,
1306 payload: ExecutionPayloadV3,
1307 versioned_hashes: Vec<B256>,
1308 parent_beacon_block_root: B256,
1309 ) -> RpcResult<PayloadStatus> {
1310 trace!(target: "rpc::engine", "Serving engine_newPayloadV3");
1311 let payload = ExecutionData {
1312 payload: payload.into(),
1313 sidecar: ExecutionPayloadSidecar::v3(CancunPayloadFields {
1314 versioned_hashes,
1315 parent_beacon_block_root,
1316 }),
1317 };
1318
1319 Ok(self.new_payload_v3_metered(payload).await?)
1320 }
1321
1322 async fn new_payload_v4(
1325 &self,
1326 payload: ExecutionPayloadV3,
1327 versioned_hashes: Vec<B256>,
1328 parent_beacon_block_root: B256,
1329 requests: RequestsOrHash,
1330 ) -> RpcResult<PayloadStatus> {
1331 trace!(target: "rpc::engine", "Serving engine_newPayloadV4");
1332
1333 if requests.is_hash() && !self.inner.accept_execution_requests_hash {
1335 return Err(EngineApiError::UnexpectedRequestsHash.into());
1336 }
1337
1338 let payload = ExecutionData {
1339 payload: payload.into(),
1340 sidecar: ExecutionPayloadSidecar::v4(
1341 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1342 PraguePayloadFields { requests },
1343 ),
1344 };
1345
1346 Ok(self.new_payload_v4_metered(payload).await?)
1347 }
1348
1349 async fn new_payload_v5(
1355 &self,
1356 payload: ExecutionPayloadV4,
1357 versioned_hashes: Vec<B256>,
1358 parent_beacon_block_root: B256,
1359 requests: RequestsOrHash,
1360 ) -> RpcResult<PayloadStatus> {
1361 trace!(target: "rpc::engine", "Serving engine_newPayloadV5");
1362 if requests.is_hash() && !self.inner.accept_execution_requests_hash {
1364 return Err(EngineApiError::UnexpectedRequestsHash.into());
1365 }
1366
1367 let payload = ExecutionData {
1368 payload: payload.into(),
1369 sidecar: ExecutionPayloadSidecar::v4(
1370 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1371 PraguePayloadFields { requests },
1372 ),
1373 };
1374
1375 Ok(self.new_payload_v5_metered(payload).await?)
1376 }
1377
1378 async fn new_payload_v6(
1382 &self,
1383 payload: ExecutionPayloadV4,
1384 versioned_hashes: Vec<B256>,
1385 parent_beacon_block_root: B256,
1386 execution_requests: RequestsOrHash,
1387 inclusion_list_transactions: Vec<Bytes>,
1388 ) -> RpcResult<PayloadStatusV2> {
1389 trace!(target: "rpc::engine", "Serving engine_newPayloadV6");
1390 if execution_requests.is_hash() && !self.inner.accept_execution_requests_hash {
1391 return Err(EngineApiError::UnexpectedRequestsHash.into());
1392 }
1393
1394 let payload = ExecutionData {
1395 payload: payload.into(),
1396 sidecar: ExecutionPayloadSidecar::v6(
1397 CancunPayloadFields { versioned_hashes, parent_beacon_block_root },
1398 PraguePayloadFields { requests: execution_requests },
1399 BogotaPayloadFields { inclusion_list_transactions },
1400 ),
1401 };
1402
1403 Ok(self.new_payload_v6_metered(payload).await?)
1406 }
1407
1408 async fn fork_choice_updated_v1(
1413 &self,
1414 fork_choice_state: ForkchoiceState,
1415 payload_attributes: Option<EngineT::PayloadAttributes>,
1416 ) -> RpcResult<ForkchoiceUpdated> {
1417 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV1");
1418 Ok(self.fork_choice_updated_v1_metered(fork_choice_state, payload_attributes).await?)
1419 }
1420
1421 async fn fork_choice_updated_v2(
1424 &self,
1425 fork_choice_state: ForkchoiceState,
1426 payload_attributes: Option<EngineT::PayloadAttributes>,
1427 ) -> RpcResult<ForkchoiceUpdated> {
1428 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV2");
1429 Ok(self.fork_choice_updated_v2_metered(fork_choice_state, payload_attributes).await?)
1430 }
1431
1432 async fn fork_choice_updated_v3(
1436 &self,
1437 fork_choice_state: ForkchoiceState,
1438 payload_attributes: Option<EngineT::PayloadAttributes>,
1439 ) -> RpcResult<ForkchoiceUpdated> {
1440 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV3");
1441 Ok(self.fork_choice_updated_v3_metered(fork_choice_state, payload_attributes).await?)
1442 }
1443
1444 async fn fork_choice_updated_v4(
1448 &self,
1449 fork_choice_state: ForkchoiceState,
1450 payload_attributes: Option<EngineT::PayloadAttributes>,
1451 custody_columns: Option<B128>,
1452 ) -> RpcResult<ForkchoiceUpdated> {
1453 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV4");
1454 Ok(self
1455 .fork_choice_updated_v4_metered(fork_choice_state, payload_attributes, custody_columns)
1456 .await?)
1457 }
1458
1459 async fn fork_choice_updated_v5(
1463 &self,
1464 fork_choice_state: ForkchoiceState,
1465 payload_attributes: Option<EngineT::PayloadAttributes>,
1466 custody_columns: Option<B128>,
1467 ) -> RpcResult<ForkchoiceUpdatedResponseV2> {
1468 trace!(target: "rpc::engine", "Serving engine_forkchoiceUpdatedV5");
1469 Ok(self
1470 .fork_choice_updated_v5_metered(fork_choice_state, payload_attributes, custody_columns)
1471 .await?)
1472 }
1473
1474 async fn get_payload_v1(
1486 &self,
1487 payload_id: PayloadId,
1488 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV1> {
1489 trace!(target: "rpc::engine", "Serving engine_getPayloadV1");
1490 Ok(self.get_payload_v1_metered(payload_id).await?)
1491 }
1492
1493 async fn get_payload_v2(
1503 &self,
1504 payload_id: PayloadId,
1505 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV2> {
1506 debug!(target: "rpc::engine", id = %payload_id, "Serving engine_getPayloadV2");
1507 Ok(self.get_payload_v2_metered(payload_id).await?)
1508 }
1509
1510 async fn get_payload_v3(
1520 &self,
1521 payload_id: PayloadId,
1522 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV3> {
1523 trace!(target: "rpc::engine", "Serving engine_getPayloadV3");
1524 Ok(self.get_payload_v3_metered(payload_id).await?)
1525 }
1526
1527 async fn get_payload_v4(
1537 &self,
1538 payload_id: PayloadId,
1539 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV4> {
1540 trace!(target: "rpc::engine", "Serving engine_getPayloadV4");
1541 Ok(self.get_payload_v4_metered(payload_id).await?)
1542 }
1543
1544 async fn get_payload_v5(
1554 &self,
1555 payload_id: PayloadId,
1556 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV5> {
1557 trace!(target: "rpc::engine", "Serving engine_getPayloadV5");
1558 Ok(self.get_payload_v5_metered(payload_id).await?)
1559 }
1560
1561 async fn get_payload_v6(
1567 &self,
1568 payload_id: PayloadId,
1569 ) -> RpcResult<EngineT::ExecutionPayloadEnvelopeV6> {
1570 trace!(target: "rpc::engine", "Serving engine_getPayloadV6");
1571 Ok(self.get_payload_v6_metered(payload_id).await?)
1572 }
1573
1574 async fn get_inclusion_list_v1(&self) -> RpcResult<Vec<Bytes>> {
1578 trace!(target: "rpc::engine", "Serving engine_getInclusionListV1");
1579 Ok(self.get_inclusion_list_v1_metered()?)
1580 }
1581
1582 async fn get_payload_bodies_by_hash_v1(
1585 &self,
1586 block_hashes: Vec<BlockHash>,
1587 ) -> RpcResult<ExecutionPayloadBodiesV1> {
1588 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByHashV1");
1589 Ok(self.get_payload_bodies_by_hash_v1_metered(block_hashes).await?)
1590 }
1591
1592 async fn get_payload_bodies_by_hash_v2(
1598 &self,
1599 block_hashes: Vec<BlockHash>,
1600 ) -> RpcResult<ExecutionPayloadBodiesV2> {
1601 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByHashV2");
1602 Ok(self.get_payload_bodies_by_hash_v2_metered(block_hashes).await?)
1603 }
1604
1605 async fn get_payload_bodies_by_range_v1(
1622 &self,
1623 start: U64,
1624 count: U64,
1625 ) -> RpcResult<ExecutionPayloadBodiesV1> {
1626 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByRangeV1");
1627 Ok(self.get_payload_bodies_by_range_v1_metered(start.to(), count.to()).await?)
1628 }
1629
1630 async fn get_payload_bodies_by_range_v2(
1636 &self,
1637 start: U64,
1638 count: U64,
1639 ) -> RpcResult<ExecutionPayloadBodiesV2> {
1640 trace!(target: "rpc::engine", "Serving engine_getPayloadBodiesByRangeV2");
1641 Ok(self.get_payload_bodies_by_range_v2_metered(start.to(), count.to()).await?)
1642 }
1643
1644 async fn get_client_version_v1(
1648 &self,
1649 client: ClientVersionV1,
1650 ) -> RpcResult<Vec<ClientVersionV1>> {
1651 trace!(target: "rpc::engine", "Serving engine_getClientVersionV1");
1652 Ok(Self::get_client_version_v1(self, client)?)
1653 }
1654
1655 async fn exchange_capabilities(&self, capabilities: Vec<String>) -> RpcResult<Vec<String>> {
1658 trace!(target: "rpc::engine", "Serving engine_exchangeCapabilities");
1659
1660 let el_caps = self.capabilities();
1661 el_caps.log_capability_mismatches(&capabilities);
1662
1663 Ok(el_caps.list())
1664 }
1665
1666 async fn has_blobs(&self, versioned_hashes: Vec<B256>) -> RpcResult<Vec<bool>> {
1667 trace!(target: "rpc::engine", "Serving engine_hasBlobs");
1668 Ok(self.has_blobs_metered(versioned_hashes)?)
1669 }
1670
1671 async fn get_blobs_v1(
1672 &self,
1673 versioned_hashes: Vec<B256>,
1674 ) -> RpcResult<Vec<Option<BlobAndProofV1>>> {
1675 trace!(target: "rpc::engine", "Serving engine_getBlobsV1");
1676 Ok(self.get_blobs_v1_metered(versioned_hashes)?)
1677 }
1678
1679 async fn get_blobs_v2(
1680 &self,
1681 versioned_hashes: Vec<B256>,
1682 ) -> RpcResult<Option<Vec<BlobAndProofV2>>> {
1683 trace!(target: "rpc::engine", "Serving engine_getBlobsV2");
1684 Ok(self.get_blobs_v2_metered(versioned_hashes)?)
1685 }
1686
1687 async fn get_blobs_v3(
1688 &self,
1689 versioned_hashes: Vec<B256>,
1690 ) -> RpcResult<Option<Vec<Option<BlobAndProofV2>>>> {
1691 trace!(target: "rpc::engine", "Serving engine_getBlobsV3");
1692 Ok(self.get_blobs_v3_metered(versioned_hashes)?)
1693 }
1694
1695 async fn get_blobs_v4(
1696 &self,
1697 versioned_hashes: Vec<B256>,
1698 indices_bitarray: B128,
1699 ) -> RpcResult<Option<Vec<Option<BlobCellsAndProofsV1>>>> {
1700 trace!(target: "rpc::engine", "Serving engine_getBlobsV4");
1701 Ok(self.get_blobs_v4_metered(versioned_hashes, indices_bitarray)?)
1702 }
1703}
1704
1705impl<Provider, EngineT, Pool, Validator, ChainSpec> IntoEngineApiRpcModule
1706 for EngineApi<Provider, EngineT, Pool, Validator, ChainSpec>
1707where
1708 EngineT: EngineTypes,
1709 Self: EngineApiServer<EngineT>,
1710{
1711 fn into_rpc_module(self) -> RpcModule<()> {
1712 EngineApiServer::<EngineT>::into_rpc(self).remove_context()
1713 }
1714}
1715
1716impl<Provider, PayloadT, Pool, Validator, ChainSpec> std::fmt::Debug
1717 for EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
1718where
1719 PayloadT: PayloadTypes,
1720{
1721 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1722 f.debug_struct("EngineApi").finish_non_exhaustive()
1723 }
1724}
1725
1726impl<Provider, PayloadT, Pool, Validator, ChainSpec> Clone
1727 for EngineApi<Provider, PayloadT, Pool, Validator, ChainSpec>
1728where
1729 PayloadT: PayloadTypes,
1730{
1731 fn clone(&self) -> Self {
1732 Self { inner: Arc::clone(&self.inner) }
1733 }
1734}
1735
1736struct EngineApiInner<Provider, PayloadT: PayloadTypes, Pool, Validator, ChainSpec> {
1738 provider: Provider,
1740 chain_spec: Arc<ChainSpec>,
1742 beacon_consensus: ConsensusEngineHandle<PayloadT>,
1744 payload_store: PayloadStore<PayloadT>,
1746 task_spawner: Runtime,
1748 metrics: EngineApiMetrics,
1750 client: ClientVersionV1,
1752 capabilities: EngineCapabilities,
1754 tx_pool: Pool,
1756 validator: Validator,
1758 accept_execution_requests_hash: bool,
1759 cell_custody: CellCustody,
1761 is_syncing: Arc<dyn Fn() -> bool + Send + Sync>,
1763}
1764
1765#[cfg(test)]
1766mod tests {
1767 use super::*;
1768 use alloy_eips::{eip7685::Requests, NumHash};
1769 use alloy_primitives::{Address, Bytes, B256};
1770 use alloy_rpc_types_engine::{
1771 ClientCode, ClientVersionV1, ExecutionPayloadV2, PayloadAttributes, PayloadStatusEnum,
1772 };
1773 use assert_matches::assert_matches;
1774 use reth_chainspec::{ChainSpec, ChainSpecBuilder, MAINNET};
1775 use reth_engine_primitives::{BeaconEngineMessage, OnForkChoiceUpdated};
1776 use reth_ethereum_engine_primitives::EthEngineTypes;
1777 use reth_ethereum_primitives::{Block, TransactionSigned};
1778 use reth_network_api::{
1779 noop::NoopNetwork, EthProtocolInfo, NetworkError, NetworkInfo, NetworkStatus,
1780 };
1781 use reth_node_ethereum::EthereumEngineValidator;
1782 use reth_payload_builder::test_utils::spawn_test_payload_service;
1783 use reth_primitives_traits::SignedTransaction;
1784 use reth_provider::{test_utils::MockEthProvider, BalStoreHandle, InMemoryBalStore, RawBal};
1785 use reth_tasks::Runtime;
1786 use reth_transaction_pool::{
1787 blobstore::InMemoryBlobStore,
1788 noop::NoopTransactionPool,
1789 test_utils::{OkValidator, TransactionBuilder},
1790 CoinbaseTipOrdering, EthPooledTransaction, Pool, PoolTransaction, TransactionOrigin,
1791 };
1792 use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver};
1793
1794 type EthTestPool = Pool<
1795 OkValidator<EthPooledTransaction>,
1796 CoinbaseTipOrdering<EthPooledTransaction>,
1797 InMemoryBlobStore,
1798 >;
1799
1800 fn eth_test_pool() -> EthTestPool {
1801 Pool::new(
1802 OkValidator::default(),
1803 CoinbaseTipOrdering::default(),
1804 InMemoryBlobStore::default(),
1805 Default::default(),
1806 )
1807 }
1808
1809 fn pooled_transaction(transaction: TransactionSigned) -> EthPooledTransaction {
1810 let transaction = transaction.try_into_recovered().unwrap();
1811 let encoded_length = transaction.encode_2718_len();
1812 EthPooledTransaction::new(transaction, encoded_length)
1813 }
1814
1815 fn setup_engine_api() -> (
1816 EngineApiTestHandle,
1817 EngineApi<
1818 Arc<MockEthProvider>,
1819 EthEngineTypes,
1820 NoopTransactionPool,
1821 EthereumEngineValidator,
1822 ChainSpec,
1823 >,
1824 ) {
1825 setup_engine_api_with_pool(NoopTransactionPool::default())
1826 }
1827
1828 fn setup_engine_api_with_pool<Pool>(
1829 tx_pool: Pool,
1830 ) -> (
1831 EngineApiTestHandle,
1832 EngineApi<Arc<MockEthProvider>, EthEngineTypes, Pool, EthereumEngineValidator, ChainSpec>,
1833 )
1834 where
1835 Pool: TransactionPool + 'static,
1836 {
1837 let client = ClientVersionV1 {
1838 code: ClientCode::RH,
1839 name: "Reth".to_string(),
1840 version: "v0.2.0-beta.5".to_string(),
1841 commit: "defa64b2".to_string(),
1842 };
1843
1844 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
1845 let provider = Arc::new(MockEthProvider::default());
1846 let payload_store = spawn_test_payload_service();
1847 let (to_engine, engine_rx) = unbounded_channel();
1848 let task_executor = Runtime::test();
1849 let api = EngineApi::new(
1850 provider.clone(),
1851 chain_spec.clone(),
1852 ConsensusEngineHandle::new(to_engine),
1853 payload_store.into(),
1854 tx_pool,
1855 task_executor,
1856 client,
1857 EngineCapabilities::default(),
1858 EthereumEngineValidator::new(chain_spec.clone()),
1859 false,
1860 NoopNetwork::default(),
1861 );
1862 let handle = EngineApiTestHandle { chain_spec, provider, from_api: engine_rx };
1863 (handle, api)
1864 }
1865
1866 #[tokio::test]
1867 async fn engine_client_version_v1() {
1868 let client = ClientVersionV1 {
1869 code: ClientCode::RH,
1870 name: "Reth".to_string(),
1871 version: "v0.2.0-beta.5".to_string(),
1872 commit: "defa64b2".to_string(),
1873 };
1874 let (_, api) = setup_engine_api();
1875 let res = api.get_client_version_v1(client.clone());
1876 assert_eq!(res.unwrap(), vec![client]);
1877 }
1878
1879 #[tokio::test]
1880 async fn get_inclusion_list_v1_returns_empty_list() {
1881 let (_, api) = setup_engine_api();
1882
1883 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
1884 assert!(res.is_empty());
1885 }
1886
1887 #[tokio::test]
1888 async fn get_inclusion_list_v1_stops_at_size_limit() {
1889 let pool = eth_test_pool();
1890 let first = pooled_transaction(
1891 TransactionBuilder::default()
1892 .max_fee_per_gas(3_000_000_000u128)
1893 .input(vec![0; 4_200])
1894 .into_legacy(),
1895 );
1896 let second = pooled_transaction(
1897 TransactionBuilder::default()
1898 .max_fee_per_gas(2_000_000_000u128)
1899 .input(vec![0; 4_200])
1900 .into_legacy(),
1901 );
1902 let third = pooled_transaction(
1903 TransactionBuilder::default().max_fee_per_gas(1_000_000_000u128).into_legacy(),
1904 );
1905 let expected: Bytes = first.consensus_ref().encoded_2718().into();
1906
1907 pool.add_transaction(TransactionOrigin::External, first).await.unwrap();
1908 pool.add_transaction(TransactionOrigin::External, second).await.unwrap();
1909 pool.add_transaction(TransactionOrigin::External, third).await.unwrap();
1910 let (_, api) = setup_engine_api_with_pool(pool);
1911
1912 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
1913 assert_eq!(res, vec![expected]);
1914 assert!(
1915 alloy_rlp::list_length::<Bytes, [u8]>(&res) <= MAX_BYTES_PER_INCLUSION_LIST as usize
1916 );
1917 }
1918
1919 #[tokio::test]
1920 async fn get_inclusion_list_v1_excludes_blob_transactions() {
1921 let pool = eth_test_pool();
1922 let blob = pooled_transaction(
1923 TransactionBuilder::default()
1924 .max_fee_per_gas(2_000_000_000u128)
1925 .max_priority_fee_per_gas(1_000_000_000u128)
1926 .into_eip4844(),
1927 );
1928 let non_blob = pooled_transaction(
1929 TransactionBuilder::default()
1930 .max_fee_per_gas(1_000_000_000u128)
1931 .max_priority_fee_per_gas(1_000_000_000u128)
1932 .into_eip1559(),
1933 );
1934 let expected: Bytes = non_blob.consensus_ref().encoded_2718().into();
1935
1936 pool.add_transaction(TransactionOrigin::External, blob).await.unwrap();
1937 pool.add_transaction(TransactionOrigin::External, non_blob).await.unwrap();
1938 let (_, api) = setup_engine_api_with_pool(pool);
1939
1940 let res = EngineApiServer::get_inclusion_list_v1(&api).await.unwrap();
1941 assert_eq!(res, vec![expected]);
1942 }
1943
1944 #[tokio::test]
1945 async fn has_blobs_returns_ordered_availability() {
1946 let (_, api) = setup_engine_api();
1947
1948 let res = api.has_blobs_metered(vec![B256::ZERO, B256::with_last_byte(1)]).unwrap();
1949 assert_eq!(res, vec![false, false]);
1950 }
1951
1952 #[tokio::test]
1953 async fn has_blobs_rejects_large_requests() {
1954 let (_, api) = setup_engine_api();
1955
1956 let res = api.has_blobs_metered(vec![B256::ZERO; MAX_BLOB_LIMIT + 1]);
1957 assert_matches!(
1958 res,
1959 Err(EngineApiError::BlobRequestTooLarge { len }) if len == MAX_BLOB_LIMIT + 1
1960 );
1961 }
1962
1963 #[tokio::test]
1964 async fn get_payload_bodies_by_hash_v2_returns_block_access_list_from_store() {
1965 let bal_store = BalStoreHandle::new(InMemoryBalStore::default());
1966 let mut provider = MockEthProvider::default();
1967 provider.bal_store = bal_store.clone();
1968 let provider = Arc::new(provider);
1969
1970 let client = ClientVersionV1 {
1971 code: ClientCode::RH,
1972 name: "Reth".to_string(),
1973 version: "v0.2.0-beta.5".to_string(),
1974 commit: "defa64b2".to_string(),
1975 };
1976 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
1977 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
1978 let (to_engine, _engine_rx) = unbounded_channel();
1979 let api = EngineApi::new(
1980 provider.clone(),
1981 chain_spec.clone(),
1982 ConsensusEngineHandle::new(to_engine),
1983 payload_store.into(),
1984 NoopTransactionPool::default(),
1985 Runtime::test(),
1986 client,
1987 EngineCapabilities::default(),
1988 EthereumEngineValidator::new(chain_spec),
1989 false,
1990 NoopNetwork::default(),
1991 );
1992
1993 let mut block = Block::default();
1994 block.header.number = 1;
1995 let block_hash = block.header.hash_slow();
1996 provider.add_block(block_hash, block);
1997
1998 let mut block_without_bal = Block::default();
1999 block_without_bal.header.number = 2;
2000 let block_without_bal_hash = block_without_bal.header.hash_slow();
2001 provider.add_block(block_without_bal_hash, block_without_bal);
2002
2003 let raw_bal = Bytes::from_static(&[alloy_rlp::EMPTY_LIST_CODE]);
2004 bal_store.insert(NumHash::new(1, block_hash), RawBal::new(raw_bal.clone())).unwrap();
2005
2006 let missing_hash = B256::with_last_byte(3);
2007 let response = api
2008 .get_payload_bodies_by_hash_v2(vec![block_hash, block_without_bal_hash, missing_hash])
2009 .await
2010 .unwrap();
2011
2012 assert_eq!(response.len(), 3);
2013 assert_eq!(response[0].as_ref().unwrap().block_access_list, Some(raw_bal));
2014 assert_eq!(response[1].as_ref().unwrap().block_access_list, None);
2015 assert!(response[2].is_none());
2016 }
2017
2018 #[tokio::test]
2019 async fn get_payload_bodies_by_range_v2_returns_block_access_lists_from_store() {
2020 let bal_store = BalStoreHandle::new(InMemoryBalStore::default());
2021 let mut provider = MockEthProvider::default();
2022 provider.bal_store = bal_store.clone();
2023 let provider = Arc::new(provider);
2024
2025 let client = ClientVersionV1 {
2026 code: ClientCode::RH,
2027 name: "Reth".to_string(),
2028 version: "v0.2.0-beta.5".to_string(),
2029 commit: "defa64b2".to_string(),
2030 };
2031 let chain_spec: Arc<ChainSpec> = MAINNET.clone();
2032 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2033 let (to_engine, _engine_rx) = unbounded_channel();
2034 let api = EngineApi::new(
2035 provider.clone(),
2036 chain_spec.clone(),
2037 ConsensusEngineHandle::new(to_engine),
2038 payload_store.into(),
2039 NoopTransactionPool::default(),
2040 Runtime::test(),
2041 client,
2042 EngineCapabilities::default(),
2043 EthereumEngineValidator::new(chain_spec),
2044 false,
2045 NoopNetwork::default(),
2046 );
2047
2048 let mut block = Block::default();
2049 block.header.number = 1;
2050 let block_hash = block.header.hash_slow();
2051 provider.add_block(block_hash, block);
2052
2053 let mut block_without_bal = Block::default();
2054 block_without_bal.header.number = 3;
2055 let block_without_bal_hash = block_without_bal.header.hash_slow();
2056 provider.add_block(block_without_bal_hash, block_without_bal);
2057
2058 let raw_bal = Bytes::from_static(&[alloy_rlp::EMPTY_LIST_CODE]);
2059 bal_store.insert(NumHash::new(1, block_hash), RawBal::new(raw_bal.clone())).unwrap();
2060
2061 let response = api.get_payload_bodies_by_range_v2(1, 3).await.unwrap();
2062
2063 assert_eq!(response.len(), 3);
2064 assert_eq!(response[0].as_ref().unwrap().block_access_list, Some(raw_bal));
2065 assert!(response[1].is_none());
2066 assert_eq!(response[2].as_ref().unwrap().block_access_list, None);
2067 }
2068
2069 struct EngineApiTestHandle {
2070 #[allow(dead_code)]
2071 chain_spec: Arc<ChainSpec>,
2072 provider: Arc<MockEthProvider>,
2073 from_api: UnboundedReceiver<BeaconEngineMessage<EthEngineTypes>>,
2074 }
2075
2076 #[tokio::test]
2077 async fn forwards_responses_to_consensus_engine() {
2078 let (mut handle, api) = setup_engine_api();
2079
2080 tokio::spawn(async move {
2081 let payload_v1 = ExecutionPayloadV1::from_block_slow(&Block::default());
2082 let execution_data = ExecutionData {
2083 payload: payload_v1.into(),
2084 sidecar: ExecutionPayloadSidecar::none(),
2085 };
2086
2087 api.new_payload_v1(execution_data).await.unwrap();
2088 });
2089 assert_matches!(handle.from_api.recv().await, Some(BeaconEngineMessage::NewPayload { .. }));
2090 }
2091
2092 #[tokio::test]
2093 async fn new_payload_v5_accepts_amsterdam_payloads() {
2094 let chain_spec = Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2095 let provider = Arc::new(MockEthProvider::default());
2096 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2097 let (to_engine, mut engine_rx) = unbounded_channel();
2098
2099 let api = EngineApi::new(
2100 provider,
2101 chain_spec.clone(),
2102 ConsensusEngineHandle::new(to_engine),
2103 payload_store.into(),
2104 NoopTransactionPool::default(),
2105 Runtime::test(),
2106 ClientVersionV1 {
2107 code: ClientCode::RH,
2108 name: "Reth".to_string(),
2109 version: "v0.0.0-test".to_string(),
2110 commit: "test".to_string(),
2111 },
2112 EngineCapabilities::default(),
2113 EthereumEngineValidator::new(chain_spec),
2114 false,
2115 NoopNetwork::default(),
2116 );
2117
2118 tokio::spawn(async move {
2119 let payload_v1 = ExecutionPayloadV1::from_block_slow(&Block::default());
2120 let payload = ExecutionPayloadV4 {
2121 payload_inner: ExecutionPayloadV3 {
2122 payload_inner: ExecutionPayloadV2 {
2123 payload_inner: payload_v1,
2124 withdrawals: Vec::new(),
2125 },
2126 blob_gas_used: 0,
2127 excess_blob_gas: 0,
2128 },
2129 block_access_list: Bytes::from_static(b"bal"),
2130 slot_number: 1,
2131 };
2132 let execution_data = ExecutionData {
2133 payload: payload.into(),
2134 sidecar: ExecutionPayloadSidecar::v4(
2135 CancunPayloadFields {
2136 versioned_hashes: Vec::new(),
2137 parent_beacon_block_root: B256::ZERO,
2138 },
2139 PraguePayloadFields { requests: RequestsOrHash::Requests(Requests::default()) },
2140 ),
2141 };
2142
2143 api.new_payload_v5(execution_data).await.unwrap();
2144 });
2145
2146 assert_matches!(engine_rx.recv().await, Some(BeaconEngineMessage::NewPayload { .. }));
2147 }
2148
2149 #[derive(Clone)]
2150 struct TestNetworkInfo {
2151 syncing: bool,
2152 }
2153
2154 impl NetworkInfo for TestNetworkInfo {
2155 fn local_addr(&self) -> std::net::SocketAddr {
2156 (std::net::Ipv4Addr::UNSPECIFIED, 0).into()
2157 }
2158
2159 async fn network_status(&self) -> Result<NetworkStatus, NetworkError> {
2160 #[allow(deprecated)]
2161 Ok(NetworkStatus {
2162 client_version: "test".to_string(),
2163 protocol_version: 5,
2164 eth_protocol_info: EthProtocolInfo {
2165 network: 1,
2166 difficulty: None,
2167 genesis: Default::default(),
2168 config: Default::default(),
2169 head: Default::default(),
2170 },
2171 capabilities: vec![],
2172 })
2173 }
2174
2175 fn chain_id(&self) -> u64 {
2176 1
2177 }
2178
2179 fn cell_custody(&self) -> &CellCustody {
2180 static CELL_CUSTODY: std::sync::OnceLock<CellCustody> = std::sync::OnceLock::new();
2181 CELL_CUSTODY.get_or_init(CellCustody::default)
2182 }
2183
2184 fn is_syncing(&self) -> bool {
2185 self.syncing
2186 }
2187
2188 fn is_initially_syncing(&self) -> bool {
2189 self.syncing
2190 }
2191 }
2192
2193 #[tokio::test]
2194 async fn get_blobs_v3_returns_null_when_syncing() {
2195 let chain_spec: Arc<ChainSpec> =
2196 Arc::new(ChainSpecBuilder::mainnet().osaka_activated().build());
2197 let provider = Arc::new(MockEthProvider::default());
2198 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2199 let (to_engine, _engine_rx) = unbounded_channel::<BeaconEngineMessage<EthEngineTypes>>();
2200
2201 let api = EngineApi::new(
2202 provider,
2203 chain_spec.clone(),
2204 ConsensusEngineHandle::new(to_engine),
2205 payload_store.into(),
2206 NoopTransactionPool::default(),
2207 Runtime::test(),
2208 ClientVersionV1 {
2209 code: ClientCode::RH,
2210 name: "Reth".to_string(),
2211 version: "v0.0.0-test".to_string(),
2212 commit: "test".to_string(),
2213 },
2214 EngineCapabilities::default(),
2215 EthereumEngineValidator::new(chain_spec),
2216 false,
2217 TestNetworkInfo { syncing: true },
2218 );
2219
2220 let res = api.get_blobs_v3_metered(vec![B256::ZERO]);
2221 assert_matches!(res, Ok(None));
2222 }
2223
2224 #[tokio::test]
2225 async fn get_blobs_v4_returns_null_when_syncing() {
2226 let chain_spec: Arc<ChainSpec> =
2227 Arc::new(ChainSpecBuilder::mainnet().osaka_activated().build());
2228 let provider = Arc::new(MockEthProvider::default());
2229 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2230 let (to_engine, _engine_rx) = unbounded_channel::<BeaconEngineMessage<EthEngineTypes>>();
2231
2232 let api = EngineApi::new(
2233 provider,
2234 chain_spec.clone(),
2235 ConsensusEngineHandle::new(to_engine),
2236 payload_store.into(),
2237 NoopTransactionPool::default(),
2238 Runtime::test(),
2239 ClientVersionV1 {
2240 code: ClientCode::RH,
2241 name: "Reth".to_string(),
2242 version: "v0.0.0-test".to_string(),
2243 commit: "test".to_string(),
2244 },
2245 EngineCapabilities::default(),
2246 EthereumEngineValidator::new(chain_spec),
2247 false,
2248 TestNetworkInfo { syncing: true },
2249 );
2250
2251 let res = api.get_blobs_v4_metered(vec![B256::ZERO], B128::from(1u128.to_le_bytes()));
2252 assert_matches!(res, Ok(None));
2253 }
2254
2255 #[test]
2256 fn engine_bitvector_uses_little_endian_cell_indices() {
2257 assert_eq!(
2258 B128::from(u128::from_le_bytes(B128::from(1u128.to_le_bytes()).into())),
2259 B128::from(1u128)
2260 );
2261 assert_eq!(
2262 B128::from(u128::from_le_bytes(B128::from((1u128 << 127).to_le_bytes()).into())),
2263 B128::from(1u128 << 127)
2264 );
2265 assert_eq!(
2266 B128::from(u128::from_le_bytes(B128::from(((1u128 << 64) - 1).to_le_bytes()).into(),)),
2267 B128::from((1u128 << 64) - 1)
2268 );
2269 }
2270
2271 #[tokio::test]
2272 async fn fcu_v4_updates_shared_cell_custody_before_forkchoice_result() {
2273 let chain_spec: Arc<ChainSpec> =
2274 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2275 let provider = Arc::new(MockEthProvider::default());
2276 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2277 let (to_engine, mut engine_rx) = unbounded_channel();
2278 let network = NoopNetwork::default();
2279 let cell_custody = network.cell_custody().clone();
2280
2281 let api = EngineApi::new(
2282 provider,
2283 chain_spec.clone(),
2284 ConsensusEngineHandle::new(to_engine),
2285 payload_store.into(),
2286 NoopTransactionPool::default(),
2287 Runtime::test(),
2288 ClientVersionV1 {
2289 code: ClientCode::RH,
2290 name: "Reth".to_string(),
2291 version: "v0.0.0-test".to_string(),
2292 commit: "test".to_string(),
2293 },
2294 EngineCapabilities::default(),
2295 EthereumEngineValidator::new(chain_spec),
2296 false,
2297 network,
2298 );
2299
2300 let state = ForkchoiceState {
2301 head_block_hash: B256::from([0x33; 32]),
2302 safe_block_hash: B256::ZERO,
2303 finalized_block_hash: B256::ZERO,
2304 };
2305 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2306 let expected_custody_columns = B128::from(0b1010u128);
2307
2308 let api_task = tokio::spawn(async move {
2309 api.fork_choice_updated_v4(state, None, Some(custody_columns)).await
2310 });
2311
2312 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2313 .await
2314 .expect("timed out waiting for forkchoiceUpdated request")
2315 .expect("expected forkchoiceUpdated request");
2316 let response_tx = match request {
2317 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2318 assert!(payload_attrs.is_none());
2319 tx
2320 }
2321 other => panic!("unexpected engine message: {other:?}"),
2322 };
2323 assert_eq!(cell_custody.get(), expected_custody_columns);
2324
2325 response_tx
2326 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2327 PayloadStatusEnum::Valid,
2328 ))))
2329 .expect("send valid response");
2330
2331 api_task
2332 .await
2333 .expect("api task should not panic")
2334 .expect("forkchoiceUpdatedV4 should succeed");
2335 assert_eq!(cell_custody.get(), expected_custody_columns);
2336 }
2337
2338 #[tokio::test]
2339 async fn fcu_v4_updates_shared_cell_custody_when_payload_attrs_invalid() {
2340 let chain_spec: Arc<ChainSpec> =
2341 Arc::new(ChainSpecBuilder::mainnet().amsterdam_activated().build());
2342 let provider = Arc::new(MockEthProvider::default());
2343 let payload_store = spawn_test_payload_service::<EthEngineTypes>();
2344 let (to_engine, mut engine_rx) = unbounded_channel();
2345 let network = NoopNetwork::default();
2346 let cell_custody = network.cell_custody().clone();
2347
2348 let api = EngineApi::new(
2349 provider,
2350 chain_spec.clone(),
2351 ConsensusEngineHandle::new(to_engine),
2352 payload_store.into(),
2353 NoopTransactionPool::default(),
2354 Runtime::test(),
2355 ClientVersionV1 {
2356 code: ClientCode::RH,
2357 name: "Reth".to_string(),
2358 version: "v0.0.0-test".to_string(),
2359 commit: "test".to_string(),
2360 },
2361 EngineCapabilities::default(),
2362 EthereumEngineValidator::new(chain_spec),
2363 false,
2364 network,
2365 );
2366
2367 let state = ForkchoiceState {
2368 head_block_hash: B256::from([0x44; 32]),
2369 safe_block_hash: B256::ZERO,
2370 finalized_block_hash: B256::ZERO,
2371 };
2372 let payload_attributes = PayloadAttributes {
2373 timestamp: 1,
2374 prev_randao: B256::ZERO,
2375 suggested_fee_recipient: Address::ZERO,
2376 withdrawals: Some(vec![]),
2377 parent_beacon_block_root: None,
2378 slot_number: None,
2379 ..Default::default()
2380 };
2381 let custody_columns = B128::from(0b1010u128.to_le_bytes());
2382 let expected_custody_columns = B128::from(0b1010u128);
2383
2384 let api_task = tokio::spawn(async move {
2385 api.fork_choice_updated_v4(state, Some(payload_attributes), Some(custody_columns)).await
2386 });
2387
2388 let request = tokio::time::timeout(std::time::Duration::from_secs(1), engine_rx.recv())
2389 .await
2390 .expect("timed out waiting for forkchoiceUpdated request")
2391 .expect("expected forkchoiceUpdated request");
2392 let response_tx = match request {
2393 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2394 assert!(
2395 payload_attrs.is_none(),
2396 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2397 );
2398 tx
2399 }
2400 other => panic!("unexpected engine message: {other:?}"),
2401 };
2402 assert_eq!(cell_custody.get(), expected_custody_columns);
2403
2404 response_tx
2405 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2406 PayloadStatusEnum::Valid,
2407 ))))
2408 .expect("send valid response");
2409
2410 let response = api_task.await.expect("api task should not panic");
2411 assert_matches!(
2412 response,
2413 Err(EngineApiError::EngineObjectValidationError(
2414 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2415 ))
2416 );
2417 }
2418
2419 #[tokio::test]
2420 async fn fcu_v3_syncing_precedes_invalid_payload_attributes_validation() {
2421 let (mut handle, api) = setup_engine_api();
2422
2423 let state = ForkchoiceState {
2424 head_block_hash: B256::from([0x11; 32]),
2425 safe_block_hash: B256::ZERO,
2426 finalized_block_hash: B256::ZERO,
2427 };
2428 let payload_attributes = PayloadAttributes {
2429 timestamp: 1,
2430 prev_randao: B256::ZERO,
2431 suggested_fee_recipient: Address::ZERO,
2432 withdrawals: Some(vec![]),
2433 parent_beacon_block_root: None,
2435 slot_number: None,
2436 ..Default::default()
2437 };
2438
2439 let api_task = tokio::spawn(async move {
2440 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2441 });
2442
2443 let request =
2444 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2445 .await
2446 .expect("timed out waiting for forkchoiceUpdated request")
2447 .expect("expected forkchoiceUpdated request");
2448 let response_tx = match request {
2449 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2450 assert!(
2451 payload_attrs.is_none(),
2452 "FCU for syncing state should be evaluated before payload attributes"
2453 );
2454 tx
2455 }
2456 other => panic!("unexpected engine message: {other:?}"),
2457 };
2458
2459 response_tx.send(Ok(OnForkChoiceUpdated::syncing())).expect("send syncing response");
2460
2461 let response = api_task
2462 .await
2463 .expect("api task should not panic")
2464 .expect("forkchoiceUpdatedV3 should return a syncing response");
2465 assert!(response.payload_status.is_syncing());
2466 assert!(response.payload_id.is_none());
2467 }
2468
2469 #[tokio::test]
2470 async fn fcu_v3_valid_forkchoice_missing_beacon_root_returns_invalid_attributes() {
2471 let (mut handle, api) = setup_engine_api();
2472
2473 let state = ForkchoiceState {
2474 head_block_hash: B256::from([0x22; 32]),
2475 safe_block_hash: B256::ZERO,
2476 finalized_block_hash: B256::ZERO,
2477 };
2478 let payload_attributes = PayloadAttributes {
2479 timestamp: 1,
2480 prev_randao: B256::ZERO,
2481 suggested_fee_recipient: Address::ZERO,
2482 withdrawals: Some(vec![]),
2483 parent_beacon_block_root: None,
2484 slot_number: None,
2485 ..Default::default()
2486 };
2487
2488 let api_task = tokio::spawn(async move {
2489 api.fork_choice_updated_v3(state, Some(payload_attributes)).await
2490 });
2491
2492 let request =
2493 tokio::time::timeout(std::time::Duration::from_secs(1), handle.from_api.recv())
2494 .await
2495 .expect("timed out waiting for forkchoiceUpdated request")
2496 .expect("expected forkchoiceUpdated request");
2497
2498 let response_tx = match request {
2499 BeaconEngineMessage::ForkchoiceUpdated { payload_attrs, tx, .. } => {
2500 assert!(
2501 payload_attrs.is_none(),
2502 "when attrs are invalid, API should first evaluate forkchoice without attrs"
2503 );
2504 tx
2505 }
2506 other => panic!("unexpected engine message: {other:?}"),
2507 };
2508
2509 response_tx
2510 .send(Ok(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
2511 PayloadStatusEnum::Valid,
2512 ))))
2513 .expect("send valid response");
2514
2515 let response = api_task.await.expect("api task should not panic");
2516 assert_matches!(
2517 response,
2518 Err(EngineApiError::EngineObjectValidationError(
2519 reth_payload_primitives::EngineObjectValidationError::PayloadAttributes(_)
2520 ))
2521 );
2522
2523 match tokio::time::timeout(std::time::Duration::from_millis(100), handle.from_api.recv())
2524 .await
2525 {
2526 Err(_) | Ok(None) => {}
2527 Ok(Some(BeaconEngineMessage::ForkchoiceUpdated { .. })) => {
2528 panic!("no second forkchoiceUpdated call should be sent when attrs are invalid")
2529 }
2530 Ok(Some(other)) => panic!("unexpected engine message: {other:?}"),
2531 }
2532 }
2533
2534 mod get_payload_bodies {
2536 use super::*;
2537 use alloy_rpc_types_engine::ExecutionPayloadBodyV1;
2538 use reth_testing_utils::generators::{self, random_block_range, BlockRangeParams};
2539
2540 #[tokio::test]
2541 async fn invalid_params() {
2542 let (_, api) = setup_engine_api();
2543
2544 let by_range_tests = [
2545 (0, 0),
2547 (0, 1),
2548 (1, 0),
2549 ];
2550
2551 for (start, count) in by_range_tests {
2553 let res = api.get_payload_bodies_by_range_v1(start, count).await;
2554 assert_matches!(res, Err(EngineApiError::InvalidBodiesRange { .. }));
2555 }
2556 }
2557
2558 #[tokio::test]
2559 async fn request_too_large() {
2560 let (_, api) = setup_engine_api();
2561
2562 let request_count = MAX_PAYLOAD_BODIES_LIMIT + 1;
2563 let res = api.get_payload_bodies_by_range_v1(0, request_count).await;
2564 assert_matches!(res, Err(EngineApiError::PayloadRequestTooLarge { .. }));
2565 }
2566
2567 #[tokio::test]
2568 async fn returns_payload_bodies() {
2569 let mut rng = generators::rng();
2570 let (handle, api) = setup_engine_api();
2571
2572 let (start, count) = (1, 10);
2573 let blocks = random_block_range(
2574 &mut rng,
2575 start..=start + count - 1,
2576 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2577 );
2578 handle
2579 .provider
2580 .extend_blocks(blocks.iter().cloned().map(|b| (b.hash(), b.into_block())));
2581
2582 let expected = blocks
2583 .iter()
2584 .cloned()
2585 .map(|b| Some(ExecutionPayloadBodyV1::from_block(b.into_block())))
2586 .collect::<Vec<_>>();
2587
2588 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2589 assert_eq!(res, expected);
2590 }
2591
2592 #[tokio::test]
2593 async fn returns_payload_bodies_with_gaps() {
2594 let mut rng = generators::rng();
2595 let (handle, api) = setup_engine_api();
2596
2597 let (start, count) = (1, 100);
2598 let blocks = random_block_range(
2599 &mut rng,
2600 start..=start + count - 1,
2601 BlockRangeParams { tx_count: 0..2, ..Default::default() },
2602 );
2603
2604 let first_missing_range = 26..=50;
2606 let second_missing_range = 76..=100;
2607 handle.provider.extend_blocks(
2608 blocks
2609 .iter()
2610 .filter(|b| {
2611 !first_missing_range.contains(&b.number) &&
2612 !second_missing_range.contains(&b.number)
2613 })
2614 .map(|b| (b.hash(), b.clone().into_block())),
2615 );
2616
2617 let expected = blocks
2618 .iter()
2619 .filter(|b| !second_missing_range.contains(&b.number))
2622 .cloned()
2623 .map(|b| {
2624 if first_missing_range.contains(&b.number) {
2625 None
2626 } else {
2627 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2628 }
2629 })
2630 .collect::<Vec<_>>();
2631
2632 let res = api.get_payload_bodies_by_range_v1(start, count).await.unwrap();
2633 assert_eq!(res, expected);
2634
2635 let expected = blocks
2636 .iter()
2637 .cloned()
2638 .map(|b| {
2641 if first_missing_range.contains(&b.number) ||
2642 second_missing_range.contains(&b.number)
2643 {
2644 None
2645 } else {
2646 Some(ExecutionPayloadBodyV1::from_block(b.into_block()))
2647 }
2648 })
2649 .collect::<Vec<_>>();
2650
2651 let hashes = blocks.iter().map(|b| b.hash()).collect();
2652 let res = api.get_payload_bodies_by_hash_v1(hashes).await.unwrap();
2653 assert_eq!(res, expected);
2654 }
2655 }
2656}