1use std::sync::Arc;
9
10use crossbeam_channel::Sender;
11use futures::Stream;
12use reth_consensus::FullConsensus;
13use reth_engine_primitives::BeaconEngineMessage;
14use reth_engine_tree::{
15 backfill::PipelineSync,
16 chain::ChainOrchestrator,
17 download::BasicBlockDownloader,
18 engine::{EngineApiKind, EngineApiRequest, EngineApiRequestHandler, EngineHandler, FromEngine},
19 persistence::PersistenceHandle,
20 tree::{EngineApiTreeHandler, EngineValidator, TreeConfig, WaitForCaches},
21};
22use reth_evm::ConfigureEvm;
23use reth_network_p2p::BlockClient;
24use reth_payload_builder::PayloadBuilderHandle;
25use reth_primitives_traits::NodePrimitives;
26use reth_provider::{
27 ProviderFactory, StorageSettingsCache,
28 providers::{BlockchainProvider, ProviderNodeTypes},
29};
30use reth_prune::PrunerWithFactory;
31use reth_stages_api::{MetricEventsSender, Pipeline};
32use reth_tasks::Runtime;
33use reth_trie_db::ChangesetCache;
34
35pub type TreeSender<T, N> =
37 Sender<FromEngine<EngineApiRequest<T, N>, <N as NodePrimitives>::Block>>;
38
39#[expect(clippy::too_many_arguments, clippy::type_complexity)]
46pub fn build_arb_engine_orchestrator<N, Client, S, V, C>(
47 engine_kind: EngineApiKind,
48 consensus: Arc<dyn FullConsensus<N::Primitives>>,
49 client: Client,
50 incoming_requests: S,
51 pipeline: Pipeline<N>,
52 pipeline_task_spawner: Runtime,
53 provider: ProviderFactory<N>,
54 blockchain_db: BlockchainProvider<N>,
55 pruner: PrunerWithFactory<ProviderFactory<N>>,
56 payload_builder: PayloadBuilderHandle<N::Payload>,
57 payload_validator: V,
58 tree_config: TreeConfig,
59 sync_metrics_tx: MetricEventsSender,
60 evm_config: C,
61 changeset_cache: ChangesetCache,
62) -> (
63 ChainOrchestrator<
64 EngineHandler<
65 EngineApiRequestHandler<EngineApiRequest<N::Payload, N::Primitives>, N::Primitives>,
66 S,
67 BasicBlockDownloader<Client, <N::Primitives as NodePrimitives>::Block>,
68 >,
69 PipelineSync<N>,
70 >,
71 TreeSender<N::Payload, N::Primitives>,
72)
73where
74 N: ProviderNodeTypes,
75 Client: BlockClient<Block = <N::Primitives as NodePrimitives>::Block> + 'static,
76 S: Stream<Item = BeaconEngineMessage<N::Payload>> + Send + Sync + Unpin + 'static,
77 V: EngineValidator<N::Payload> + WaitForCaches,
78 C: ConfigureEvm<Primitives = N::Primitives> + 'static,
79{
80 let downloader = BasicBlockDownloader::new(client, consensus.clone());
81 let use_hashed_state = provider.cached_storage_settings().use_hashed_state();
82
83 let persistence_handle =
84 PersistenceHandle::<N::Primitives>::spawn_service(provider, pruner, sync_metrics_tx);
85
86 let canonical_in_memory_state = blockchain_db.canonical_in_memory_state();
87
88 let (to_tree_tx, from_tree) = EngineApiTreeHandler::spawn_new(
89 blockchain_db,
90 consensus,
91 payload_validator,
92 persistence_handle,
93 payload_builder,
94 canonical_in_memory_state,
95 tree_config,
96 engine_kind,
97 evm_config,
98 changeset_cache,
99 use_hashed_state,
100 );
101
102 let tree_sender = to_tree_tx.clone();
105
106 let engine_handler = EngineApiRequestHandler::new(to_tree_tx, from_tree);
107 let handler = EngineHandler::new(engine_handler, downloader, incoming_requests);
108
109 let backfill_sync = PipelineSync::new(pipeline, pipeline_task_spawner);
110
111 (ChainOrchestrator::new(handler, backfill_sync), tree_sender)
112}