1pub use jsonrpsee::{
4 core::middleware::layer::Either,
5 server::middleware::rpc::{RpcService, RpcServiceBuilder},
6};
7use reth_engine_tree::tree::WaitForCaches;
8pub use reth_engine_tree::tree::{BasicEngineValidator, EngineValidator};
9pub use reth_rpc_builder::{
10 middleware::{RethAuthHttpMiddleware, RethRpcMiddleware},
11 Identity, Stack,
12};
13use reth_storage_overlay::OverlayManager;
14
15use crate::{
16 invalid_block_hook::InvalidBlockHookExt, txpool_prewarm, ConfigureEngineEvm,
17 ConsensusEngineEvent, ConsensusEngineHandle,
18};
19use alloy_consensus::BlockHeader;
20use alloy_eips::BlockNumberOrTag;
21use alloy_rpc_types::engine::ClientVersionV1;
22use alloy_rpc_types_engine::ExecutionData;
23use futures::{Stream, StreamExt};
24use jsonrpsee::RpcModule;
25use parking_lot::Mutex;
26use reth_chain_state::{CanonStateNotification, CanonStateSubscriptions};
27use reth_chainspec::{ChainSpecProvider, EthChainSpec, EthereumHardforks, Hardforks};
28use reth_node_api::{
29 AddOnsContext, BlockTy, EngineApiValidator, EngineTypes, FullNodeComponents, FullNodeTypes,
30 NodeAddOns, NodeTypes, PayloadTypes, PayloadValidator, PrimitivesTy, TreeConfig,
31};
32use reth_node_core::{
33 cli::config::RethTransactionPoolConfig,
34 node_config::NodeConfig,
35 version::{version_metadata, CLIENT_CODE},
36};
37use reth_payload_builder::{PayloadBuilderHandle, PayloadStore};
38use reth_primitives_traits::NodePrimitives;
39use reth_rpc::{
40 eth::{core::EthRpcConverterFor, DevSigner, EthApiTypes, FullEthApiServer},
41 AdminApi,
42};
43use reth_rpc_api::{
44 eth::helpers::{EthTransactions, GetBlockAccessList},
45 IntoEngineApiRpcModule,
46};
47use reth_rpc_builder::{
48 auth::{AuthRpcModule, AuthServerHandle},
49 config::RethRpcServerConfig,
50 RpcModuleBuilder, RpcRegistryInner, RpcServerConfig, RpcServerHandle, TransportRpcModules,
51};
52use reth_rpc_engine_api::{capabilities::EngineCapabilities, EngineApi};
53use reth_rpc_eth_types::{cache::cache_new_blocks_task, EthConfig, EthStateCache};
54use reth_tokio_util::EventSender;
55use reth_tracing::tracing::{debug, info};
56use std::{
57 fmt::{self, Debug},
58 future::Future,
59 ops::{Deref, DerefMut},
60 sync::Arc,
61};
62use tokio::sync::oneshot;
63
64#[derive(Debug, Clone)]
68pub struct RethRpcServerHandles {
69 pub rpc: RpcServerHandle,
71 pub auth: AuthServerHandle,
73}
74
75pub struct RpcHooks<Node: FullNodeComponents, EthApi> {
77 pub on_rpc_started: Box<dyn OnRpcStarted<Node, EthApi>>,
79 pub extend_rpc_modules: Box<dyn ExtendRpcModules<Node, EthApi>>,
81}
82
83impl<Node, EthApi> Default for RpcHooks<Node, EthApi>
84where
85 Node: FullNodeComponents,
86 EthApi: EthApiTypes,
87{
88 fn default() -> Self {
89 Self { on_rpc_started: Box::<()>::default(), extend_rpc_modules: Box::<()>::default() }
90 }
91}
92
93impl<Node, EthApi> RpcHooks<Node, EthApi>
94where
95 Node: FullNodeComponents,
96 EthApi: EthApiTypes,
97{
98 pub(crate) fn set_on_rpc_started<F>(&mut self, hook: F) -> &mut Self
100 where
101 F: OnRpcStarted<Node, EthApi> + 'static,
102 {
103 self.on_rpc_started = Box::new(hook);
104 self
105 }
106
107 #[expect(unused)]
109 pub(crate) fn on_rpc_started<F>(mut self, hook: F) -> Self
110 where
111 F: OnRpcStarted<Node, EthApi> + 'static,
112 {
113 self.set_on_rpc_started(hook);
114 self
115 }
116
117 pub(crate) fn set_extend_rpc_modules<F>(&mut self, hook: F) -> &mut Self
119 where
120 F: ExtendRpcModules<Node, EthApi> + 'static,
121 {
122 self.extend_rpc_modules = Box::new(hook);
123 self
124 }
125
126 #[expect(unused)]
128 pub(crate) fn extend_rpc_modules<F>(mut self, hook: F) -> Self
129 where
130 F: ExtendRpcModules<Node, EthApi> + 'static,
131 {
132 self.set_extend_rpc_modules(hook);
133 self
134 }
135}
136
137impl<Node, EthApi> fmt::Debug for RpcHooks<Node, EthApi>
138where
139 Node: FullNodeComponents,
140 EthApi: EthApiTypes,
141{
142 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
143 f.debug_struct("RpcHooks")
144 .field("on_rpc_started", &"...")
145 .field("extend_rpc_modules", &"...")
146 .finish()
147 }
148}
149
150pub trait OnRpcStarted<Node: FullNodeComponents, EthApi: EthApiTypes>: Send {
152 fn on_rpc_started(
154 self: Box<Self>,
155 ctx: RpcContext<'_, Node, EthApi>,
156 handles: RethRpcServerHandles,
157 ) -> eyre::Result<()>;
158}
159
160impl<Node, EthApi, F> OnRpcStarted<Node, EthApi> for F
161where
162 F: FnOnce(RpcContext<'_, Node, EthApi>, RethRpcServerHandles) -> eyre::Result<()> + Send,
163 Node: FullNodeComponents,
164 EthApi: EthApiTypes,
165{
166 fn on_rpc_started(
167 self: Box<Self>,
168 ctx: RpcContext<'_, Node, EthApi>,
169 handles: RethRpcServerHandles,
170 ) -> eyre::Result<()> {
171 (*self)(ctx, handles)
172 }
173}
174
175impl<Node, EthApi> OnRpcStarted<Node, EthApi> for ()
176where
177 Node: FullNodeComponents,
178 EthApi: EthApiTypes,
179{
180 fn on_rpc_started(
181 self: Box<Self>,
182 _: RpcContext<'_, Node, EthApi>,
183 _: RethRpcServerHandles,
184 ) -> eyre::Result<()> {
185 Ok(())
186 }
187}
188
189pub trait ExtendRpcModules<Node: FullNodeComponents, EthApi: EthApiTypes>: Send {
191 fn extend_rpc_modules(self: Box<Self>, ctx: RpcContext<'_, Node, EthApi>) -> eyre::Result<()>;
193}
194
195impl<Node, EthApi, F> ExtendRpcModules<Node, EthApi> for F
196where
197 F: FnOnce(RpcContext<'_, Node, EthApi>) -> eyre::Result<()> + Send,
198 Node: FullNodeComponents,
199 EthApi: EthApiTypes,
200{
201 fn extend_rpc_modules(self: Box<Self>, ctx: RpcContext<'_, Node, EthApi>) -> eyre::Result<()> {
202 (*self)(ctx)
203 }
204}
205
206impl<Node, EthApi> ExtendRpcModules<Node, EthApi> for ()
207where
208 Node: FullNodeComponents,
209 EthApi: EthApiTypes,
210{
211 fn extend_rpc_modules(self: Box<Self>, _: RpcContext<'_, Node, EthApi>) -> eyre::Result<()> {
212 Ok(())
213 }
214}
215
216#[derive(Debug, Clone)]
218#[expect(clippy::type_complexity)]
219pub struct RpcRegistry<Node: FullNodeComponents, EthApi: EthApiTypes> {
220 pub(crate) registry: RpcRegistryInner<
221 Node::Provider,
222 Node::Pool,
223 Node::Network,
224 EthApi,
225 Node::Evm,
226 Node::Consensus,
227 >,
228}
229
230impl<Node, EthApi> Deref for RpcRegistry<Node, EthApi>
231where
232 Node: FullNodeComponents,
233 EthApi: EthApiTypes,
234{
235 type Target = RpcRegistryInner<
236 Node::Provider,
237 Node::Pool,
238 Node::Network,
239 EthApi,
240 Node::Evm,
241 Node::Consensus,
242 >;
243
244 fn deref(&self) -> &Self::Target {
245 &self.registry
246 }
247}
248
249impl<Node, EthApi> DerefMut for RpcRegistry<Node, EthApi>
250where
251 Node: FullNodeComponents,
252 EthApi: EthApiTypes,
253{
254 fn deref_mut(&mut self) -> &mut Self::Target {
255 &mut self.registry
256 }
257}
258
259#[expect(missing_debug_implementations)]
261pub struct RpcModuleContainer<'a, Node: FullNodeComponents, EthApi: EthApiTypes> {
262 pub modules: &'a mut TransportRpcModules,
264 pub auth_module: &'a mut AuthRpcModule,
266 pub registry: &'a mut RpcRegistry<Node, EthApi>,
268}
269
270#[expect(missing_debug_implementations)]
278pub struct RpcContext<'a, Node: FullNodeComponents, EthApi: EthApiTypes> {
279 pub(crate) node: Node,
281
282 pub(crate) config: &'a NodeConfig<<Node::Types as NodeTypes>::ChainSpec>,
284
285 pub registry: &'a mut RpcRegistry<Node, EthApi>,
289 pub modules: &'a mut TransportRpcModules,
294 pub auth_module: &'a mut AuthRpcModule,
298}
299
300impl<Node, EthApi> RpcContext<'_, Node, EthApi>
301where
302 Node: FullNodeComponents,
303 EthApi: EthApiTypes,
304{
305 pub const fn config(&self) -> &NodeConfig<<Node::Types as NodeTypes>::ChainSpec> {
307 self.config
308 }
309
310 pub const fn node(&self) -> &Node {
314 &self.node
315 }
316
317 pub fn pool(&self) -> &Node::Pool {
319 self.node.pool()
320 }
321
322 pub fn provider(&self) -> &Node::Provider {
324 self.node.provider()
325 }
326
327 pub fn network(&self) -> &Node::Network {
329 self.node.network()
330 }
331
332 pub fn payload_builder_handle(
334 &self,
335 ) -> &PayloadBuilderHandle<<Node::Types as NodeTypes>::Payload> {
336 self.node.payload_builder_handle()
337 }
338}
339
340pub struct RpcHandle<Node: FullNodeComponents, EthApi: EthApiTypes> {
342 pub rpc_server_handles: RethRpcServerHandles,
344 pub rpc_registry: RpcRegistry<Node, EthApi>,
346 pub engine_events: EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>>,
351 pub beacon_engine_handle: ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload>,
353 pub engine_shutdown: EngineShutdown,
355}
356
357impl<Node: FullNodeComponents, EthApi: EthApiTypes> Clone for RpcHandle<Node, EthApi> {
358 fn clone(&self) -> Self {
359 Self {
360 rpc_server_handles: self.rpc_server_handles.clone(),
361 rpc_registry: self.rpc_registry.clone(),
362 engine_events: self.engine_events.clone(),
363 beacon_engine_handle: self.beacon_engine_handle.clone(),
364 engine_shutdown: self.engine_shutdown.clone(),
365 }
366 }
367}
368
369impl<Node: FullNodeComponents, EthApi: EthApiTypes> Deref for RpcHandle<Node, EthApi> {
370 type Target = RpcRegistry<Node, EthApi>;
371
372 fn deref(&self) -> &Self::Target {
373 &self.rpc_registry
374 }
375}
376
377impl<Node: FullNodeComponents, EthApi: EthApiTypes> Debug for RpcHandle<Node, EthApi>
378where
379 RpcRegistry<Node, EthApi>: Debug,
380{
381 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
382 f.debug_struct("RpcHandle")
383 .field("rpc_server_handles", &self.rpc_server_handles)
384 .field("rpc_registry", &self.rpc_registry)
385 .field("engine_shutdown", &self.engine_shutdown)
386 .finish()
387 }
388}
389
390impl<Node: FullNodeComponents, EthApi: EthApiTypes> RpcHandle<Node, EthApi> {
391 pub const fn rpc_server_handles(&self) -> &RethRpcServerHandles {
393 &self.rpc_server_handles
394 }
395
396 pub const fn consensus_engine_handle(
400 &self,
401 ) -> &ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload> {
402 &self.beacon_engine_handle
403 }
404
405 pub const fn consensus_engine_events(
407 &self,
408 ) -> &EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>> {
409 &self.engine_events
410 }
411
412 pub const fn eth_api(&self) -> &EthApi {
414 self.rpc_registry.registry.eth_api()
415 }
416
417 pub fn admin_api(
419 &self,
420 ) -> AdminApi<Node::Network, <Node::Types as NodeTypes>::ChainSpec, Node::Pool>
421 where
422 <Node::Types as NodeTypes>::ChainSpec: EthereumHardforks,
423 {
424 self.rpc_registry.registry.admin_api()
425 }
426}
427
428#[derive(Debug, Clone)]
434pub struct RpcServerOnlyHandle<Node: FullNodeComponents, EthApi: EthApiTypes> {
435 pub rpc_server_handle: RpcServerHandle,
437 pub rpc_registry: RpcRegistry<Node, EthApi>,
439 pub engine_events: EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>>,
441 pub engine_handle: ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload>,
443}
444
445impl<Node: FullNodeComponents, EthApi: EthApiTypes> RpcServerOnlyHandle<Node, EthApi> {
446 pub const fn rpc_server_handle(&self) -> &RpcServerHandle {
448 &self.rpc_server_handle
449 }
450
451 pub const fn consensus_engine_handle(
455 &self,
456 ) -> &ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload> {
457 &self.engine_handle
458 }
459
460 pub const fn consensus_engine_events(
462 &self,
463 ) -> &EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>> {
464 &self.engine_events
465 }
466}
467
468#[derive(Debug, Clone)]
474pub struct AuthServerOnlyHandle<Node: FullNodeComponents, EthApi: EthApiTypes> {
475 pub auth_server_handle: AuthServerHandle,
477 pub rpc_registry: RpcRegistry<Node, EthApi>,
479 pub engine_events: EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>>,
481 pub engine_handle: ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload>,
483}
484
485impl<Node: FullNodeComponents, EthApi: EthApiTypes> AuthServerOnlyHandle<Node, EthApi> {
486 pub const fn consensus_engine_handle(
490 &self,
491 ) -> &ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload> {
492 &self.engine_handle
493 }
494
495 pub const fn consensus_engine_events(
497 &self,
498 ) -> &EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>> {
499 &self.engine_events
500 }
501}
502
503struct RpcSetupContext<'a, Node: FullNodeComponents, EthApi: EthApiTypes> {
505 node: Node,
506 config: &'a NodeConfig<<Node::Types as NodeTypes>::ChainSpec>,
507 modules: TransportRpcModules,
508 auth_module: AuthRpcModule,
509 auth_config: reth_rpc_builder::auth::AuthServerConfig,
510 registry: RpcRegistry<Node, EthApi>,
511 on_rpc_started: Box<dyn OnRpcStarted<Node, EthApi>>,
512 engine_events: EventSender<ConsensusEngineEvent<<Node::Types as NodeTypes>::Primitives>>,
513 engine_handle: ConsensusEngineHandle<<Node::Types as NodeTypes>::Payload>,
514}
515
516pub struct RpcAddOns<
527 Node: FullNodeComponents,
528 EthB: EthApiBuilder<Node>,
529 PVB,
530 EB = BasicEngineApiBuilder<PVB>,
531 EVB = BasicEngineValidatorBuilder<PVB>,
532 RpcMiddleware = Identity,
533 AuthHttpMiddleware = Identity,
534> {
535 pub hooks: RpcHooks<Node, EthB::EthApi>,
537 eth_api_builder: EthB,
539 payload_validator_builder: PVB,
541 engine_api_builder: EB,
543 engine_validator_builder: EVB,
545 rpc_middleware: RpcMiddleware,
550 auth_http_middleware: AuthHttpMiddleware,
555 tokio_runtime: Option<tokio::runtime::Handle>,
557}
558
559impl<Node, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware> Debug
560 for RpcAddOns<Node, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
561where
562 Node: FullNodeComponents,
563 EthB: EthApiBuilder<Node>,
564 PVB: Debug,
565 EB: Debug,
566 EVB: Debug,
567{
568 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
569 f.debug_struct("RpcAddOns")
570 .field("hooks", &self.hooks)
571 .field("eth_api_builder", &"...")
572 .field("payload_validator_builder", &self.payload_validator_builder)
573 .field("engine_api_builder", &self.engine_api_builder)
574 .field("engine_validator_builder", &self.engine_validator_builder)
575 .field("rpc_middleware", &"...")
576 .finish()
577 }
578}
579
580impl<Node, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
581 RpcAddOns<Node, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
582where
583 Node: FullNodeComponents,
584 EthB: EthApiBuilder<Node>,
585{
586 pub fn new(
588 eth_api_builder: EthB,
589 payload_validator_builder: PVB,
590 engine_api_builder: EB,
591 engine_validator_builder: EVB,
592 rpc_middleware: RpcMiddleware,
593 auth_http_middleware: AuthHttpMiddleware,
594 ) -> Self {
595 Self {
596 hooks: RpcHooks::default(),
597 eth_api_builder,
598 payload_validator_builder,
599 engine_api_builder,
600 engine_validator_builder,
601 rpc_middleware,
602 auth_http_middleware,
603 tokio_runtime: None,
604 }
605 }
606
607 pub fn with_engine_api<T>(
609 self,
610 engine_api_builder: T,
611 ) -> RpcAddOns<Node, EthB, PVB, T, EVB, RpcMiddleware, AuthHttpMiddleware> {
612 self.map_engine_api(|_| engine_api_builder)
613 }
614
615 pub fn map_engine_api<T>(
617 self,
618 f: impl FnOnce(EB) -> T,
619 ) -> RpcAddOns<Node, EthB, PVB, T, EVB, RpcMiddleware, AuthHttpMiddleware> {
620 let Self {
621 hooks,
622 eth_api_builder,
623 payload_validator_builder,
624 engine_api_builder,
625 engine_validator_builder,
626 rpc_middleware,
627 auth_http_middleware,
628 tokio_runtime,
629 } = self;
630 RpcAddOns {
631 hooks,
632 eth_api_builder,
633 payload_validator_builder,
634 engine_api_builder: f(engine_api_builder),
635 engine_validator_builder,
636 rpc_middleware,
637 auth_http_middleware,
638 tokio_runtime,
639 }
640 }
641
642 pub fn with_payload_validator<T>(
644 self,
645 payload_validator_builder: T,
646 ) -> RpcAddOns<Node, EthB, T, EB, EVB, RpcMiddleware, AuthHttpMiddleware> {
647 let Self {
648 hooks,
649 eth_api_builder,
650 engine_api_builder,
651 engine_validator_builder,
652 rpc_middleware,
653 auth_http_middleware,
654 tokio_runtime,
655 ..
656 } = self;
657 RpcAddOns {
658 hooks,
659 eth_api_builder,
660 payload_validator_builder,
661 engine_api_builder,
662 engine_validator_builder,
663 rpc_middleware,
664 auth_http_middleware,
665 tokio_runtime,
666 }
667 }
668
669 pub fn with_engine_validator<T>(
671 self,
672 engine_validator_builder: T,
673 ) -> RpcAddOns<Node, EthB, PVB, EB, T, RpcMiddleware, AuthHttpMiddleware> {
674 let Self {
675 hooks,
676 eth_api_builder,
677 payload_validator_builder,
678 engine_api_builder,
679 rpc_middleware,
680 auth_http_middleware,
681 tokio_runtime,
682 ..
683 } = self;
684 RpcAddOns {
685 hooks,
686 eth_api_builder,
687 payload_validator_builder,
688 engine_api_builder,
689 engine_validator_builder,
690 rpc_middleware,
691 auth_http_middleware,
692 tokio_runtime,
693 }
694 }
695
696 pub fn with_rpc_middleware<T>(
735 self,
736 rpc_middleware: T,
737 ) -> RpcAddOns<Node, EthB, PVB, EB, EVB, T, AuthHttpMiddleware> {
738 let Self {
739 hooks,
740 eth_api_builder,
741 payload_validator_builder,
742 engine_api_builder,
743 engine_validator_builder,
744 auth_http_middleware,
745 tokio_runtime,
746 ..
747 } = self;
748 RpcAddOns {
749 hooks,
750 eth_api_builder,
751 payload_validator_builder,
752 engine_api_builder,
753 engine_validator_builder,
754 rpc_middleware,
755 auth_http_middleware,
756 tokio_runtime,
757 }
758 }
759
760 pub fn with_auth_http_middleware<T>(
765 self,
766 auth_http_middleware: T,
767 ) -> RpcAddOns<Node, EthB, PVB, EB, EVB, RpcMiddleware, T> {
768 let Self {
769 hooks,
770 eth_api_builder,
771 payload_validator_builder,
772 engine_api_builder,
773 engine_validator_builder,
774 rpc_middleware,
775 tokio_runtime,
776 ..
777 } = self;
778 RpcAddOns {
779 hooks,
780 eth_api_builder,
781 payload_validator_builder,
782 engine_api_builder,
783 engine_validator_builder,
784 rpc_middleware,
785 auth_http_middleware,
786 tokio_runtime,
787 }
788 }
789
790 pub fn layer_auth_http_middleware<T>(
792 self,
793 layer: T,
794 ) -> RpcAddOns<Node, EthB, PVB, EB, EVB, RpcMiddleware, Stack<AuthHttpMiddleware, T>> {
795 let Self {
796 hooks,
797 eth_api_builder,
798 payload_validator_builder,
799 engine_api_builder,
800 engine_validator_builder,
801 rpc_middleware,
802 auth_http_middleware,
803 tokio_runtime,
804 } = self;
805 let auth_http_middleware = Stack::new(auth_http_middleware, layer);
806 RpcAddOns {
807 hooks,
808 eth_api_builder,
809 payload_validator_builder,
810 engine_api_builder,
811 engine_validator_builder,
812 rpc_middleware,
813 auth_http_middleware,
814 tokio_runtime,
815 }
816 }
817
818 pub fn map_auth_http_middleware<T>(
820 self,
821 f: impl FnOnce(AuthHttpMiddleware) -> T,
822 ) -> RpcAddOns<Node, EthB, PVB, EB, EVB, RpcMiddleware, T> {
823 let Self {
824 hooks,
825 eth_api_builder,
826 payload_validator_builder,
827 engine_api_builder,
828 engine_validator_builder,
829 rpc_middleware,
830 auth_http_middleware,
831 tokio_runtime,
832 } = self;
833 RpcAddOns {
834 hooks,
835 eth_api_builder,
836 payload_validator_builder,
837 engine_api_builder,
838 engine_validator_builder,
839 rpc_middleware,
840 auth_http_middleware: f(auth_http_middleware),
841 tokio_runtime,
842 }
843 }
844
845 #[expect(clippy::type_complexity)]
847 pub fn option_layer_auth_http_middleware<T>(
848 self,
849 layer: Option<T>,
850 ) -> RpcAddOns<
851 Node,
852 EthB,
853 PVB,
854 EB,
855 EVB,
856 RpcMiddleware,
857 Stack<AuthHttpMiddleware, Either<T, Identity>>,
858 > {
859 let layer = layer.map(Either::Left).unwrap_or(Either::Right(Identity::new()));
860 self.layer_auth_http_middleware(layer)
861 }
862
863 pub fn with_tokio_runtime(self, tokio_runtime: Option<tokio::runtime::Handle>) -> Self {
867 let Self {
868 hooks,
869 eth_api_builder,
870 payload_validator_builder,
871 engine_validator_builder,
872 engine_api_builder,
873 rpc_middleware,
874 auth_http_middleware,
875 ..
876 } = self;
877 Self {
878 hooks,
879 eth_api_builder,
880 payload_validator_builder,
881 engine_validator_builder,
882 engine_api_builder,
883 rpc_middleware,
884 auth_http_middleware,
885 tokio_runtime,
886 }
887 }
888
889 pub fn layer_rpc_middleware<T>(
891 self,
892 layer: T,
893 ) -> RpcAddOns<Node, EthB, PVB, EB, EVB, Stack<RpcMiddleware, T>, AuthHttpMiddleware> {
894 let Self {
895 hooks,
896 eth_api_builder,
897 payload_validator_builder,
898 engine_api_builder,
899 engine_validator_builder,
900 rpc_middleware,
901 auth_http_middleware,
902 tokio_runtime,
903 } = self;
904 let rpc_middleware = Stack::new(rpc_middleware, layer);
905 RpcAddOns {
906 hooks,
907 eth_api_builder,
908 payload_validator_builder,
909 engine_api_builder,
910 engine_validator_builder,
911 rpc_middleware,
912 auth_http_middleware,
913 tokio_runtime,
914 }
915 }
916
917 #[expect(clippy::type_complexity)]
919 pub fn option_layer_rpc_middleware<T>(
920 self,
921 layer: Option<T>,
922 ) -> RpcAddOns<
923 Node,
924 EthB,
925 PVB,
926 EB,
927 EVB,
928 Stack<RpcMiddleware, Either<T, Identity>>,
929 AuthHttpMiddleware,
930 > {
931 let layer = layer.map(Either::Left).unwrap_or(Either::Right(Identity::new()));
932 self.layer_rpc_middleware(layer)
933 }
934
935 pub fn on_rpc_started<F>(mut self, hook: F) -> Self
937 where
938 F: FnOnce(RpcContext<'_, Node, EthB::EthApi>, RethRpcServerHandles) -> eyre::Result<()>
939 + Send
940 + 'static,
941 {
942 self.hooks.set_on_rpc_started(hook);
943 self
944 }
945
946 pub fn extend_rpc_modules<F>(mut self, hook: F) -> Self
948 where
949 F: FnOnce(RpcContext<'_, Node, EthB::EthApi>) -> eyre::Result<()> + Send + 'static,
950 {
951 self.hooks.set_extend_rpc_modules(hook);
952 self
953 }
954}
955
956impl<Node, EthB, EV, EB, Engine> Default
957 for RpcAddOns<Node, EthB, EV, EB, Engine, Identity, Identity>
958where
959 Node: FullNodeComponents,
960 EthB: EthApiBuilder<Node>,
961 EV: Default,
962 EB: Default,
963 Engine: Default,
964{
965 fn default() -> Self {
966 Self::new(
967 EthB::default(),
968 EV::default(),
969 EB::default(),
970 Engine::default(),
971 Default::default(),
972 Identity::new(),
973 )
974 }
975}
976
977impl<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
978 RpcAddOns<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
979where
980 N: FullNodeComponents,
981 N::Provider: ChainSpecProvider<ChainSpec: EthereumHardforks>,
982 EthB: EthApiBuilder<N>,
983 EB: EngineApiBuilder<N>,
984 EVB: EngineValidatorBuilder<N>,
985 RpcMiddleware: RethRpcMiddleware,
986 AuthHttpMiddleware: RethAuthHttpMiddleware<Identity>,
987{
988 pub async fn launch_rpc_server<F>(
994 self,
995 ctx: AddOnsContext<'_, N>,
996 ext: F,
997 ) -> eyre::Result<RpcServerOnlyHandle<N, EthB::EthApi>>
998 where
999 F: FnOnce(RpcModuleContainer<'_, N, EthB::EthApi>) -> eyre::Result<()>,
1000 {
1001 let rpc_middleware = self.rpc_middleware.clone();
1002 let tokio_runtime = self.tokio_runtime.clone();
1003 let setup_ctx = self.setup_rpc_components(ctx, ext).await?;
1004 let RpcSetupContext {
1005 node,
1006 config,
1007 mut modules,
1008 mut auth_module,
1009 auth_config: _,
1010 mut registry,
1011 on_rpc_started,
1012 engine_events,
1013 engine_handle,
1014 } = setup_ctx;
1015
1016 let server_config = config
1017 .rpc
1018 .rpc_server_config()
1019 .set_rpc_middleware(rpc_middleware)
1020 .with_tokio_runtime(tokio_runtime);
1021 let rpc_server_handle = Self::launch_rpc_server_internal(server_config, &modules).await?;
1022
1023 let handles =
1024 RethRpcServerHandles { rpc: rpc_server_handle.clone(), auth: AuthServerHandle::noop() };
1025 Self::finalize_rpc_setup(
1026 &mut registry,
1027 &mut modules,
1028 &mut auth_module,
1029 &node,
1030 config,
1031 on_rpc_started,
1032 handles,
1033 )?;
1034
1035 Ok(RpcServerOnlyHandle {
1036 rpc_server_handle,
1037 rpc_registry: registry,
1038 engine_events,
1039 engine_handle,
1040 })
1041 }
1042
1043 pub async fn launch_add_ons_with<F>(
1046 self,
1047 ctx: AddOnsContext<'_, N>,
1048 ext: F,
1049 ) -> eyre::Result<RpcHandle<N, EthB::EthApi>>
1050 where
1051 F: FnOnce(RpcModuleContainer<'_, N, EthB::EthApi>) -> eyre::Result<()>,
1052 {
1053 let disable_auth = ctx.config.rpc.disable_auth_server;
1055 self.launch_add_ons_with_opt_engine(ctx, ext, disable_auth).await
1056 }
1057
1058 pub async fn launch_add_ons_with_opt_engine<F>(
1064 self,
1065 ctx: AddOnsContext<'_, N>,
1066 ext: F,
1067 disable_auth: bool,
1068 ) -> eyre::Result<RpcHandle<N, EthB::EthApi>>
1069 where
1070 F: FnOnce(RpcModuleContainer<'_, N, EthB::EthApi>) -> eyre::Result<()>,
1071 {
1072 let rpc_middleware = self.rpc_middleware.clone();
1073 let auth_http_middleware = self.auth_http_middleware.clone();
1074 let tokio_runtime = self.tokio_runtime.clone();
1075 let setup_ctx = self.setup_rpc_components(ctx, ext).await?;
1076 let RpcSetupContext {
1077 node,
1078 config,
1079 mut modules,
1080 mut auth_module,
1081 auth_config,
1082 mut registry,
1083 on_rpc_started,
1084 engine_events,
1085 engine_handle,
1086 } = setup_ctx;
1087
1088 let server_config = config
1089 .rpc
1090 .rpc_server_config()
1091 .set_rpc_middleware(rpc_middleware)
1092 .with_tokio_runtime(tokio_runtime);
1093
1094 let auth_config = auth_config.with_http_middleware(auth_http_middleware);
1095
1096 let (rpc, auth) = if disable_auth {
1097 let rpc = Self::launch_rpc_server_internal(server_config, &modules).await?;
1099 (rpc, AuthServerHandle::noop())
1100 } else {
1101 let auth_module_clone = auth_module.clone();
1102 let (rpc, auth) = futures::future::try_join(
1104 Self::launch_rpc_server_internal(server_config, &modules),
1105 Self::launch_auth_server_internal(auth_config.start(auth_module_clone)),
1106 )
1107 .await?;
1108 (rpc, auth)
1109 };
1110
1111 let handles = RethRpcServerHandles { rpc, auth };
1112
1113 Self::finalize_rpc_setup(
1114 &mut registry,
1115 &mut modules,
1116 &mut auth_module,
1117 &node,
1118 config,
1119 on_rpc_started,
1120 handles.clone(),
1121 )?;
1122
1123 Ok(RpcHandle {
1124 rpc_server_handles: handles,
1125 rpc_registry: registry,
1126 engine_events,
1127 beacon_engine_handle: engine_handle,
1128 engine_shutdown: EngineShutdown::default(),
1129 })
1130 }
1131
1132 async fn setup_rpc_components<'a, F>(
1134 self,
1135 ctx: AddOnsContext<'a, N>,
1136 ext: F,
1137 ) -> eyre::Result<RpcSetupContext<'a, N, EthB::EthApi>>
1138 where
1139 F: FnOnce(RpcModuleContainer<'_, N, EthB::EthApi>) -> eyre::Result<()>,
1140 {
1141 let Self { eth_api_builder, engine_api_builder, hooks, .. } = self;
1142
1143 let engine_api = engine_api_builder.build_engine_api(&ctx).await?;
1144 let AddOnsContext {
1145 node,
1146 config,
1147 beacon_engine_handle,
1148 jwt_secret,
1149 engine_events,
1150 sender_recovery_cache,
1151 } = ctx;
1152
1153 info!(target: "reth::cli", "Engine API handler initialized");
1154
1155 let cache = EthStateCache::spawn_with(
1156 node.provider().clone(),
1157 config.rpc.eth_config().cache,
1158 node.task_executor().clone(),
1159 );
1160
1161 let new_canonical_blocks = node.provider().canonical_state_stream();
1162 let c = cache.clone();
1163 node.task_executor().spawn_critical_task("cache canonical blocks task", async move {
1164 cache_new_blocks_task(c, new_canonical_blocks).await;
1165 });
1166
1167 let prewarm_bals = config
1168 .rpc
1169 .eth_config()
1170 .cache
1171 .prewarm_bals
1172 .map(|count| (count, node.provider().canonical_state_stream()));
1173
1174 let eth_config = config.rpc.eth_config().max_batch_size(config.txpool.max_batch_size());
1175 let ctx = EthApiCtx {
1176 components: &node,
1177 config: eth_config,
1178 cache,
1179 engine_handle: beacon_engine_handle.clone(),
1180 sender_recovery_cache,
1181 };
1182 let eth_api = eth_api_builder.build_eth_api(ctx).await?;
1183
1184 if let Some((count, events)) = prewarm_bals {
1185 node.task_executor().spawn_task(prewarm_new_block_bals_task(
1186 eth_api.clone(),
1187 events,
1188 count,
1189 ));
1190 }
1191
1192 let auth_config = config.rpc.auth_server_config(jwt_secret)?;
1193 let module_config = config.rpc.transport_rpc_module_config();
1194 debug!(target: "reth::cli", http=?module_config.http(), ws=?module_config.ws(), "Using RPC module config");
1195
1196 let (mut modules, mut auth_module, registry) = RpcModuleBuilder::default()
1197 .with_provider(node.provider().clone())
1198 .with_pool(node.pool().clone())
1199 .with_network(node.network().clone())
1200 .with_executor(node.task_executor().clone())
1201 .with_evm_config(node.evm_config().clone())
1202 .with_consensus(node.consensus().clone())
1203 .build_with_auth_server(
1204 module_config,
1205 engine_api,
1206 eth_api,
1207 engine_events.clone(),
1208 beacon_engine_handle.clone(),
1209 );
1210
1211 if config.dev.dev {
1213 let signers = DevSigner::from_mnemonic(config.dev.dev_mnemonic.as_str(), 20);
1214 registry.eth_api().signers().write().extend(signers);
1215 }
1216
1217 let mut registry = RpcRegistry { registry };
1218 let ctx = RpcContext {
1219 node: node.clone(),
1220 config,
1221 registry: &mut registry,
1222 modules: &mut modules,
1223 auth_module: &mut auth_module,
1224 };
1225
1226 let RpcHooks { on_rpc_started, extend_rpc_modules } = hooks;
1227
1228 ext(RpcModuleContainer {
1229 modules: ctx.modules,
1230 auth_module: ctx.auth_module,
1231 registry: ctx.registry,
1232 })?;
1233 extend_rpc_modules.extend_rpc_modules(ctx)?;
1234
1235 Ok(RpcSetupContext {
1236 node,
1237 config,
1238 modules,
1239 auth_module,
1240 auth_config,
1241 registry,
1242 on_rpc_started,
1243 engine_events,
1244 engine_handle: beacon_engine_handle,
1245 })
1246 }
1247
1248 async fn launch_rpc_server_internal<M>(
1250 server_config: RpcServerConfig<M>,
1251 modules: &TransportRpcModules,
1252 ) -> eyre::Result<RpcServerHandle>
1253 where
1254 M: RethRpcMiddleware,
1255 {
1256 let handle = server_config.start(modules).await?;
1257
1258 if let Some(path) = handle.ipc_endpoint() {
1259 info!(target: "reth::cli", %path, "RPC IPC server started");
1260 }
1261 if let Some(addr) = handle.http_local_addr() {
1262 info!(target: "reth::cli", url=%addr, "RPC HTTP server started");
1263 }
1264 if let Some(addr) = handle.ws_local_addr() {
1265 info!(target: "reth::cli", url=%addr, "RPC WS server started");
1266 }
1267
1268 Ok(handle)
1269 }
1270
1271 async fn launch_auth_server_internal(
1273 start_fut: impl Future<Output = Result<AuthServerHandle, reth_rpc_builder::error::RpcError>>,
1274 ) -> eyre::Result<AuthServerHandle> {
1275 start_fut
1276 .await
1277 .map_err(Into::into)
1278 .inspect(|handle| {
1279 let addr = handle.local_addr();
1280 if let Some(ipc_endpoint) = handle.ipc_endpoint() {
1281 info!(target: "reth::cli", url=%addr, ipc_endpoint=%ipc_endpoint, "RPC auth server started");
1282 } else {
1283 info!(target: "reth::cli", url=%addr, "RPC auth server started");
1284 }
1285 })
1286 }
1287
1288 fn finalize_rpc_setup(
1290 registry: &mut RpcRegistry<N, EthB::EthApi>,
1291 modules: &mut TransportRpcModules,
1292 auth_module: &mut AuthRpcModule,
1293 node: &N,
1294 config: &NodeConfig<<N::Types as NodeTypes>::ChainSpec>,
1295 on_rpc_started: Box<dyn OnRpcStarted<N, EthB::EthApi>>,
1296 handles: RethRpcServerHandles,
1297 ) -> eyre::Result<()> {
1298 let ctx = RpcContext { node: node.clone(), config, registry, modules, auth_module };
1299
1300 on_rpc_started.on_rpc_started(ctx, handles)?;
1301 Ok(())
1302 }
1303}
1304
1305impl<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware> NodeAddOns<N>
1306 for RpcAddOns<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
1307where
1308 N: FullNodeComponents,
1309 <N as FullNodeTypes>::Provider: ChainSpecProvider<ChainSpec: EthereumHardforks>,
1310 EthB: EthApiBuilder<N>,
1311 PVB: PayloadValidatorBuilder<N>,
1312 EB: EngineApiBuilder<N>,
1313 EVB: EngineValidatorBuilder<N>,
1314 RpcMiddleware: RethRpcMiddleware,
1315 AuthHttpMiddleware: RethAuthHttpMiddleware<Identity>,
1316{
1317 type Handle = RpcHandle<N, EthB::EthApi>;
1318
1319 async fn launch_add_ons(self, ctx: AddOnsContext<'_, N>) -> eyre::Result<Self::Handle> {
1320 self.launch_add_ons_with(ctx, |_| Ok(())).await
1321 }
1322}
1323
1324pub trait RethRpcAddOns<N: FullNodeComponents>:
1327 NodeAddOns<N, Handle = RpcHandle<N, Self::EthApi>>
1328{
1329 type EthApi: EthApiTypes;
1331
1332 fn hooks_mut(&mut self) -> &mut RpcHooks<N, Self::EthApi>;
1334}
1335
1336impl<N: FullNodeComponents, EthB, EV, EB, Engine, RpcMiddleware, AuthHttpMiddleware>
1337 RethRpcAddOns<N> for RpcAddOns<N, EthB, EV, EB, Engine, RpcMiddleware, AuthHttpMiddleware>
1338where
1339 Self: NodeAddOns<N, Handle = RpcHandle<N, EthB::EthApi>>,
1340 EthB: EthApiBuilder<N>,
1341{
1342 type EthApi = EthB::EthApi;
1343
1344 fn hooks_mut(&mut self) -> &mut RpcHooks<N, Self::EthApi> {
1345 &mut self.hooks
1346 }
1347}
1348
1349#[derive(Debug)]
1352pub struct EthApiCtx<'a, N: FullNodeTypes> {
1353 pub components: &'a N,
1355 pub config: EthConfig,
1357 pub cache: EthStateCache<PrimitivesTy<N::Types>>,
1359 pub sender_recovery_cache: Option<reth_evm::SenderRecoveryCache>,
1361 pub engine_handle: ConsensusEngineHandle<<N::Types as NodeTypes>::Payload>,
1363}
1364
1365impl<'a, N: FullNodeComponents<Types: NodeTypes<ChainSpec: Hardforks + EthereumHardforks>>>
1366 EthApiCtx<'a, N>
1367{
1368 pub fn eth_api_builder(self) -> reth_rpc::EthApiBuilder<N, EthRpcConverterFor<N>> {
1370 reth_rpc::EthApiBuilder::new_with_components(self.components.clone())
1371 .eth_cache(self.cache)
1372 .sender_recovery_cache(self.sender_recovery_cache)
1373 .eth_state_cache_config(self.config.cache)
1374 .task_spawner(self.components.task_executor().clone())
1375 .gas_cap(self.config.rpc_gas_cap.into())
1376 .max_simulate_blocks(self.config.rpc_max_simulate_blocks)
1377 .compute_state_root_for_eth_simulate(self.config.compute_state_root_for_eth_simulate)
1378 .eth_proof_window(self.config.eth_proof_window)
1379 .fee_history_cache_config(self.config.fee_history_cache)
1380 .proof_permits(self.config.proof_permits)
1381 .gas_oracle_config(self.config.gas_oracle)
1382 .max_batch_size(self.config.max_batch_size)
1383 .max_blocking_io_requests(self.config.max_blocking_io_requests)
1384 .pending_block_kind(self.config.pending_block_kind)
1385 .raw_tx_forwarder(self.config.raw_tx_forwarder)
1386 .evm_memory_limit(self.config.rpc_evm_memory_limit)
1387 .force_blob_sidecar_upcasting(self.config.force_blob_sidecar_upcasting)
1388 }
1389}
1390
1391pub trait EthApiBuilder<N: FullNodeComponents>: Default + Send + 'static {
1393 type EthApi: FullEthApiServer<Provider = N::Provider, Pool = N::Pool>;
1395
1396 fn build_eth_api(
1398 self,
1399 ctx: EthApiCtx<'_, N>,
1400 ) -> impl Future<Output = eyre::Result<Self::EthApi>> + Send;
1401}
1402
1403pub trait EngineValidatorAddOn<Node: FullNodeComponents>: Send {
1405 type ValidatorBuilder: EngineValidatorBuilder<Node>;
1407
1408 fn engine_validator_builder(&self) -> Self::ValidatorBuilder;
1410}
1411
1412impl<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware> EngineValidatorAddOn<N>
1413 for RpcAddOns<N, EthB, PVB, EB, EVB, RpcMiddleware, AuthHttpMiddleware>
1414where
1415 N: FullNodeComponents,
1416 EthB: EthApiBuilder<N>,
1417 PVB: Send,
1418 EB: EngineApiBuilder<N>,
1419 EVB: EngineValidatorBuilder<N>,
1420 RpcMiddleware: Send,
1421 AuthHttpMiddleware: Send,
1422{
1423 type ValidatorBuilder = EVB;
1424
1425 fn engine_validator_builder(&self) -> Self::ValidatorBuilder {
1426 self.engine_validator_builder.clone()
1427 }
1428}
1429
1430pub trait EngineApiBuilder<Node: FullNodeComponents>: Send + Sync {
1437 type EngineApi: IntoEngineApiRpcModule + Send + Sync;
1439
1440 fn build_engine_api(
1445 self,
1446 ctx: &AddOnsContext<'_, Node>,
1447 ) -> impl Future<Output = eyre::Result<Self::EngineApi>> + Send;
1448}
1449
1450pub trait PayloadValidatorBuilder<Node: FullNodeComponents>: Send + Sync + Clone {
1455 type Validator: PayloadValidator<<Node::Types as NodeTypes>::Payload>;
1457
1458 fn build(
1463 self,
1464 ctx: &AddOnsContext<'_, Node>,
1465 ) -> impl Future<Output = eyre::Result<Self::Validator>> + Send;
1466}
1467
1468pub trait EngineValidatorBuilder<Node: FullNodeComponents>: Send + Sync + Clone {
1473 type EngineValidator: EngineValidator<<Node::Types as NodeTypes>::Payload, <Node::Types as NodeTypes>::Primitives>
1475 + WaitForCaches;
1476
1477 fn build_tree_validator(
1481 self,
1482 ctx: &AddOnsContext<'_, Node>,
1483 tree_config: TreeConfig,
1484 overlay_manager: OverlayManager<PrimitivesTy<Node::Types>>,
1485 ) -> impl Future<Output = eyre::Result<Self::EngineValidator>> + Send;
1486}
1487
1488#[derive(Debug, Clone)]
1492pub struct BasicEngineValidatorBuilder<EV> {
1493 payload_validator_builder: EV,
1495}
1496
1497impl<EV> BasicEngineValidatorBuilder<EV> {
1498 pub const fn new(payload_validator_builder: EV) -> Self {
1500 Self { payload_validator_builder }
1501 }
1502}
1503
1504impl<EV> Default for BasicEngineValidatorBuilder<EV>
1505where
1506 EV: Default,
1507{
1508 fn default() -> Self {
1509 Self::new(EV::default())
1510 }
1511}
1512
1513impl<Node, EV> EngineValidatorBuilder<Node> for BasicEngineValidatorBuilder<EV>
1514where
1515 Node: FullNodeComponents<
1516 Evm: ConfigureEngineEvm<
1517 <<Node::Types as NodeTypes>::Payload as PayloadTypes>::ExecutionData,
1518 >,
1519 >,
1520 EV: PayloadValidatorBuilder<Node>,
1521 EV::Validator: reth_engine_primitives::PayloadValidator<
1522 <Node::Types as NodeTypes>::Payload,
1523 Block = BlockTy<Node::Types>,
1524 > + Clone,
1525{
1526 type EngineValidator = BasicEngineValidator<Node::Provider, Node::Evm, EV::Validator>;
1527
1528 async fn build_tree_validator(
1529 self,
1530 ctx: &AddOnsContext<'_, Node>,
1531 tree_config: TreeConfig,
1532 overlay_manager: OverlayManager<PrimitivesTy<Node::Types>>,
1533 ) -> eyre::Result<Self::EngineValidator> {
1534 let validator = self.payload_validator_builder.build(ctx).await?;
1535 let data_dir = ctx.config.datadir.clone().resolve_datadir(ctx.config.chain.chain());
1536 let invalid_block_hook = ctx.create_invalid_block_hook(&data_dir).await?;
1537
1538 let txpool_prewarming = tree_config.txpool_prewarming();
1539 let mut validator = BasicEngineValidator::new(
1540 ctx.node.provider().clone(),
1541 std::sync::Arc::new(ctx.node.consensus().clone()),
1542 ctx.node.evm_config().clone(),
1543 validator,
1544 tree_config,
1545 invalid_block_hook,
1546 overlay_manager,
1547 ctx.node.task_executor().clone(),
1548 );
1549
1550 if txpool_prewarming {
1551 validator = validator
1552 .with_txpool_prewarming(txpool_prewarm::Source::new(ctx.node.pool().clone()));
1553 }
1554
1555 Ok(validator)
1556 }
1557}
1558
1559#[derive(Debug, Default)]
1565pub struct BasicEngineApiBuilder<PVB> {
1566 payload_validator_builder: PVB,
1567}
1568
1569impl<N, PVB> EngineApiBuilder<N> for BasicEngineApiBuilder<PVB>
1570where
1571 N: FullNodeComponents<
1572 Types: NodeTypes<
1573 ChainSpec: EthereumHardforks,
1574 Payload: PayloadTypes<ExecutionData = ExecutionData> + EngineTypes,
1575 >,
1576 >,
1577 PVB: PayloadValidatorBuilder<N>,
1578 PVB::Validator: EngineApiValidator<<N::Types as NodeTypes>::Payload>,
1579{
1580 type EngineApi = EngineApi<
1581 N::Provider,
1582 <N::Types as NodeTypes>::Payload,
1583 N::Pool,
1584 PVB::Validator,
1585 <N::Types as NodeTypes>::ChainSpec,
1586 >;
1587
1588 async fn build_engine_api(self, ctx: &AddOnsContext<'_, N>) -> eyre::Result<Self::EngineApi> {
1589 let Self { payload_validator_builder } = self;
1590
1591 let engine_validator = payload_validator_builder.build(ctx).await?;
1592 let client = ClientVersionV1 {
1593 code: CLIENT_CODE,
1594 name: version_metadata().name_client.to_string(),
1595 version: version_metadata().cargo_pkg_version.to_string(),
1596 commit: version_metadata().vergen_git_sha.to_string(),
1597 };
1598
1599 Ok(EngineApi::new(
1600 ctx.node.provider().clone(),
1601 ctx.config.chain.clone(),
1602 ctx.beacon_engine_handle.clone(),
1603 PayloadStore::new(ctx.node.payload_builder_handle().clone()),
1604 ctx.node.pool().clone(),
1605 ctx.node.task_executor().clone(),
1606 client,
1607 EngineCapabilities::default(),
1608 engine_validator,
1609 ctx.config.engine.accept_execution_requests_hash,
1610 ctx.node.network().clone(),
1611 ))
1612 }
1613}
1614
1615#[derive(Debug, Clone, Default)]
1621#[non_exhaustive]
1622pub struct NoopEngineApiBuilder;
1623
1624impl<N: FullNodeComponents> EngineApiBuilder<N> for NoopEngineApiBuilder {
1625 type EngineApi = NoopEngineApi;
1626
1627 async fn build_engine_api(self, _ctx: &AddOnsContext<'_, N>) -> eyre::Result<Self::EngineApi> {
1628 Ok(NoopEngineApi::default())
1629 }
1630}
1631
1632#[derive(Debug, Clone, Default)]
1637#[non_exhaustive]
1638pub struct NoopEngineApi;
1639
1640impl IntoEngineApiRpcModule for NoopEngineApi {
1641 fn into_rpc_module(self) -> RpcModule<()> {
1642 RpcModule::new(())
1643 }
1644}
1645
1646#[derive(Clone, Debug)]
1651pub struct EngineShutdown {
1652 tx: Arc<Mutex<Option<oneshot::Sender<EngineShutdownRequest>>>>,
1654}
1655
1656impl EngineShutdown {
1657 pub fn new() -> (Self, oneshot::Receiver<EngineShutdownRequest>) {
1659 let (tx, rx) = oneshot::channel();
1660 (Self { tx: Arc::new(Mutex::new(Some(tx))) }, rx)
1661 }
1662
1663 pub fn shutdown(&self) -> Option<oneshot::Receiver<()>> {
1670 let mut guard = self.tx.lock();
1671 let tx = guard.take()?;
1672 let (done_tx, done_rx) = oneshot::channel();
1673 let _ = tx.send(EngineShutdownRequest { done_tx });
1674 Some(done_rx)
1675 }
1676}
1677
1678impl Default for EngineShutdown {
1679 fn default() -> Self {
1680 Self { tx: Arc::new(Mutex::new(None)) }
1681 }
1682}
1683
1684#[derive(Debug)]
1686pub struct EngineShutdownRequest {
1687 pub done_tx: oneshot::Sender<()>,
1689}
1690
1691async fn prewarm_new_block_bals_task<EthApi: GetBlockAccessList, N: NodePrimitives>(
1693 eth_api: EthApi,
1694 mut events: impl Stream<Item = CanonStateNotification<N>> + Unpin,
1695 startup_blocks: usize,
1696) {
1697 let mut startup = match eth_api.recovered_block(BlockNumberOrTag::Latest.into()).await {
1698 Ok(Some(head)) => {
1699 if head.block_access_list_hash().is_some() {
1700 debug!(target: "reth::cli", block_hash = ?head.hash(), "Stopping BAL prewarming: native BALs available");
1701 return;
1702 }
1703
1704 head.number().saturating_sub(startup_blocks as u64)..head.number()
1706 }
1707 Err(err) => {
1708 debug!(target: "reth::cli", %err, "Failed to load head for BAL prewarming");
1709 0..0
1710 }
1711 Ok(None) => 0..0,
1712 }
1713 .map(|number| number + 1);
1714
1715 loop {
1716 tokio::select! {
1717 biased;
1719 Some(event) = events.next() => {
1720 let committed = event.committed();
1721 if let Some(block) = committed.blocks_iter().find(|block| block.block_access_list_hash().is_some()) {
1722 debug!(target: "reth::cli", block_hash = ?block.hash(), "Stopping BAL prewarming: native BALs available");
1723 return;
1724 }
1725
1726 for block in committed.blocks_iter() {
1727 if let Err(err) = eth_api.get_block_access_list(block.hash().into()).await {
1728 debug!(
1729 target: "reth::cli",
1730 %err,
1731 block_hash = ?block.hash(),
1732 "Failed to prewarm BAL for canonical block",
1733 );
1734 }
1735 }
1736 }
1737 Some(number) = async { startup.next() } => {
1738 if let Err(err) = eth_api.get_block_access_list(number.into()).await {
1739 debug!(target: "reth::cli", %err, block_number = number, "Failed to prewarm BAL on startup");
1740 }
1741 }
1742 else => break,
1743 }
1744 }
1745}
1746
1747#[cfg(test)]
1748mod tests {
1749 use super::*;
1750 use alloy_consensus::{Block, BlockBody, Header};
1751 use alloy_primitives::B256;
1752 use reth_evm_ethereum::EthEvmConfig;
1753 use reth_network_api::noop::NoopNetwork;
1754 use reth_primitives_traits::RecoveredBlock;
1755 use reth_provider::{test_utils::MockEthProvider, Chain, ExecutionOutcome};
1756 use reth_rpc::EthApiBuilder;
1757 use reth_rpc_eth_types::EthStateCacheConfig;
1758 use reth_transaction_pool::noop::NoopTransactionPool;
1759
1760 #[tokio::test]
1761 async fn bal_prewarming_startup_range_and_native_stop() {
1762 for (count, native, native_event) in [
1763 (0, false, false),
1764 (2, false, false),
1765 (10, false, false),
1766 (2, true, false),
1767 (10, false, true),
1768 ] {
1769 let provider = MockEthProvider::default();
1770 let blocks: Vec<_> = (0..=3)
1771 .map(|number| {
1772 let block = RecoveredBlock::new_unhashed(
1773 Block {
1774 header: Header {
1775 number,
1776 block_access_list_hash: (native && number == 3)
1777 .then_some(B256::ZERO),
1778 ..Default::default()
1779 },
1780 body: BlockBody::default(),
1781 },
1782 Vec::new(),
1783 );
1784 provider.add_block(block.hash(), block.clone_block());
1785 block
1786 })
1787 .collect();
1788 let api = EthApiBuilder::new(
1789 provider,
1790 NoopTransactionPool::default(),
1791 NoopNetwork::default(),
1792 EthEvmConfig::mainnet(),
1793 )
1794 .eth_state_cache_config(EthStateCacheConfig {
1795 prewarm_bals: Some(count),
1796 ..Default::default()
1797 })
1798 .build();
1799
1800 let chain = Chain::new(
1801 blocks.clone(),
1802 ExecutionOutcome { receipts: vec![vec![]; blocks.len()], ..Default::default() },
1803 Default::default(),
1804 );
1805 cache_new_blocks_task(
1806 api.cache().clone(),
1807 futures::stream::iter([reth_chain_state::CanonStateNotification::Commit {
1808 new: Arc::new(chain),
1809 }]),
1810 )
1811 .await;
1812 let events = futures::stream::iter(native_event.then(|| {
1813 let mut block = blocks.last().unwrap().clone_block();
1814 block.header.number += 1;
1815 block.header.block_access_list_hash = Some(B256::ZERO);
1816 CanonStateNotification::Commit {
1817 new: Arc::new(Chain::<reth_chain_state::EthPrimitives>::new(
1818 [RecoveredBlock::new_unhashed(block, Vec::new())],
1819 Default::default(),
1820 Default::default(),
1821 )),
1822 }
1823 }));
1824 prewarm_new_block_bals_task(api.clone(), events, count).await;
1825 for (number, block) in blocks.iter().enumerate() {
1826 assert_eq!(
1827 api.cache().get_bal(block.hash()).await.unwrap().is_some(),
1828 !native &&
1829 !native_event &&
1830 number != 0 &&
1831 number > 3_usize.saturating_sub(count),
1832 "count={count}, native={native}, native_event={native_event}, block={number}",
1833 );
1834 }
1835 }
1836 }
1837}