1#![doc(
15 html_logo_url = "https://raw.githubusercontent.com/paradigmxyz/reth/main/assets/reth-docs.png",
16 html_favicon_url = "https://avatars0.githubusercontent.com/u/97369466?s=256",
17 issue_tracker_base_url = "https://github.com/paradigmxyz/reth/issues/"
18)]
19#![cfg_attr(not(test), warn(unused_crate_dependencies))]
20#![cfg_attr(docsrs, feature(doc_cfg))]
21
22use crate::{auth::AuthRpcModule, error::WsHttpSamePortError, metrics::RpcRequestMetrics};
23use alloy_network::{Ethereum, IntoWallet};
24use alloy_provider::{fillers::RecommendedFillers, Provider, ProviderBuilder};
25use core::marker::PhantomData;
26use error::{ConflictingModules, RpcError, ServerKind};
27use http::{header::AUTHORIZATION, HeaderMap};
28use jsonrpsee::{
29 core::RegisterMethodError,
30 server::{middleware::rpc::RpcServiceBuilder, AlreadyStoppedError, IdProvider, ServerHandle},
31 Methods, RpcModule,
32};
33use reth_chainspec::{ChainSpecProvider, EthereumHardforks};
34use reth_consensus::FullConsensus;
35use reth_engine_primitives::{ConsensusEngineEvent, ConsensusEngineHandle};
36use reth_evm::ConfigureEvm;
37use reth_network_api::{noop::NoopNetwork, NetworkInfo, Peers};
38use reth_payload_primitives::PayloadTypes;
39use reth_primitives_traits::{NodePrimitives, TxTy};
40use reth_rpc::{
41 AdminApi, DebugApi, EngineEthApi, EthApi, EthApiBuilder, EthBundle, MinerApi, NetApi,
42 OtterscanApi, RPCApi, RethApi, TraceApi, TxPoolApi, Web3Api,
43};
44use reth_rpc_api::servers::*;
45use reth_rpc_engine_api::RethEngineApi;
46use reth_rpc_eth_api::{
47 helpers::{
48 pending_block::PendingEnvBuilder, Call, EthApiSpec, EthTransactions, LoadPendingBlock,
49 TraceExt,
50 },
51 node::RpcNodeCoreAdapter,
52 EthApiServer, EthApiTypes, FullEthApiServer, FullEthApiTypes, RpcBlock, RpcConvert,
53 RpcConverter, RpcHeader, RpcNodeCore, RpcReceipt, RpcTransaction, RpcTxReq,
54};
55use reth_rpc_eth_types::{receipt::EthReceiptConverter, EthConfig, EthSubscriptionIdProvider};
56use reth_rpc_layer::{
57 AuthLayer, Claims, CompressionLayer, DecompressionLayer, JwtAuthValidator, JwtSecret,
58};
59pub use reth_rpc_server_types::RethRpcModule;
60use reth_storage_api::{
61 BlockReader, ChangeSetReader, FullRpcProvider, NodePrimitivesProvider, StateProviderFactory,
62};
63use reth_tasks::{pool::BlockingTaskGuard, Runtime};
64use reth_tokio_util::EventSender;
65use reth_transaction_pool::{noop::NoopTransactionPool, TransactionPool};
66use serde::{Deserialize, Serialize};
67use std::{
68 collections::HashMap,
69 fmt::Debug,
70 net::{Ipv4Addr, SocketAddr, SocketAddrV4},
71 time::{Duration, SystemTime, UNIX_EPOCH},
72};
73use tower_http::cors::CorsLayer;
74
75pub use cors::CorsDomainError;
76
77pub use jsonrpsee::server::ServerBuilder;
79use jsonrpsee::server::ServerConfigBuilder;
80pub use reth_ipc::server::{
81 Builder as IpcServerBuilder, RpcServiceBuilder as IpcRpcServiceBuilder,
82};
83pub use reth_rpc_server_types::{constants, RpcModuleSelection};
84pub use tower::layer::util::{Identity, Stack};
85
86pub mod auth;
88
89pub mod config;
91
92pub mod middleware;
94
95mod cors;
97
98pub mod error;
100
101pub mod eth;
103pub use eth::EthHandlers;
104
105mod metrics;
107use crate::middleware::RethRpcMiddleware;
108pub use metrics::{MeteredBatchRequestsFuture, MeteredRequestFuture, RpcRequestMetricsService};
109use reth_chain_state::{
110 CanonStateSubscriptions, ForkChoiceSubscriptions, PersistedBlockSubscriptions,
111};
112use reth_rpc::eth::sim_bundle::EthSimBundle;
113
114pub mod rate_limiter;
116
117#[derive(Debug, Clone)]
121pub struct RpcModuleBuilder<N, Provider, Pool, Network, EvmConfig, Consensus> {
122 provider: Provider,
124 pool: Pool,
126 network: Network,
128 executor: Option<Runtime>,
130 evm_config: EvmConfig,
132 consensus: Consensus,
134 _primitives: PhantomData<N>,
136}
137
138impl<N, Provider, Pool, Network, EvmConfig, Consensus>
141 RpcModuleBuilder<N, Provider, Pool, Network, EvmConfig, Consensus>
142{
143 pub const fn new(
145 provider: Provider,
146 pool: Pool,
147 network: Network,
148 executor: Runtime,
149 evm_config: EvmConfig,
150 consensus: Consensus,
151 ) -> Self {
152 Self {
153 provider,
154 pool,
155 network,
156 executor: Some(executor),
157 evm_config,
158 consensus,
159 _primitives: PhantomData,
160 }
161 }
162
163 pub fn with_provider<P>(
165 self,
166 provider: P,
167 ) -> RpcModuleBuilder<N, P, Pool, Network, EvmConfig, Consensus> {
168 let Self { pool, network, executor, evm_config, consensus, _primitives, .. } = self;
169 RpcModuleBuilder { provider, network, pool, executor, evm_config, consensus, _primitives }
170 }
171
172 pub fn with_pool<P>(
174 self,
175 pool: P,
176 ) -> RpcModuleBuilder<N, Provider, P, Network, EvmConfig, Consensus> {
177 let Self { provider, network, executor, evm_config, consensus, _primitives, .. } = self;
178 RpcModuleBuilder { provider, network, pool, executor, evm_config, consensus, _primitives }
179 }
180
181 pub fn with_noop_pool(
187 self,
188 ) -> RpcModuleBuilder<N, Provider, NoopTransactionPool, Network, EvmConfig, Consensus> {
189 let Self { provider, executor, network, evm_config, consensus, _primitives, .. } = self;
190 RpcModuleBuilder {
191 provider,
192 executor,
193 network,
194 evm_config,
195 pool: NoopTransactionPool::default(),
196 consensus,
197 _primitives,
198 }
199 }
200
201 pub fn with_network<Net>(
203 self,
204 network: Net,
205 ) -> RpcModuleBuilder<N, Provider, Pool, Net, EvmConfig, Consensus> {
206 let Self { provider, pool, executor, evm_config, consensus, _primitives, .. } = self;
207 RpcModuleBuilder { provider, network, pool, executor, evm_config, consensus, _primitives }
208 }
209
210 pub fn with_noop_network(
216 self,
217 ) -> RpcModuleBuilder<N, Provider, Pool, NoopNetwork, EvmConfig, Consensus> {
218 let Self { provider, pool, executor, evm_config, consensus, _primitives, .. } = self;
219 RpcModuleBuilder {
220 provider,
221 pool,
222 executor,
223 network: NoopNetwork::default(),
224 evm_config,
225 consensus,
226 _primitives,
227 }
228 }
229
230 pub fn with_executor(self, executor: Runtime) -> Self {
232 let Self { pool, network, provider, evm_config, consensus, _primitives, .. } = self;
233 Self {
234 provider,
235 network,
236 pool,
237 executor: Some(executor),
238 evm_config,
239 consensus,
240 _primitives,
241 }
242 }
243
244 pub fn with_evm_config<E>(
246 self,
247 evm_config: E,
248 ) -> RpcModuleBuilder<N, Provider, Pool, Network, E, Consensus> {
249 let Self { provider, pool, executor, network, consensus, _primitives, .. } = self;
250 RpcModuleBuilder { provider, network, pool, executor, evm_config, consensus, _primitives }
251 }
252
253 pub fn with_consensus<C>(
255 self,
256 consensus: C,
257 ) -> RpcModuleBuilder<N, Provider, Pool, Network, EvmConfig, C> {
258 let Self { provider, network, pool, executor, evm_config, _primitives, .. } = self;
259 RpcModuleBuilder { provider, network, pool, executor, evm_config, consensus, _primitives }
260 }
261
262 #[expect(clippy::type_complexity)]
264 pub fn eth_api_builder<ChainSpec>(
265 &self,
266 ) -> EthApiBuilder<
267 RpcNodeCoreAdapter<Provider, Pool, Network, EvmConfig>,
268 RpcConverter<Ethereum, EvmConfig, EthReceiptConverter<ChainSpec>>,
269 >
270 where
271 Provider: Clone,
272 Pool: Clone,
273 Network: Clone,
274 EvmConfig: Clone,
275 RpcNodeCoreAdapter<Provider, Pool, Network, EvmConfig>:
276 RpcNodeCore<Provider: ChainSpecProvider<ChainSpec = ChainSpec>, Evm = EvmConfig>,
277 {
278 EthApiBuilder::new(
279 self.provider.clone(),
280 self.pool.clone(),
281 self.network.clone(),
282 self.evm_config.clone(),
283 )
284 }
285
286 #[expect(clippy::type_complexity)]
292 pub fn bootstrap_eth_api<ChainSpec>(
293 &self,
294 ) -> EthApi<
295 RpcNodeCoreAdapter<Provider, Pool, Network, EvmConfig>,
296 RpcConverter<Ethereum, EvmConfig, EthReceiptConverter<ChainSpec>>,
297 >
298 where
299 Provider: Clone,
300 Pool: Clone,
301 Network: Clone,
302 EvmConfig: ConfigureEvm + Clone,
303 RpcNodeCoreAdapter<Provider, Pool, Network, EvmConfig>:
304 RpcNodeCore<Provider: ChainSpecProvider<ChainSpec = ChainSpec>, Evm = EvmConfig>,
305 RpcConverter<Ethereum, EvmConfig, EthReceiptConverter<ChainSpec>>: RpcConvert,
306 (): PendingEnvBuilder<EvmConfig>,
307 {
308 self.eth_api_builder().build()
309 }
310}
311
312impl<N, Provider, Pool, Network, EvmConfig, Consensus>
313 RpcModuleBuilder<N, Provider, Pool, Network, EvmConfig, Consensus>
314where
315 N: NodePrimitives,
316 Provider: FullRpcProvider<Block = N::Block, Receipt = N::Receipt, Header = N::BlockHeader>
317 + CanonStateSubscriptions<Primitives = N>
318 + ForkChoiceSubscriptions<Header = N::BlockHeader>
319 + PersistedBlockSubscriptions
320 + ChangeSetReader
321 + NodePrimitivesProvider<Primitives = N>,
322 Pool: TransactionPool + Clone + 'static,
323 Network: NetworkInfo + Peers + Clone + 'static,
324 EvmConfig: ConfigureEvm<Primitives = N> + 'static,
325 Consensus: FullConsensus<N> + Clone + 'static,
326{
327 pub fn build_with_auth_server<EthApi, Payload>(
334 self,
335 module_config: TransportRpcModuleConfig,
336 engine: impl IntoEngineApiRpcModule,
337 eth: EthApi,
338 engine_events: EventSender<ConsensusEngineEvent<N>>,
339 beacon_engine_handle: ConsensusEngineHandle<Payload>,
340 ) -> (
341 TransportRpcModules,
342 AuthRpcModule,
343 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>,
344 )
345 where
346 EthApi: FullEthApiServer<Provider = Provider, Pool = Pool>,
347 Payload: PayloadTypes,
348 {
349 let config = module_config.config.clone().unwrap_or_default();
350
351 let mut registry = self.into_registry(config, eth, engine_events);
352 let modules = registry.create_transport_rpc_modules(module_config);
353 let auth_module = registry.create_auth_module(engine, beacon_engine_handle);
354
355 (modules, auth_module, registry)
356 }
357
358 pub fn into_registry<EthApi>(
363 self,
364 config: RpcModuleConfig,
365 eth: EthApi,
366 engine_events: EventSender<ConsensusEngineEvent<N>>,
367 ) -> RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
368 where
369 EthApi: FullEthApiServer<Provider = Provider, Pool = Pool>,
370 {
371 let Self { provider, pool, network, executor, consensus, evm_config, .. } = self;
372 let executor =
373 executor.expect("RpcModuleBuilder requires a Runtime to be set via `with_executor`");
374 RpcRegistryInner::new(
375 provider,
376 pool,
377 network,
378 executor,
379 consensus,
380 config,
381 evm_config,
382 eth,
383 engine_events,
384 )
385 }
386
387 pub fn build<EthApi>(
390 self,
391 module_config: TransportRpcModuleConfig,
392 eth: EthApi,
393 engine_events: EventSender<ConsensusEngineEvent<N>>,
394 ) -> TransportRpcModules<()>
395 where
396 EthApi: FullEthApiServer<Provider = Provider, Pool = Pool>,
397 {
398 if module_config.is_empty() {
399 TransportRpcModules::default()
400 } else {
401 let config = module_config.config.clone().unwrap_or_default();
402 let mut registry = self.into_registry(config, eth, engine_events);
403 registry.create_transport_rpc_modules(module_config)
404 }
405 }
406}
407
408impl<N: NodePrimitives> Default for RpcModuleBuilder<N, (), (), (), (), ()> {
409 fn default() -> Self {
410 Self {
411 provider: (),
412 pool: (),
413 network: (),
414 executor: None,
415 evm_config: (),
416 consensus: (),
417 _primitives: PhantomData,
418 }
419 }
420}
421
422#[derive(Debug, Default, Clone, Eq, PartialEq, Serialize, Deserialize)]
424pub struct RpcModuleConfig {
425 eth: EthConfig,
427}
428
429impl RpcModuleConfig {
432 pub fn builder() -> RpcModuleConfigBuilder {
434 RpcModuleConfigBuilder::default()
435 }
436
437 pub const fn new(eth: EthConfig) -> Self {
439 Self { eth }
440 }
441
442 pub const fn eth(&self) -> &EthConfig {
444 &self.eth
445 }
446
447 pub const fn eth_mut(&mut self) -> &mut EthConfig {
449 &mut self.eth
450 }
451}
452
453#[derive(Clone, Debug, Default)]
455pub struct RpcModuleConfigBuilder {
456 eth: Option<EthConfig>,
457}
458
459impl RpcModuleConfigBuilder {
462 pub fn eth(mut self, eth: EthConfig) -> Self {
464 self.eth = Some(eth);
465 self
466 }
467
468 pub fn build(self) -> RpcModuleConfig {
470 let Self { eth } = self;
471 RpcModuleConfig { eth: eth.unwrap_or_default() }
472 }
473
474 pub const fn get_eth(&self) -> Option<&EthConfig> {
476 self.eth.as_ref()
477 }
478
479 pub const fn eth_mut(&mut self) -> &mut Option<EthConfig> {
481 &mut self.eth
482 }
483
484 pub fn eth_mut_or_default(&mut self) -> &mut EthConfig {
486 self.eth.get_or_insert_with(EthConfig::default)
487 }
488}
489
490#[derive(Debug)]
492pub struct RpcRegistryInner<Provider, Pool, Network, EthApi: EthApiTypes, EvmConfig, Consensus> {
493 provider: Provider,
494 pool: Pool,
495 network: Network,
496 executor: Runtime,
497 evm_config: EvmConfig,
498 consensus: Consensus,
499 eth: EthHandlers<EthApi>,
501 blocking_pool_guard: BlockingTaskGuard,
503 modules: HashMap<RethRpcModule, Methods>,
505 eth_config: EthConfig,
507 engine_events:
509 EventSender<ConsensusEngineEvent<<EthApi::RpcConvert as RpcConvert>::Primitives>>,
510}
511
512impl<N, Provider, Pool, Network, EthApi, EvmConfig, Consensus>
515 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
516where
517 N: NodePrimitives,
518 Provider: StateProviderFactory
519 + CanonStateSubscriptions<Primitives = N>
520 + BlockReader<Block = N::Block, Receipt = N::Receipt>
521 + Clone
522 + Unpin
523 + 'static,
524 Pool: Send + Sync + Clone + 'static,
525 Network: Clone + 'static,
526 EthApi: FullEthApiTypes + 'static,
527 EvmConfig: ConfigureEvm<Primitives = N>,
528{
529 #[expect(clippy::too_many_arguments)]
531 pub fn new(
532 provider: Provider,
533 pool: Pool,
534 network: Network,
535 executor: Runtime,
536 consensus: Consensus,
537 config: RpcModuleConfig,
538 evm_config: EvmConfig,
539 eth_api: EthApi,
540 engine_events: EventSender<
541 ConsensusEngineEvent<<EthApi::Provider as NodePrimitivesProvider>::Primitives>,
542 >,
543 ) -> Self
544 where
545 EvmConfig: ConfigureEvm<Primitives = N>,
546 {
547 let blocking_pool_guard = BlockingTaskGuard::new(config.eth.max_tracing_requests);
548
549 let eth = EthHandlers::bootstrap(config.eth.clone(), executor.clone(), eth_api);
550
551 Self {
552 provider,
553 pool,
554 network,
555 eth,
556 executor,
557 consensus,
558 modules: Default::default(),
559 blocking_pool_guard,
560 eth_config: config.eth,
561 evm_config,
562 engine_events,
563 }
564 }
565}
566
567impl<Provider, Pool, Network, EthApi, Evm, Consensus>
568 RpcRegistryInner<Provider, Pool, Network, EthApi, Evm, Consensus>
569where
570 EthApi: EthApiTypes,
571{
572 pub const fn eth_api(&self) -> &EthApi {
574 &self.eth.api
575 }
576
577 pub const fn eth_handlers(&self) -> &EthHandlers<EthApi> {
579 &self.eth
580 }
581
582 pub const fn pool(&self) -> &Pool {
584 &self.pool
585 }
586
587 pub const fn tasks(&self) -> &Runtime {
589 &self.executor
590 }
591
592 pub const fn provider(&self) -> &Provider {
594 &self.provider
595 }
596
597 pub const fn evm_config(&self) -> &Evm {
599 &self.evm_config
600 }
601
602 pub fn methods(&self) -> Vec<Methods> {
604 self.modules.values().cloned().collect()
605 }
606
607 pub fn module(&self) -> RpcModule<()> {
609 let mut module = RpcModule::new(());
610 for methods in self.modules.values().cloned() {
611 module.merge(methods).expect("No conflicts");
612 }
613 module
614 }
615}
616
617impl<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
618 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
619where
620 Network: NetworkInfo + Clone + 'static,
621 EthApi: EthApiTypes,
622 Provider: BlockReader + ChainSpecProvider<ChainSpec: EthereumHardforks>,
623 EvmConfig: ConfigureEvm,
624{
625 pub fn admin_api(&self) -> AdminApi<Network, Provider::ChainSpec, Pool>
627 where
628 Network: Peers,
629 Pool: TransactionPool + Clone + 'static,
630 {
631 AdminApi::new(self.network.clone(), self.provider.chain_spec(), self.pool.clone())
632 }
633
634 pub fn web3_api(&self) -> Web3Api<Network> {
636 Web3Api::new(self.network.clone())
637 }
638
639 pub fn register_admin(&mut self) -> &mut Self
641 where
642 Network: Peers,
643 Pool: TransactionPool + Clone + 'static,
644 {
645 let adminapi = self.admin_api();
646 self.modules.insert(RethRpcModule::Admin, adminapi.into_rpc().into());
647 self
648 }
649
650 pub fn register_web3(&mut self) -> &mut Self {
652 let web3api = self.web3_api();
653 self.modules.insert(RethRpcModule::Web3, web3api.into_rpc().into());
654 self
655 }
656}
657
658impl<N, Provider, Pool, Network, EthApi, EvmConfig, Consensus>
659 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
660where
661 N: NodePrimitives,
662 Provider: FullRpcProvider<
663 Header = N::BlockHeader,
664 Block = N::Block,
665 Receipt = N::Receipt,
666 Transaction = N::SignedTx,
667 > + ChangeSetReader
668 + CanonStateSubscriptions<Primitives = N>
669 + ForkChoiceSubscriptions<Header = N::BlockHeader>
670 + PersistedBlockSubscriptions,
671 Network: NetworkInfo + Peers + Clone + 'static,
672 EthApi: EthApiServer<
673 RpcTxReq<EthApi::NetworkTypes>,
674 RpcTransaction<EthApi::NetworkTypes>,
675 RpcBlock<EthApi::NetworkTypes>,
676 RpcReceipt<EthApi::NetworkTypes>,
677 RpcHeader<EthApi::NetworkTypes>,
678 TxTy<N>,
679 > + EthApiTypes,
680 EvmConfig: ConfigureEvm<Primitives = N> + 'static,
681{
682 pub fn register_eth(&mut self) -> &mut Self {
688 let eth_api = self.eth_api().clone();
689 self.modules.insert(RethRpcModule::Eth, eth_api.into_rpc().into());
690 self
691 }
692
693 pub fn register_ots(&mut self) -> &mut Self
699 where
700 EthApi: TraceExt + EthTransactions<Primitives = N>,
701 {
702 let otterscan_api = self.otterscan_api();
703 self.modules.insert(RethRpcModule::Ots, otterscan_api.into_rpc().into());
704 self
705 }
706
707 pub fn register_debug(&mut self) -> &mut Self
713 where
714 EthApi: EthTransactions + TraceExt,
715 {
716 let debug_api = self.debug_api();
717 self.modules.insert(RethRpcModule::Debug, debug_api.into_rpc().into());
718 self
719 }
720
721 pub fn register_trace(&mut self) -> &mut Self
727 where
728 EthApi: TraceExt,
729 {
730 let trace_api = self.trace_api();
731 self.modules.insert(RethRpcModule::Trace, trace_api.into_rpc().into());
732 self
733 }
734
735 pub fn register_net(&mut self) -> &mut Self
743 where
744 EthApi: EthApiSpec + 'static,
745 {
746 let netapi = self.net_api();
747 self.modules.insert(RethRpcModule::Net, netapi.into_rpc().into());
748 self
749 }
750
751 pub fn register_reth(&mut self) -> &mut Self {
759 let rethapi = self.reth_api();
760 self.modules.insert(RethRpcModule::Reth, rethapi.into_rpc().into());
761 self
762 }
763
764 pub fn otterscan_api(&self) -> OtterscanApi<EthApi> {
770 let eth_api = self.eth_api().clone();
771 OtterscanApi::new(eth_api)
772 }
773}
774
775impl<N, Provider, Pool, Network, EthApi, EvmConfig, Consensus>
776 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
777where
778 N: NodePrimitives,
779 Provider: FullRpcProvider<
780 Block = N::Block,
781 Header = N::BlockHeader,
782 Transaction = N::SignedTx,
783 Receipt = N::Receipt,
784 > + ChangeSetReader,
785 Network: NetworkInfo + Peers + Clone + 'static,
786 EthApi: EthApiTypes,
787 EvmConfig: ConfigureEvm<Primitives = N>,
788{
789 pub fn trace_api(&self) -> TraceApi<EthApi> {
795 TraceApi::new(
796 self.eth_api().clone(),
797 self.blocking_pool_guard.clone(),
798 self.eth_config.clone(),
799 )
800 }
801
802 pub fn bundle_api(&self) -> EthBundle<EthApi>
808 where
809 EthApi: EthTransactions + LoadPendingBlock + Call,
810 {
811 let eth_api = self.eth_api().clone();
812 EthBundle::new(eth_api, self.blocking_pool_guard.clone())
813 }
814
815 pub fn debug_api(&self) -> DebugApi<EthApi>
821 where
822 EthApi: FullEthApiTypes,
823 {
824 DebugApi::new(
825 self.eth_api().clone(),
826 self.blocking_pool_guard.clone(),
827 self.tasks(),
828 self.engine_events.new_listener(),
829 )
830 }
831
832 pub fn net_api(&self) -> NetApi<Network, EthApi>
838 where
839 EthApi: EthApiSpec + 'static,
840 {
841 let eth_api = self.eth_api().clone();
842 NetApi::new(self.network.clone(), eth_api)
843 }
844
845 pub fn reth_api(&self) -> RethApi<Provider, EvmConfig> {
847 RethApi::new(
848 self.provider.clone(),
849 self.evm_config.clone(),
850 self.blocking_pool_guard.clone(),
851 self.executor.clone(),
852 )
853 }
854}
855
856impl<N, Provider, Pool, Network, EthApi, EvmConfig, Consensus>
857 RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
858where
859 N: NodePrimitives,
860 Provider: FullRpcProvider<Block = N::Block>
861 + CanonStateSubscriptions<Primitives = N>
862 + ForkChoiceSubscriptions<Header = N::BlockHeader>
863 + PersistedBlockSubscriptions
864 + ChangeSetReader,
865 Pool: TransactionPool + Clone + 'static,
866 Network: NetworkInfo + Peers + Clone + 'static,
867 EthApi: FullEthApiServer,
868 EvmConfig: ConfigureEvm<Primitives = N> + 'static,
869 Consensus: FullConsensus<N> + Clone + 'static,
870{
871 pub fn create_auth_module<Payload>(
878 &self,
879 engine_api: impl IntoEngineApiRpcModule,
880 beacon_engine_handle: ConsensusEngineHandle<Payload>,
881 ) -> AuthRpcModule
882 where
883 Payload: PayloadTypes,
884 {
885 let mut module = engine_api.into_rpc_module();
886
887 let reth_engine_api = RethEngineApi::new(beacon_engine_handle);
889 module
890 .merge(RethEngineApiServer::into_rpc(reth_engine_api).remove_context())
891 .expect("No conflicting methods");
892
893 let eth_handlers = self.eth_handlers();
895 let engine_eth = EngineEthApi::new(eth_handlers.api.clone(), eth_handlers.filter.clone());
896
897 module.merge(engine_eth.into_rpc()).expect("No conflicting methods");
898
899 AuthRpcModule { inner: module }
900 }
901
902 fn maybe_module(&mut self, config: Option<&RpcModuleSelection>) -> Option<RpcModule<()>> {
904 config.map(|config| self.module_for(config))
905 }
906
907 pub fn create_transport_rpc_modules(
911 &mut self,
912 config: TransportRpcModuleConfig,
913 ) -> TransportRpcModules<()> {
914 let mut modules = TransportRpcModules::default();
915 let http = self.maybe_module(config.http.as_ref());
916 let ws = self.maybe_module(config.ws.as_ref());
917 let ipc = self.maybe_module(config.ipc.as_ref());
918
919 modules.config = config;
920 modules.http = http;
921 modules.ws = ws;
922 modules.ipc = ipc;
923 modules
924 }
925
926 pub fn module_for(&mut self, config: &RpcModuleSelection) -> RpcModule<()> {
929 let mut module = RpcModule::new(());
930 let all_methods = self.reth_methods(config.iter_selection());
931 for methods in all_methods {
932 module.merge(methods).expect("No conflicts");
933 }
934 module
935 }
936
937 pub fn reth_methods(
946 &mut self,
947 namespaces: impl Iterator<Item = RethRpcModule>,
948 ) -> Vec<Methods> {
949 let EthHandlers { api: eth_api, filter: eth_filter, pubsub: eth_pubsub, .. } =
950 self.eth_handlers().clone();
951
952 let namespaces: Vec<_> = namespaces.collect();
954 namespaces
955 .iter()
956 .map(|namespace| {
957 self.modules
958 .entry(namespace.clone())
959 .or_insert_with(|| match namespace.clone() {
960 RethRpcModule::Admin => AdminApi::new(
961 self.network.clone(),
962 self.provider.chain_spec(),
963 self.pool.clone(),
964 )
965 .into_rpc()
966 .into(),
967 RethRpcModule::Debug => DebugApi::new(
968 eth_api.clone(),
969 self.blocking_pool_guard.clone(),
970 &self.executor,
971 self.engine_events.new_listener(),
972 )
973 .into_rpc()
974 .into(),
975 RethRpcModule::Eth => {
976 let mut module = eth_api.clone().into_rpc();
978 module.merge(eth_filter.clone().into_rpc()).expect("No conflicts");
979 module.merge(eth_pubsub.clone().into_rpc()).expect("No conflicts");
980 module
981 .merge(
982 EthBundle::new(
983 eth_api.clone(),
984 self.blocking_pool_guard.clone(),
985 )
986 .into_rpc(),
987 )
988 .expect("No conflicts");
989
990 module.into()
991 }
992 RethRpcModule::Net => {
993 NetApi::new(self.network.clone(), eth_api.clone()).into_rpc().into()
994 }
995 RethRpcModule::Trace => TraceApi::new(
996 eth_api.clone(),
997 self.blocking_pool_guard.clone(),
998 self.eth_config.clone(),
999 )
1000 .into_rpc()
1001 .into(),
1002 RethRpcModule::Web3 => Web3Api::new(self.network.clone()).into_rpc().into(),
1003 RethRpcModule::Txpool => TxPoolApi::new(
1004 self.eth.api.pool().clone(),
1005 dyn_clone::clone(self.eth.api.converter()),
1006 )
1007 .into_rpc()
1008 .into(),
1009 RethRpcModule::Rpc => RPCApi::new(
1010 namespaces
1011 .iter()
1012 .map(|module| (module.to_string(), "1.0".to_string()))
1013 .collect(),
1014 )
1015 .into_rpc()
1016 .into(),
1017 RethRpcModule::Ots => OtterscanApi::new(eth_api.clone()).into_rpc().into(),
1018 RethRpcModule::Reth => RethApi::new(
1019 self.provider.clone(),
1020 self.evm_config.clone(),
1021 self.blocking_pool_guard.clone(),
1022 self.executor.clone(),
1023 )
1024 .into_rpc()
1025 .into(),
1026 RethRpcModule::Miner => MinerApi::default().into_rpc().into(),
1027 RethRpcModule::Mev => {
1028 EthSimBundle::new(eth_api.clone(), self.blocking_pool_guard.clone())
1029 .into_rpc()
1030 .into()
1031 }
1032 RethRpcModule::Flashbots |
1036 RethRpcModule::Testing |
1037 RethRpcModule::Other(_) => Default::default(),
1038 })
1039 .clone()
1040 })
1041 .collect::<Vec<_>>()
1042 }
1043}
1044
1045impl<Provider, Pool, Network, EthApi, EvmConfig, Consensus> Clone
1046 for RpcRegistryInner<Provider, Pool, Network, EthApi, EvmConfig, Consensus>
1047where
1048 EthApi: EthApiTypes,
1049 Provider: Clone,
1050 Pool: Clone,
1051 Network: Clone,
1052 EvmConfig: Clone,
1053 Consensus: Clone,
1054{
1055 fn clone(&self) -> Self {
1056 Self {
1057 provider: self.provider.clone(),
1058 pool: self.pool.clone(),
1059 network: self.network.clone(),
1060 executor: self.executor.clone(),
1061 evm_config: self.evm_config.clone(),
1062 consensus: self.consensus.clone(),
1063 eth: self.eth.clone(),
1064 blocking_pool_guard: self.blocking_pool_guard.clone(),
1065 modules: self.modules.clone(),
1066 eth_config: self.eth_config.clone(),
1067 engine_events: self.engine_events.clone(),
1068 }
1069 }
1070}
1071
1072#[derive(Debug)]
1084pub struct RpcServerConfig<RpcMiddleware = Identity> {
1085 http_server_config: Option<ServerConfigBuilder>,
1087 http_cors_domains: Option<String>,
1089 http_addr: Option<SocketAddr>,
1091 http_disable_compression: bool,
1093 http_compression_algorithms: Option<Vec<String>>,
1097 http_decompression_algorithms: Option<Vec<String>>,
1101 http_max_request_body_size: Option<u32>,
1103 ws_server_config: Option<ServerConfigBuilder>,
1105 ws_cors_domains: Option<String>,
1107 ws_addr: Option<SocketAddr>,
1109 ipc_server_config: Option<IpcServerBuilder<Identity, Identity>>,
1111 ipc_endpoint: Option<String>,
1113 jwt_secret: Option<JwtSecret>,
1115 rpc_metrics_enabled: bool,
1117 rpc_middleware: RpcMiddleware,
1119}
1120
1121impl Default for RpcServerConfig<Identity> {
1124 fn default() -> Self {
1126 Self {
1127 http_server_config: None,
1128 http_cors_domains: None,
1129 http_addr: None,
1130 http_disable_compression: false,
1131 http_compression_algorithms: None,
1132 http_decompression_algorithms: None,
1133 http_max_request_body_size: None,
1134 ws_server_config: None,
1135 ws_cors_domains: None,
1136 ws_addr: None,
1137 ipc_server_config: None,
1138 ipc_endpoint: None,
1139 jwt_secret: None,
1140 rpc_metrics_enabled: true,
1141 rpc_middleware: Default::default(),
1142 }
1143 }
1144}
1145
1146impl RpcServerConfig {
1147 pub fn http(config: ServerConfigBuilder) -> Self {
1149 Self::default().with_http(config)
1150 }
1151
1152 pub fn ws(config: ServerConfigBuilder) -> Self {
1154 Self::default().with_ws(config)
1155 }
1156
1157 pub fn ipc(config: IpcServerBuilder<Identity, Identity>) -> Self {
1159 Self::default().with_ipc(config)
1160 }
1161
1162 pub fn with_http(mut self, config: ServerConfigBuilder) -> Self {
1167 self.http_server_config =
1168 Some(config.set_id_provider(EthSubscriptionIdProvider::default()));
1169 self
1170 }
1171
1172 pub fn with_ws(mut self, config: ServerConfigBuilder) -> Self {
1177 self.ws_server_config = Some(config.set_id_provider(EthSubscriptionIdProvider::default()));
1178 self
1179 }
1180
1181 pub fn with_ipc(mut self, config: IpcServerBuilder<Identity, Identity>) -> Self {
1186 self.ipc_server_config = Some(config.set_id_provider(EthSubscriptionIdProvider::default()));
1187 self
1188 }
1189}
1190
1191impl<RpcMiddleware> RpcServerConfig<RpcMiddleware> {
1192 pub fn set_rpc_middleware<T>(self, rpc_middleware: T) -> RpcServerConfig<T> {
1194 RpcServerConfig {
1195 http_server_config: self.http_server_config,
1196 http_cors_domains: self.http_cors_domains,
1197 http_addr: self.http_addr,
1198 http_disable_compression: self.http_disable_compression,
1199 http_compression_algorithms: self.http_compression_algorithms,
1200 http_decompression_algorithms: self.http_decompression_algorithms,
1201 http_max_request_body_size: self.http_max_request_body_size,
1202 ws_server_config: self.ws_server_config,
1203 ws_cors_domains: self.ws_cors_domains,
1204 ws_addr: self.ws_addr,
1205 ipc_server_config: self.ipc_server_config,
1206 ipc_endpoint: self.ipc_endpoint,
1207 jwt_secret: self.jwt_secret,
1208 rpc_metrics_enabled: self.rpc_metrics_enabled,
1209 rpc_middleware,
1210 }
1211 }
1212
1213 pub const fn with_rpc_metrics_enabled(mut self, enabled: bool) -> Self {
1215 self.rpc_metrics_enabled = enabled;
1216 self
1217 }
1218
1219 pub fn with_cors(self, cors_domain: Option<String>) -> Self {
1221 self.with_http_cors(cors_domain.clone()).with_ws_cors(cors_domain)
1222 }
1223
1224 pub fn with_ws_cors(mut self, cors_domain: Option<String>) -> Self {
1226 self.ws_cors_domains = cors_domain;
1227 self
1228 }
1229
1230 pub fn with_http_cors(mut self, cors_domain: Option<String>) -> Self {
1232 self.http_cors_domains = cors_domain;
1233 self
1234 }
1235
1236 pub const fn with_http_disable_compression(mut self, http_disable_compression: bool) -> Self {
1238 self.http_disable_compression = http_disable_compression;
1239 self
1240 }
1241
1242 pub fn with_http_compression_algorithms(mut self, algos: Option<Vec<String>>) -> Self {
1247 self.http_compression_algorithms = algos;
1248 self
1249 }
1250
1251 pub fn with_http_decompression(mut self, algos: Option<Vec<String>>, max_size: u32) -> Self {
1255 self.http_decompression_algorithms = algos;
1256 self.http_max_request_body_size = Some(max_size);
1257 self
1258 }
1259
1260 pub const fn with_http_address(mut self, addr: SocketAddr) -> Self {
1265 self.http_addr = Some(addr);
1266 self
1267 }
1268
1269 pub const fn with_ws_address(mut self, addr: SocketAddr) -> Self {
1274 self.ws_addr = Some(addr);
1275 self
1276 }
1277
1278 pub fn with_id_provider<I>(mut self, id_provider: I) -> Self
1282 where
1283 I: IdProvider + Clone + 'static,
1284 {
1285 if let Some(config) = self.http_server_config {
1286 self.http_server_config = Some(config.set_id_provider(id_provider.clone()));
1287 }
1288 if let Some(config) = self.ws_server_config {
1289 self.ws_server_config = Some(config.set_id_provider(id_provider.clone()));
1290 }
1291 if let Some(ipc) = self.ipc_server_config {
1292 self.ipc_server_config = Some(ipc.set_id_provider(id_provider));
1293 }
1294
1295 self
1296 }
1297
1298 pub fn with_ipc_endpoint(mut self, path: impl Into<String>) -> Self {
1302 self.ipc_endpoint = Some(path.into());
1303 self
1304 }
1305
1306 pub const fn with_jwt_secret(mut self, secret: Option<JwtSecret>) -> Self {
1308 self.jwt_secret = secret;
1309 self
1310 }
1311
1312 pub fn with_tokio_runtime(mut self, tokio_runtime: Option<tokio::runtime::Handle>) -> Self {
1314 let Some(tokio_runtime) = tokio_runtime else { return self };
1315 if let Some(http_server_config) = self.http_server_config {
1316 self.http_server_config =
1317 Some(http_server_config.custom_tokio_runtime(tokio_runtime.clone()));
1318 }
1319 if let Some(ws_server_config) = self.ws_server_config {
1320 self.ws_server_config =
1321 Some(ws_server_config.custom_tokio_runtime(tokio_runtime.clone()));
1322 }
1323 if let Some(ipc_server_config) = self.ipc_server_config {
1324 self.ipc_server_config = Some(ipc_server_config.custom_tokio_runtime(tokio_runtime));
1325 }
1326 self
1327 }
1328
1329 pub const fn has_server(&self) -> bool {
1333 self.http_server_config.is_some() ||
1334 self.ws_server_config.is_some() ||
1335 self.ipc_server_config.is_some()
1336 }
1337
1338 pub const fn http_address(&self) -> Option<SocketAddr> {
1340 self.http_addr
1341 }
1342
1343 pub const fn ws_address(&self) -> Option<SocketAddr> {
1345 self.ws_addr
1346 }
1347
1348 pub fn ipc_endpoint(&self) -> Option<String> {
1350 self.ipc_endpoint.clone()
1351 }
1352
1353 pub const fn rpc_metrics_enabled(&self) -> bool {
1355 self.rpc_metrics_enabled
1356 }
1357
1358 fn maybe_cors_layer(cors: Option<String>) -> Result<Option<CorsLayer>, CorsDomainError> {
1360 cors.as_deref().map(cors::create_cors_layer).transpose()
1361 }
1362
1363 fn maybe_jwt_layer(jwt_secret: Option<JwtSecret>) -> Option<AuthLayer<JwtAuthValidator>> {
1365 jwt_secret.map(|secret| AuthLayer::new(JwtAuthValidator::new(secret)))
1366 }
1367
1368 fn maybe_compression_layer(
1371 disable_compression: bool,
1372 algos: Option<&[String]>,
1373 ) -> Option<CompressionLayer> {
1374 if disable_compression {
1375 None
1376 } else {
1377 match algos {
1378 None => Some(CompressionLayer::new()),
1380 Some(algos) => Some(CompressionLayer::with_algorithms(algos)),
1381 }
1382 }
1383 }
1384
1385 fn maybe_decompression_layer(
1391 max_request_body_size: Option<u32>,
1392 algos: Option<&[String]>,
1393 ) -> Option<DecompressionLayer> {
1394 let algos = algos?;
1395 if algos.is_empty() {
1396 return None;
1397 }
1398
1399 let max = max_request_body_size?;
1400 Some(DecompressionLayer::new(algos, max as usize))
1401 }
1402
1403 pub async fn start(self, modules: &TransportRpcModules) -> Result<RpcServerHandle, RpcError>
1409 where
1410 RpcMiddleware: RethRpcMiddleware,
1411 {
1412 let mut http_handle = None;
1413 let mut ws_handle = None;
1414 let mut ipc_handle = None;
1415
1416 let http_socket_addr = self.http_addr.unwrap_or(SocketAddr::V4(SocketAddrV4::new(
1417 Ipv4Addr::LOCALHOST,
1418 constants::DEFAULT_HTTP_RPC_PORT,
1419 )));
1420
1421 let ws_socket_addr = self.ws_addr.unwrap_or(SocketAddr::V4(SocketAddrV4::new(
1422 Ipv4Addr::LOCALHOST,
1423 constants::DEFAULT_WS_RPC_PORT,
1424 )));
1425
1426 let rpc_metrics_enabled = self.rpc_metrics_enabled;
1427 let ipc_path =
1428 self.ipc_endpoint.clone().unwrap_or_else(|| constants::DEFAULT_IPC_ENDPOINT.into());
1429
1430 if let Some(builder) = self.ipc_server_config {
1431 let ipc = builder
1432 .set_rpc_middleware(
1433 IpcRpcServiceBuilder::new().option_layer(
1434 rpc_metrics_enabled
1435 .then(|| modules.ipc.as_ref().map(RpcRequestMetrics::ipc))
1436 .flatten(),
1437 ),
1438 )
1439 .build(ipc_path);
1440 ipc_handle = Some(ipc.start(modules.ipc.clone().expect("ipc server error")).await?);
1441 }
1442
1443 if self.http_addr == self.ws_addr &&
1445 self.http_server_config.is_some() &&
1446 self.ws_server_config.is_some()
1447 {
1448 let cors = match (self.ws_cors_domains.as_ref(), self.http_cors_domains.as_ref()) {
1449 (Some(ws_cors), Some(http_cors)) => {
1450 if ws_cors.trim() != http_cors.trim() {
1451 return Err(WsHttpSamePortError::ConflictingCorsDomains {
1452 http_cors_domains: Some(http_cors.clone()),
1453 ws_cors_domains: Some(ws_cors.clone()),
1454 }
1455 .into());
1456 }
1457 Some(ws_cors)
1458 }
1459 (a, b) => a.or(b),
1460 }
1461 .cloned();
1462
1463 modules.config.ensure_ws_http_identical()?;
1465
1466 if let Some(config) = self.http_server_config {
1467 let server = ServerBuilder::new()
1468 .set_http_middleware(
1469 tower::ServiceBuilder::new()
1470 .option_layer(Self::maybe_cors_layer(cors)?)
1471 .option_layer(Self::maybe_jwt_layer(self.jwt_secret))
1472 .option_layer(Self::maybe_decompression_layer(
1473 self.http_max_request_body_size,
1474 self.http_decompression_algorithms.as_deref(),
1475 ))
1476 .option_layer(Self::maybe_compression_layer(
1477 self.http_disable_compression,
1478 self.http_compression_algorithms.as_deref(),
1479 )),
1480 )
1481 .set_rpc_middleware(
1482 RpcServiceBuilder::default()
1483 .option_layer(
1484 rpc_metrics_enabled
1485 .then(|| {
1486 modules
1487 .http
1488 .as_ref()
1489 .or(modules.ws.as_ref())
1490 .map(RpcRequestMetrics::same_port)
1491 })
1492 .flatten(),
1493 )
1494 .layer(self.rpc_middleware.clone()),
1495 )
1496 .set_config(config.build())
1497 .build(http_socket_addr)
1498 .await
1499 .map_err(|err| {
1500 RpcError::server_error(err, ServerKind::WsHttp(http_socket_addr))
1501 })?;
1502 let addr = server.local_addr().map_err(|err| {
1503 RpcError::server_error(err, ServerKind::WsHttp(http_socket_addr))
1504 })?;
1505 if let Some(module) = modules.http.as_ref().or(modules.ws.as_ref()) {
1506 let handle = server.start(module.clone());
1507 http_handle = Some(handle.clone());
1508 ws_handle = Some(handle);
1509 }
1510 return Ok(RpcServerHandle {
1511 http_local_addr: Some(addr),
1512 ws_local_addr: Some(addr),
1513 http: http_handle,
1514 ws: ws_handle,
1515 ipc_endpoint: self.ipc_endpoint.clone(),
1516 ipc: ipc_handle,
1517 jwt_secret: self.jwt_secret,
1518 });
1519 }
1520 }
1521
1522 let mut ws_local_addr = None;
1523 let mut ws_server = None;
1524 let mut http_local_addr = None;
1525 let mut http_server = None;
1526
1527 if let Some(config) = self.ws_server_config {
1528 let server = ServerBuilder::new()
1529 .set_config(config.ws_only().build())
1530 .set_http_middleware(
1531 tower::ServiceBuilder::new()
1532 .option_layer(Self::maybe_cors_layer(self.ws_cors_domains.clone())?)
1533 .option_layer(Self::maybe_jwt_layer(self.jwt_secret)),
1534 )
1535 .set_rpc_middleware(
1536 RpcServiceBuilder::default()
1537 .option_layer(
1538 rpc_metrics_enabled
1539 .then(|| modules.ws.as_ref().map(RpcRequestMetrics::ws))
1540 .flatten(),
1541 )
1542 .layer(self.rpc_middleware.clone()),
1543 )
1544 .build(ws_socket_addr)
1545 .await
1546 .map_err(|err| RpcError::server_error(err, ServerKind::WS(ws_socket_addr)))?;
1547
1548 let addr = server
1549 .local_addr()
1550 .map_err(|err| RpcError::server_error(err, ServerKind::WS(ws_socket_addr)))?;
1551
1552 ws_local_addr = Some(addr);
1553 ws_server = Some(server);
1554 }
1555
1556 if let Some(config) = self.http_server_config {
1557 let server = ServerBuilder::new()
1558 .set_config(config.http_only().build())
1559 .set_http_middleware(
1560 tower::ServiceBuilder::new()
1561 .option_layer(Self::maybe_cors_layer(self.http_cors_domains.clone())?)
1562 .option_layer(Self::maybe_jwt_layer(self.jwt_secret))
1563 .option_layer(Self::maybe_decompression_layer(
1564 self.http_max_request_body_size,
1565 self.http_decompression_algorithms.as_deref(),
1566 ))
1567 .option_layer(Self::maybe_compression_layer(
1568 self.http_disable_compression,
1569 self.http_compression_algorithms.as_deref(),
1570 )),
1571 )
1572 .set_rpc_middleware(
1573 RpcServiceBuilder::default()
1574 .option_layer(
1575 rpc_metrics_enabled
1576 .then(|| modules.http.as_ref().map(RpcRequestMetrics::http))
1577 .flatten(),
1578 )
1579 .layer(self.rpc_middleware.clone()),
1580 )
1581 .build(http_socket_addr)
1582 .await
1583 .map_err(|err| RpcError::server_error(err, ServerKind::Http(http_socket_addr)))?;
1584 let local_addr = server
1585 .local_addr()
1586 .map_err(|err| RpcError::server_error(err, ServerKind::Http(http_socket_addr)))?;
1587 http_local_addr = Some(local_addr);
1588 http_server = Some(server);
1589 }
1590
1591 http_handle = http_server
1592 .map(|http_server| http_server.start(modules.http.clone().expect("http server error")));
1593 ws_handle = ws_server
1594 .map(|ws_server| ws_server.start(modules.ws.clone().expect("ws server error")));
1595 Ok(RpcServerHandle {
1596 http_local_addr,
1597 ws_local_addr,
1598 http: http_handle,
1599 ws: ws_handle,
1600 ipc_endpoint: self.ipc_endpoint.clone(),
1601 ipc: ipc_handle,
1602 jwt_secret: self.jwt_secret,
1603 })
1604 }
1605}
1606
1607#[derive(Debug, Clone, Default, Eq, PartialEq)]
1619pub struct TransportRpcModuleConfig {
1620 http: Option<RpcModuleSelection>,
1622 ws: Option<RpcModuleSelection>,
1624 ipc: Option<RpcModuleSelection>,
1626 config: Option<RpcModuleConfig>,
1628}
1629
1630impl TransportRpcModuleConfig {
1633 pub fn set_http(http: impl Into<RpcModuleSelection>) -> Self {
1635 Self::default().with_http(http)
1636 }
1637
1638 pub fn set_ws(ws: impl Into<RpcModuleSelection>) -> Self {
1640 Self::default().with_ws(ws)
1641 }
1642
1643 pub fn set_ipc(ipc: impl Into<RpcModuleSelection>) -> Self {
1645 Self::default().with_ipc(ipc)
1646 }
1647
1648 pub fn with_http(mut self, http: impl Into<RpcModuleSelection>) -> Self {
1650 self.http = Some(http.into());
1651 self
1652 }
1653
1654 pub fn with_ws(mut self, ws: impl Into<RpcModuleSelection>) -> Self {
1656 self.ws = Some(ws.into());
1657 self
1658 }
1659
1660 pub fn with_ipc(mut self, ipc: impl Into<RpcModuleSelection>) -> Self {
1662 self.ipc = Some(ipc.into());
1663 self
1664 }
1665
1666 pub fn with_config(mut self, config: RpcModuleConfig) -> Self {
1668 self.config = Some(config);
1669 self
1670 }
1671
1672 pub const fn http_mut(&mut self) -> &mut Option<RpcModuleSelection> {
1674 &mut self.http
1675 }
1676
1677 pub const fn ws_mut(&mut self) -> &mut Option<RpcModuleSelection> {
1679 &mut self.ws
1680 }
1681
1682 pub const fn ipc_mut(&mut self) -> &mut Option<RpcModuleSelection> {
1684 &mut self.ipc
1685 }
1686
1687 pub const fn config_mut(&mut self) -> &mut Option<RpcModuleConfig> {
1689 &mut self.config
1690 }
1691
1692 pub const fn is_empty(&self) -> bool {
1694 self.http.is_none() && self.ws.is_none() && self.ipc.is_none()
1695 }
1696
1697 pub const fn http(&self) -> Option<&RpcModuleSelection> {
1699 self.http.as_ref()
1700 }
1701
1702 pub const fn ws(&self) -> Option<&RpcModuleSelection> {
1704 self.ws.as_ref()
1705 }
1706
1707 pub const fn ipc(&self) -> Option<&RpcModuleSelection> {
1709 self.ipc.as_ref()
1710 }
1711
1712 pub const fn config(&self) -> Option<&RpcModuleConfig> {
1714 self.config.as_ref()
1715 }
1716
1717 pub fn contains_any(&self, module: &RethRpcModule) -> bool {
1719 self.contains_http(module) || self.contains_ws(module) || self.contains_ipc(module)
1720 }
1721
1722 pub fn contains_http(&self, module: &RethRpcModule) -> bool {
1724 self.http.as_ref().is_some_and(|http| http.contains(module))
1725 }
1726
1727 pub fn contains_ws(&self, module: &RethRpcModule) -> bool {
1729 self.ws.as_ref().is_some_and(|ws| ws.contains(module))
1730 }
1731
1732 pub fn contains_ipc(&self, module: &RethRpcModule) -> bool {
1734 self.ipc.as_ref().is_some_and(|ipc| ipc.contains(module))
1735 }
1736
1737 fn ensure_ws_http_identical(&self) -> Result<(), WsHttpSamePortError> {
1740 if RpcModuleSelection::are_identical(self.http.as_ref(), self.ws.as_ref()) {
1741 Ok(())
1742 } else {
1743 let http_modules =
1744 self.http.as_ref().map(RpcModuleSelection::to_selection).unwrap_or_default();
1745 let ws_modules =
1746 self.ws.as_ref().map(RpcModuleSelection::to_selection).unwrap_or_default();
1747
1748 let http_not_ws = http_modules.difference(&ws_modules).cloned().collect();
1749 let ws_not_http = ws_modules.difference(&http_modules).cloned().collect();
1750 let overlap = http_modules.intersection(&ws_modules).cloned().collect();
1751
1752 Err(WsHttpSamePortError::ConflictingModules(Box::new(ConflictingModules {
1753 overlap,
1754 http_not_ws,
1755 ws_not_http,
1756 })))
1757 }
1758 }
1759}
1760
1761#[derive(Debug, Clone, Default)]
1763pub struct TransportRpcModules<Context = ()> {
1764 config: TransportRpcModuleConfig,
1766 http: Option<RpcModule<Context>>,
1768 ws: Option<RpcModule<Context>>,
1770 ipc: Option<RpcModule<Context>>,
1772}
1773
1774impl TransportRpcModules {
1777 pub fn with_config(mut self, config: TransportRpcModuleConfig) -> Self {
1780 self.config = config;
1781 self
1782 }
1783
1784 pub fn with_http(mut self, http: RpcModule<()>) -> Self {
1787 self.http = Some(http);
1788 self
1789 }
1790
1791 pub fn with_ws(mut self, ws: RpcModule<()>) -> Self {
1794 self.ws = Some(ws);
1795 self
1796 }
1797
1798 pub fn with_ipc(mut self, ipc: RpcModule<()>) -> Self {
1801 self.ipc = Some(ipc);
1802 self
1803 }
1804
1805 pub const fn module_config(&self) -> &TransportRpcModuleConfig {
1807 &self.config
1808 }
1809
1810 pub fn merge_if_module_configured(
1815 &mut self,
1816 module: RethRpcModule,
1817 other: impl Into<Methods>,
1818 ) -> Result<(), RegisterMethodError> {
1819 let other = other.into();
1820 if self.module_config().contains_http(&module) {
1821 self.merge_http(other.clone())?;
1822 }
1823 if self.module_config().contains_ws(&module) {
1824 self.merge_ws(other.clone())?;
1825 }
1826 if self.module_config().contains_ipc(&module) {
1827 self.merge_ipc(other)?;
1828 }
1829
1830 Ok(())
1831 }
1832
1833 pub fn merge_if_module_configured_with<F>(
1840 &mut self,
1841 module: RethRpcModule,
1842 f: F,
1843 ) -> Result<(), RegisterMethodError>
1844 where
1845 F: FnOnce() -> Methods,
1846 {
1847 if !self.module_config().contains_any(&module) {
1849 return Ok(());
1850 }
1851 self.merge_if_module_configured(module, f())
1852 }
1853
1854 pub fn merge_http(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
1860 if let Some(ref mut http) = self.http {
1861 return http.merge(other.into()).map(|_| true)
1862 }
1863 Ok(false)
1864 }
1865
1866 pub fn merge_ws(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
1872 if let Some(ref mut ws) = self.ws {
1873 return ws.merge(other.into()).map(|_| true)
1874 }
1875 Ok(false)
1876 }
1877
1878 pub fn merge_ipc(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
1884 if let Some(ref mut ipc) = self.ipc {
1885 return ipc.merge(other.into()).map(|_| true)
1886 }
1887 Ok(false)
1888 }
1889
1890 pub fn merge_configured(
1894 &mut self,
1895 other: impl Into<Methods>,
1896 ) -> Result<(), RegisterMethodError> {
1897 let other = other.into();
1898 self.merge_http(other.clone())?;
1899 self.merge_ws(other.clone())?;
1900 self.merge_ipc(other)?;
1901 Ok(())
1902 }
1903
1904 pub fn methods_by_module(&self, module: RethRpcModule) -> Methods {
1908 self.methods_by(|name| name.starts_with(module.as_str()))
1909 }
1910
1911 pub fn methods_by<F>(&self, mut filter: F) -> Methods
1915 where
1916 F: FnMut(&str) -> bool,
1917 {
1918 let mut methods = Methods::new();
1919
1920 let mut f =
1922 |name: &str, mm: &Methods| filter(name) && !mm.method_names().any(|m| m == name);
1923
1924 if let Some(m) = self.http_methods(|name| f(name, &methods)) {
1925 let _ = methods.merge(m);
1926 }
1927 if let Some(m) = self.ws_methods(|name| f(name, &methods)) {
1928 let _ = methods.merge(m);
1929 }
1930 if let Some(m) = self.ipc_methods(|name| f(name, &methods)) {
1931 let _ = methods.merge(m);
1932 }
1933 methods
1934 }
1935
1936 pub fn http_methods<F>(&self, filter: F) -> Option<Methods>
1940 where
1941 F: FnMut(&str) -> bool,
1942 {
1943 self.http.as_ref().map(|module| methods_by(module, filter))
1944 }
1945
1946 pub fn ws_methods<F>(&self, filter: F) -> Option<Methods>
1950 where
1951 F: FnMut(&str) -> bool,
1952 {
1953 self.ws.as_ref().map(|module| methods_by(module, filter))
1954 }
1955
1956 pub fn ipc_methods<F>(&self, filter: F) -> Option<Methods>
1960 where
1961 F: FnMut(&str) -> bool,
1962 {
1963 self.ipc.as_ref().map(|module| methods_by(module, filter))
1964 }
1965
1966 pub fn remove_http_method(&mut self, method_name: &'static str) -> bool {
1974 if let Some(http_module) = &mut self.http {
1975 http_module.remove_method(method_name).is_some()
1976 } else {
1977 false
1978 }
1979 }
1980
1981 pub fn remove_http_methods(&mut self, methods: impl IntoIterator<Item = &'static str>) {
1983 for name in methods {
1984 self.remove_http_method(name);
1985 }
1986 }
1987
1988 pub fn remove_ws_method(&mut self, method_name: &'static str) -> bool {
1996 if let Some(ws_module) = &mut self.ws {
1997 ws_module.remove_method(method_name).is_some()
1998 } else {
1999 false
2000 }
2001 }
2002
2003 pub fn remove_ws_methods(&mut self, methods: impl IntoIterator<Item = &'static str>) {
2005 for name in methods {
2006 self.remove_ws_method(name);
2007 }
2008 }
2009
2010 pub fn remove_ipc_method(&mut self, method_name: &'static str) -> bool {
2018 if let Some(ipc_module) = &mut self.ipc {
2019 ipc_module.remove_method(method_name).is_some()
2020 } else {
2021 false
2022 }
2023 }
2024
2025 pub fn remove_ipc_methods(&mut self, methods: impl IntoIterator<Item = &'static str>) {
2027 for name in methods {
2028 self.remove_ipc_method(name);
2029 }
2030 }
2031
2032 pub fn remove_method_from_configured(&mut self, method_name: &'static str) -> bool {
2036 let http_removed = self.remove_http_method(method_name);
2037 let ws_removed = self.remove_ws_method(method_name);
2038 let ipc_removed = self.remove_ipc_method(method_name);
2039
2040 http_removed || ws_removed || ipc_removed
2041 }
2042
2043 pub fn rename(
2047 &mut self,
2048 old_name: &'static str,
2049 new_method: impl Into<Methods>,
2050 ) -> Result<(), RegisterMethodError> {
2051 self.remove_method_from_configured(old_name);
2053
2054 self.merge_configured(new_method)
2056 }
2057
2058 pub fn replace_http(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
2065 let other = other.into();
2066 self.remove_http_methods(other.method_names());
2067 self.merge_http(other)
2068 }
2069
2070 pub fn replace_ipc(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
2077 let other = other.into();
2078 self.remove_ipc_methods(other.method_names());
2079 self.merge_ipc(other)
2080 }
2081
2082 pub fn replace_ws(&mut self, other: impl Into<Methods>) -> Result<bool, RegisterMethodError> {
2089 let other = other.into();
2090 self.remove_ws_methods(other.method_names());
2091 self.merge_ws(other)
2092 }
2093
2094 pub fn replace_configured(
2098 &mut self,
2099 other: impl Into<Methods>,
2100 ) -> Result<bool, RegisterMethodError> {
2101 let other = other.into();
2102 self.replace_http(other.clone())?;
2103 self.replace_ws(other.clone())?;
2104 self.replace_ipc(other)?;
2105 Ok(true)
2106 }
2107
2108 pub fn add_or_replace_http(
2112 &mut self,
2113 other: impl Into<Methods>,
2114 ) -> Result<bool, RegisterMethodError> {
2115 let other = other.into();
2116 self.remove_http_methods(other.method_names());
2117 self.merge_http(other)
2118 }
2119
2120 pub fn add_or_replace_ws(
2124 &mut self,
2125 other: impl Into<Methods>,
2126 ) -> Result<bool, RegisterMethodError> {
2127 let other = other.into();
2128 self.remove_ws_methods(other.method_names());
2129 self.merge_ws(other)
2130 }
2131
2132 pub fn add_or_replace_ipc(
2136 &mut self,
2137 other: impl Into<Methods>,
2138 ) -> Result<bool, RegisterMethodError> {
2139 let other = other.into();
2140 self.remove_ipc_methods(other.method_names());
2141 self.merge_ipc(other)
2142 }
2143
2144 pub fn add_or_replace_configured(
2146 &mut self,
2147 other: impl Into<Methods>,
2148 ) -> Result<(), RegisterMethodError> {
2149 let other = other.into();
2150 self.add_or_replace_http(other.clone())?;
2151 self.add_or_replace_ws(other.clone())?;
2152 self.add_or_replace_ipc(other)?;
2153 Ok(())
2154 }
2155 pub fn add_or_replace_if_module_configured(
2158 &mut self,
2159 module: RethRpcModule,
2160 other: impl Into<Methods>,
2161 ) -> Result<(), RegisterMethodError> {
2162 let other = other.into();
2163 if self.module_config().contains_http(&module) {
2164 self.add_or_replace_http(other.clone())?;
2165 }
2166 if self.module_config().contains_ws(&module) {
2167 self.add_or_replace_ws(other.clone())?;
2168 }
2169 if self.module_config().contains_ipc(&module) {
2170 self.add_or_replace_ipc(other)?;
2171 }
2172 Ok(())
2173 }
2174}
2175
2176fn methods_by<T, F>(module: &RpcModule<T>, mut filter: F) -> Methods
2178where
2179 F: FnMut(&str) -> bool,
2180{
2181 let mut methods = Methods::new();
2182 let method_names = module.method_names().filter(|name| filter(name));
2183
2184 for name in method_names {
2185 if let Some(matched_method) = module.method(name).cloned() {
2186 let _ = methods.verify_and_insert(name, matched_method);
2187 }
2188 }
2189
2190 methods
2191}
2192
2193#[derive(Clone, Debug)]
2198#[must_use = "Server stops if dropped"]
2199pub struct RpcServerHandle {
2200 http_local_addr: Option<SocketAddr>,
2202 ws_local_addr: Option<SocketAddr>,
2203 http: Option<ServerHandle>,
2204 ws: Option<ServerHandle>,
2205 ipc_endpoint: Option<String>,
2206 ipc: Option<jsonrpsee::server::ServerHandle>,
2207 jwt_secret: Option<JwtSecret>,
2208}
2209
2210impl RpcServerHandle {
2213 fn bearer_token(&self) -> Option<String> {
2215 self.jwt_secret.as_ref().map(|secret| {
2216 format!(
2217 "Bearer {}",
2218 secret
2219 .encode(&Claims {
2220 iat: (SystemTime::now().duration_since(UNIX_EPOCH).unwrap() +
2221 Duration::from_secs(60))
2222 .as_secs(),
2223 exp: None,
2224 })
2225 .unwrap()
2226 )
2227 })
2228 }
2229 pub const fn http_local_addr(&self) -> Option<SocketAddr> {
2231 self.http_local_addr
2232 }
2233
2234 pub const fn ws_local_addr(&self) -> Option<SocketAddr> {
2236 self.ws_local_addr
2237 }
2238
2239 pub fn stop(self) -> Result<(), AlreadyStoppedError> {
2241 if let Some(handle) = self.http {
2242 handle.stop()?
2243 }
2244
2245 if let Some(handle) = self.ws {
2246 handle.stop()?
2247 }
2248
2249 if let Some(handle) = self.ipc {
2250 handle.stop()?
2251 }
2252
2253 Ok(())
2254 }
2255
2256 pub fn ipc_endpoint(&self) -> Option<String> {
2258 self.ipc_endpoint.clone()
2259 }
2260
2261 pub fn http_url(&self) -> Option<String> {
2263 self.http_local_addr.map(|addr| format!("http://{addr}"))
2264 }
2265
2266 pub fn ws_url(&self) -> Option<String> {
2268 self.ws_local_addr.map(|addr| format!("ws://{addr}"))
2269 }
2270
2271 pub fn http_client(&self) -> Option<jsonrpsee::http_client::HttpClient> {
2273 let url = self.http_url()?;
2274
2275 let client = if let Some(token) = self.bearer_token() {
2276 jsonrpsee::http_client::HttpClientBuilder::default()
2277 .set_headers(HeaderMap::from_iter([(AUTHORIZATION, token.parse().unwrap())]))
2278 .build(url)
2279 } else {
2280 jsonrpsee::http_client::HttpClientBuilder::default().build(url)
2281 };
2282
2283 client.expect("failed to create http client").into()
2284 }
2285
2286 pub async fn ws_client(&self) -> Option<jsonrpsee::ws_client::WsClient> {
2288 let url = self.ws_url()?;
2289 let mut builder = jsonrpsee::ws_client::WsClientBuilder::default();
2290
2291 if let Some(token) = self.bearer_token() {
2292 let headers = HeaderMap::from_iter([(AUTHORIZATION, token.parse().unwrap())]);
2293 builder = builder.set_headers(headers);
2294 }
2295
2296 let client = builder.build(url).await.expect("failed to create ws client");
2297 Some(client)
2298 }
2299
2300 pub fn eth_http_provider(
2302 &self,
2303 ) -> Option<impl Provider<alloy_network::Ethereum> + Clone + Unpin + 'static> {
2304 self.new_http_provider_for()
2305 }
2306
2307 pub fn eth_http_provider_with_wallet<W>(
2310 &self,
2311 wallet: W,
2312 ) -> Option<impl Provider<alloy_network::Ethereum> + Clone + Unpin + 'static>
2313 where
2314 W: IntoWallet<alloy_network::Ethereum, NetworkWallet: Clone + Unpin + 'static>,
2315 {
2316 let rpc_url = self.http_url()?;
2317 let provider =
2318 ProviderBuilder::new().wallet(wallet).connect_http(rpc_url.parse().expect("valid url"));
2319 Some(provider)
2320 }
2321
2322 pub fn new_http_provider_for<N>(&self) -> Option<impl Provider<N> + Clone + Unpin + 'static>
2327 where
2328 N: RecommendedFillers<RecommendedFillers: Unpin>,
2329 {
2330 let rpc_url = self.http_url()?;
2331 let provider = ProviderBuilder::default()
2332 .with_recommended_fillers()
2333 .connect_http(rpc_url.parse().expect("valid url"));
2334 Some(provider)
2335 }
2336
2337 pub async fn eth_ws_provider(
2339 &self,
2340 ) -> Option<impl Provider<alloy_network::Ethereum> + Clone + Unpin + 'static> {
2341 self.new_ws_provider_for().await
2342 }
2343
2344 pub async fn eth_ws_provider_with_wallet<W>(
2347 &self,
2348 wallet: W,
2349 ) -> Option<impl Provider<alloy_network::Ethereum> + Clone + Unpin + 'static>
2350 where
2351 W: IntoWallet<alloy_network::Ethereum, NetworkWallet: Clone + Unpin + 'static>,
2352 {
2353 let rpc_url = self.ws_url()?;
2354 let provider = ProviderBuilder::new()
2355 .wallet(wallet)
2356 .connect(&rpc_url)
2357 .await
2358 .expect("failed to create ws client");
2359 Some(provider)
2360 }
2361
2362 pub async fn new_ws_provider_for<N>(&self) -> Option<impl Provider<N> + Clone + Unpin + 'static>
2367 where
2368 N: RecommendedFillers<RecommendedFillers: Unpin>,
2369 {
2370 let rpc_url = self.ws_url()?;
2371 let provider = ProviderBuilder::default()
2372 .with_recommended_fillers()
2373 .connect(&rpc_url)
2374 .await
2375 .expect("failed to create ws client");
2376 Some(provider)
2377 }
2378
2379 pub async fn eth_ipc_provider(
2381 &self,
2382 ) -> Option<impl Provider<alloy_network::Ethereum> + Clone + Unpin + 'static> {
2383 self.new_ipc_provider_for().await
2384 }
2385
2386 pub async fn new_ipc_provider_for<N>(
2391 &self,
2392 ) -> Option<impl Provider<N> + Clone + Unpin + 'static>
2393 where
2394 N: RecommendedFillers<RecommendedFillers: Unpin>,
2395 {
2396 let rpc_url = self.ipc_endpoint()?;
2397 let provider = ProviderBuilder::default()
2398 .with_recommended_fillers()
2399 .connect(&rpc_url)
2400 .await
2401 .expect("failed to create ipc client");
2402 Some(provider)
2403 }
2404}
2405
2406#[cfg(test)]
2407mod tests {
2408 use super::*;
2409
2410 #[test]
2411 fn parse_eth_call_bundle_selection() {
2412 let selection = "eth,admin,debug".parse::<RpcModuleSelection>().unwrap();
2413 assert_eq!(
2414 selection,
2415 RpcModuleSelection::Selection(
2416 [RethRpcModule::Eth, RethRpcModule::Admin, RethRpcModule::Debug,].into()
2417 )
2418 );
2419 }
2420
2421 #[test]
2422 fn parse_rpc_module_selection() {
2423 let selection = "all".parse::<RpcModuleSelection>().unwrap();
2424 assert_eq!(selection, RpcModuleSelection::All);
2425 }
2426
2427 #[test]
2428 fn parse_rpc_module_selection_none() {
2429 let selection = "none".parse::<RpcModuleSelection>().unwrap();
2430 assert_eq!(selection, RpcModuleSelection::Selection(Default::default()));
2431 }
2432
2433 #[test]
2434 fn parse_rpc_unique_module_selection() {
2435 let selection = "eth,admin,eth,net".parse::<RpcModuleSelection>().unwrap();
2436 assert_eq!(
2437 selection,
2438 RpcModuleSelection::Selection(
2439 [RethRpcModule::Eth, RethRpcModule::Admin, RethRpcModule::Net,].into()
2440 )
2441 );
2442 }
2443
2444 #[test]
2445 fn identical_selection() {
2446 assert!(RpcModuleSelection::are_identical(
2447 Some(&RpcModuleSelection::All),
2448 Some(&RpcModuleSelection::All),
2449 ));
2450 assert!(!RpcModuleSelection::are_identical(
2451 Some(&RpcModuleSelection::All),
2452 Some(&RpcModuleSelection::Standard),
2453 ));
2454 assert!(RpcModuleSelection::are_identical(
2455 Some(&RpcModuleSelection::Selection(RpcModuleSelection::Standard.to_selection())),
2456 Some(&RpcModuleSelection::Standard),
2457 ));
2458 assert!(RpcModuleSelection::are_identical(
2459 Some(&RpcModuleSelection::Selection([RethRpcModule::Eth].into())),
2460 Some(&RpcModuleSelection::Selection([RethRpcModule::Eth].into())),
2461 ));
2462 assert!(RpcModuleSelection::are_identical(
2463 None,
2464 Some(&RpcModuleSelection::Selection(Default::default())),
2465 ));
2466 assert!(RpcModuleSelection::are_identical(
2467 Some(&RpcModuleSelection::Selection(Default::default())),
2468 None,
2469 ));
2470 assert!(RpcModuleSelection::are_identical(None, None));
2471 }
2472
2473 #[test]
2474 fn test_rpc_module_str() {
2475 macro_rules! assert_rpc_module {
2476 ($($s:expr => $v:expr,)*) => {
2477 $(
2478 let val: RethRpcModule = $s.parse().unwrap();
2479 assert_eq!(val, $v);
2480 assert_eq!(val.to_string(), $s);
2481 )*
2482 };
2483 }
2484 assert_rpc_module!
2485 (
2486 "admin" => RethRpcModule::Admin,
2487 "debug" => RethRpcModule::Debug,
2488 "eth" => RethRpcModule::Eth,
2489 "net" => RethRpcModule::Net,
2490 "trace" => RethRpcModule::Trace,
2491 "web3" => RethRpcModule::Web3,
2492 "rpc" => RethRpcModule::Rpc,
2493 "ots" => RethRpcModule::Ots,
2494 "reth" => RethRpcModule::Reth,
2495 );
2496 }
2497
2498 #[test]
2499 fn test_default_selection() {
2500 let selection = RpcModuleSelection::Standard.to_selection();
2501 assert_eq!(selection, [RethRpcModule::Eth, RethRpcModule::Net, RethRpcModule::Web3].into())
2502 }
2503
2504 #[test]
2505 fn test_create_rpc_module_config() {
2506 let selection = vec!["eth", "admin"];
2507 let config = RpcModuleSelection::try_from_selection(selection).unwrap();
2508 assert_eq!(
2509 config,
2510 RpcModuleSelection::Selection([RethRpcModule::Eth, RethRpcModule::Admin].into())
2511 );
2512 }
2513
2514 #[test]
2515 fn test_configure_transport_config() {
2516 let config = TransportRpcModuleConfig::default()
2517 .with_http([RethRpcModule::Eth, RethRpcModule::Admin]);
2518 assert_eq!(
2519 config,
2520 TransportRpcModuleConfig {
2521 http: Some(RpcModuleSelection::Selection(
2522 [RethRpcModule::Eth, RethRpcModule::Admin].into()
2523 )),
2524 ws: None,
2525 ipc: None,
2526 config: None,
2527 }
2528 )
2529 }
2530
2531 #[test]
2532 fn test_configure_transport_config_none() {
2533 let config = TransportRpcModuleConfig::default().with_http(Vec::<RethRpcModule>::new());
2534 assert_eq!(
2535 config,
2536 TransportRpcModuleConfig {
2537 http: Some(RpcModuleSelection::Selection(Default::default())),
2538 ws: None,
2539 ipc: None,
2540 config: None,
2541 }
2542 )
2543 }
2544
2545 fn create_test_module() -> RpcModule<()> {
2546 let mut module = RpcModule::new(());
2547 module.register_method("anything", |_, _, _| "succeed").unwrap();
2548 module
2549 }
2550
2551 #[test]
2552 fn test_remove_http_method() {
2553 let mut modules =
2554 TransportRpcModules { http: Some(create_test_module()), ..Default::default() };
2555 assert!(modules.remove_http_method("anything"));
2557
2558 assert!(!modules.remove_http_method("non_existent_method"));
2560
2561 assert!(modules.http.as_ref().unwrap().method("anything").is_none());
2563 }
2564
2565 #[test]
2566 fn test_remove_ws_method() {
2567 let mut modules =
2568 TransportRpcModules { ws: Some(create_test_module()), ..Default::default() };
2569
2570 assert!(modules.remove_ws_method("anything"));
2572
2573 assert!(!modules.remove_ws_method("non_existent_method"));
2575
2576 assert!(modules.ws.as_ref().unwrap().method("anything").is_none());
2578 }
2579
2580 #[test]
2581 fn test_remove_ipc_method() {
2582 let mut modules =
2583 TransportRpcModules { ipc: Some(create_test_module()), ..Default::default() };
2584
2585 assert!(modules.remove_ipc_method("anything"));
2587
2588 assert!(!modules.remove_ipc_method("non_existent_method"));
2590
2591 assert!(modules.ipc.as_ref().unwrap().method("anything").is_none());
2593 }
2594
2595 #[test]
2596 fn test_remove_method_from_configured() {
2597 let mut modules = TransportRpcModules {
2598 http: Some(create_test_module()),
2599 ws: Some(create_test_module()),
2600 ipc: Some(create_test_module()),
2601 ..Default::default()
2602 };
2603
2604 assert!(modules.remove_method_from_configured("anything"));
2606
2607 assert!(!modules.remove_method_from_configured("anything"));
2609
2610 assert!(!modules.remove_method_from_configured("non_existent_method"));
2612
2613 assert!(modules.http.as_ref().unwrap().method("anything").is_none());
2615 assert!(modules.ws.as_ref().unwrap().method("anything").is_none());
2616 assert!(modules.ipc.as_ref().unwrap().method("anything").is_none());
2617 }
2618
2619 #[test]
2620 fn test_transport_rpc_module_rename() {
2621 let mut modules = TransportRpcModules {
2622 http: Some(create_test_module()),
2623 ws: Some(create_test_module()),
2624 ipc: Some(create_test_module()),
2625 ..Default::default()
2626 };
2627
2628 assert!(modules.http.as_ref().unwrap().method("anything").is_some());
2630 assert!(modules.ws.as_ref().unwrap().method("anything").is_some());
2631 assert!(modules.ipc.as_ref().unwrap().method("anything").is_some());
2632
2633 assert!(modules.http.as_ref().unwrap().method("something").is_none());
2635 assert!(modules.ws.as_ref().unwrap().method("something").is_none());
2636 assert!(modules.ipc.as_ref().unwrap().method("something").is_none());
2637
2638 let mut other_module = RpcModule::new(());
2640 other_module.register_method("something", |_, _, _| "fails").unwrap();
2641
2642 modules.rename("anything", other_module).expect("rename failed");
2644
2645 assert!(modules.http.as_ref().unwrap().method("anything").is_none());
2647 assert!(modules.ws.as_ref().unwrap().method("anything").is_none());
2648 assert!(modules.ipc.as_ref().unwrap().method("anything").is_none());
2649
2650 assert!(modules.http.as_ref().unwrap().method("something").is_some());
2652 assert!(modules.ws.as_ref().unwrap().method("something").is_some());
2653 assert!(modules.ipc.as_ref().unwrap().method("something").is_some());
2654 }
2655
2656 #[test]
2657 fn test_replace_http_method() {
2658 let mut modules =
2659 TransportRpcModules { http: Some(create_test_module()), ..Default::default() };
2660
2661 let mut other_module = RpcModule::new(());
2662 other_module.register_method("something", |_, _, _| "fails").unwrap();
2663
2664 assert!(modules.replace_http(other_module.clone()).unwrap());
2665
2666 assert!(modules.http.as_ref().unwrap().method("something").is_some());
2667
2668 other_module.register_method("anything", |_, _, _| "fails").unwrap();
2669 assert!(modules.replace_http(other_module.clone()).unwrap());
2670
2671 assert!(modules.http.as_ref().unwrap().method("anything").is_some());
2672 }
2673 #[test]
2674 fn test_replace_ipc_method() {
2675 let mut modules =
2676 TransportRpcModules { ipc: Some(create_test_module()), ..Default::default() };
2677
2678 let mut other_module = RpcModule::new(());
2679 other_module.register_method("something", |_, _, _| "fails").unwrap();
2680
2681 assert!(modules.replace_ipc(other_module.clone()).unwrap());
2682
2683 assert!(modules.ipc.as_ref().unwrap().method("something").is_some());
2684
2685 other_module.register_method("anything", |_, _, _| "fails").unwrap();
2686 assert!(modules.replace_ipc(other_module.clone()).unwrap());
2687
2688 assert!(modules.ipc.as_ref().unwrap().method("anything").is_some());
2689 }
2690 #[test]
2691 fn test_replace_ws_method() {
2692 let mut modules =
2693 TransportRpcModules { ws: Some(create_test_module()), ..Default::default() };
2694
2695 let mut other_module = RpcModule::new(());
2696 other_module.register_method("something", |_, _, _| "fails").unwrap();
2697
2698 assert!(modules.replace_ws(other_module.clone()).unwrap());
2699
2700 assert!(modules.ws.as_ref().unwrap().method("something").is_some());
2701
2702 other_module.register_method("anything", |_, _, _| "fails").unwrap();
2703 assert!(modules.replace_ws(other_module.clone()).unwrap());
2704
2705 assert!(modules.ws.as_ref().unwrap().method("anything").is_some());
2706 }
2707
2708 #[test]
2709 fn test_replace_configured() {
2710 let mut modules = TransportRpcModules {
2711 http: Some(create_test_module()),
2712 ws: Some(create_test_module()),
2713 ipc: Some(create_test_module()),
2714 ..Default::default()
2715 };
2716 let mut other_module = RpcModule::new(());
2717 other_module.register_method("something", |_, _, _| "fails").unwrap();
2718
2719 assert!(modules.replace_configured(other_module).unwrap());
2720
2721 assert!(modules.http.as_ref().unwrap().method("something").is_some());
2723 assert!(modules.ipc.as_ref().unwrap().method("something").is_some());
2724 assert!(modules.ws.as_ref().unwrap().method("something").is_some());
2725
2726 assert!(modules.http.as_ref().unwrap().method("anything").is_some());
2727 assert!(modules.ipc.as_ref().unwrap().method("anything").is_some());
2728 assert!(modules.ws.as_ref().unwrap().method("anything").is_some());
2729 }
2730
2731 #[test]
2732 fn test_add_or_replace_if_module_configured() {
2733 let config = TransportRpcModuleConfig::default()
2735 .with_http([RethRpcModule::Eth])
2736 .with_ws([RethRpcModule::Eth]);
2737
2738 let mut http_module = RpcModule::new(());
2740 http_module.register_method("eth_existing", |_, _, _| "original").unwrap();
2741
2742 let mut ws_module = RpcModule::new(());
2744 ws_module.register_method("eth_existing", |_, _, _| "original").unwrap();
2745
2746 let ipc_module = RpcModule::new(());
2748
2749 let mut modules = TransportRpcModules {
2751 config,
2752 http: Some(http_module),
2753 ws: Some(ws_module),
2754 ipc: Some(ipc_module),
2755 };
2756
2757 let mut new_module = RpcModule::new(());
2759 new_module.register_method("eth_existing", |_, _, _| "replaced").unwrap(); new_module.register_method("eth_new", |_, _, _| "added").unwrap(); let new_methods: Methods = new_module.into();
2762
2763 let result = modules.add_or_replace_if_module_configured(RethRpcModule::Eth, new_methods);
2765 assert!(result.is_ok(), "Function should succeed");
2766
2767 let http = modules.http.as_ref().unwrap();
2769 assert!(http.method("eth_existing").is_some());
2770 assert!(http.method("eth_new").is_some());
2771
2772 let ws = modules.ws.as_ref().unwrap();
2774 assert!(ws.method("eth_existing").is_some());
2775 assert!(ws.method("eth_new").is_some());
2776
2777 let ipc = modules.ipc.as_ref().unwrap();
2779 assert!(ipc.method("eth_existing").is_none());
2780 assert!(ipc.method("eth_new").is_none());
2781 }
2782
2783 #[test]
2784 fn test_merge_if_module_configured_with_lazy_evaluation() {
2785 let config = TransportRpcModuleConfig::default().with_http([RethRpcModule::Eth]);
2787
2788 let mut modules =
2789 TransportRpcModules { config, http: Some(RpcModule::new(())), ws: None, ipc: None };
2790
2791 let mut closure_called = false;
2793
2794 let result = modules.merge_if_module_configured_with(RethRpcModule::Eth, || {
2796 closure_called = true;
2797 let mut methods = RpcModule::new(());
2798 methods.register_method("eth_test", |_, _, _| "test").unwrap();
2799 methods.into()
2800 });
2801
2802 assert!(result.is_ok());
2803 assert!(closure_called, "Closure should be called when module is configured");
2804 assert!(modules.http.as_ref().unwrap().method("eth_test").is_some());
2805
2806 closure_called = false;
2808 let result = modules.merge_if_module_configured_with(RethRpcModule::Debug, || {
2809 closure_called = true;
2810 RpcModule::new(()).into()
2811 });
2812
2813 assert!(result.is_ok());
2814 assert!(!closure_called, "Closure should NOT be called when module is not configured");
2815 }
2816}