Skip to main content

reth_engine_tree/
launch.rs

1//! Engine orchestrator launch helper.
2//!
3//! Provides [`EngineOrchestratorBuilder`] which wires
4//! together all engine components and builds a
5//! [`ChainOrchestrator`] ready to be polled as a `Stream`.
6
7use crate::{
8    backfill::PipelineSync,
9    chain::ChainOrchestrator,
10    download::BasicBlockDownloader,
11    engine::{EngineApiKind, EngineApiRequest, EngineApiRequestHandler, EngineHandler},
12    persistence::PersistenceHandle,
13    tree::{EngineApiTreeHandler, EngineValidator, TreeConfig, WaitForCaches},
14};
15use futures::Stream;
16use reth_consensus::FullConsensus;
17use reth_engine_primitives::BeaconEngineMessage;
18use reth_evm::ConfigureEvm;
19use reth_network_p2p::{BlockAccessListsClient, BlockClient};
20use reth_payload_builder::PayloadBuilderHandle;
21use reth_primitives_traits::NodePrimitives;
22use reth_provider::{
23    providers::{BlockchainProvider, ProviderNodeTypes},
24    ProviderFactory,
25};
26use reth_prune::PrunerWithFactory;
27use reth_stages_api::{MetricEventsSender, Pipeline};
28use reth_storage_overlay::OverlayManager;
29use reth_tasks::Runtime;
30use std::sync::Arc;
31
32/// The [`ChainOrchestrator`] built by [`EngineOrchestratorBuilder`].
33pub type EngineOrchestrator<T, N, Client, S, B> = ChainOrchestrator<
34    EngineHandler<
35        EngineApiRequestHandler<EngineApiRequest<T, N>, N>,
36        S,
37        BasicBlockDownloader<Client, <N as NodePrimitives>::Block>,
38    >,
39    B,
40>;
41
42/// Components needed to build the engine [`ChainOrchestrator`] that drives the chain forward.
43///
44/// [`build`](Self::build) spawns and wires together the following components:
45///
46/// - **[`BasicBlockDownloader`]** — downloads blocks on demand from the network during live sync.
47/// - **[`PersistenceHandle`]** — spawns the persistence service on a background thread for writing
48///   blocks and performing pruning outside the critical consensus path.
49/// - **[`EngineApiTreeHandler`]** — spawns the tree handler that processes engine API requests
50///   (`newPayload`, `forkchoiceUpdated`) and maintains the in-memory chain state.
51/// - **[`EngineApiRequestHandler`]** + **[`EngineHandler`]** — glue that routes incoming CL
52///   messages to the tree handler and manages download requests.
53/// - **[`PipelineSync`]** — wraps the staged sync [`Pipeline`] for backfill sync when the node
54///   needs to catch up over large block ranges.
55///
56/// The returned orchestrator implements [`Stream`] and yields
57/// [`ChainEvent`]s.
58///
59/// [`ChainEvent`]: crate::chain::ChainEvent
60#[derive(Debug)]
61pub struct EngineOrchestratorBuilder<N, Client, S, V, C>
62where
63    N: ProviderNodeTypes,
64{
65    /// The engine API flavor to run.
66    pub engine_kind: EngineApiKind,
67    /// Consensus used to validate downloaded and incoming blocks.
68    pub consensus: Arc<dyn FullConsensus<N::Primitives>>,
69    /// Network client used to download blocks during live sync.
70    pub client: Client,
71    /// Incoming consensus layer messages.
72    pub incoming_requests: S,
73    /// Staged sync pipeline used for backfill.
74    pub pipeline: Pipeline<N>,
75    /// Runtime the pipeline runs on.
76    pub pipeline_task_spawner: Runtime,
77    /// Provider factory handed to the persistence service.
78    pub provider: ProviderFactory<N>,
79    /// Blockchain provider backing the engine tree.
80    pub blockchain_db: BlockchainProvider<N>,
81    /// Pruner run by the persistence service.
82    pub pruner: PrunerWithFactory<ProviderFactory<N>>,
83    /// Handle to the payload builder service.
84    pub payload_builder: PayloadBuilderHandle<N::Payload>,
85    /// Validator for incoming payloads.
86    pub payload_validator: V,
87    /// Overlay manager for state on top of the database.
88    pub overlay_manager: OverlayManager<N::Primitives>,
89    /// Engine tree configuration.
90    pub tree_config: TreeConfig,
91    /// Sender for sync metric events.
92    pub sync_metrics_tx: MetricEventsSender,
93    /// EVM configuration used to execute payloads.
94    pub evm_config: C,
95    /// Runtime used to spawn engine tree tasks.
96    pub runtime: Runtime,
97}
98
99impl<N, Client, S, V, C> EngineOrchestratorBuilder<N, Client, S, V, C>
100where
101    N: ProviderNodeTypes,
102    Client: BlockClient<Block = <N::Primitives as NodePrimitives>::Block>
103        + BlockAccessListsClient
104        + 'static,
105    S: Stream<Item = BeaconEngineMessage<N::Payload>> + Send + Sync + Unpin + 'static,
106    V: EngineValidator<N::Payload> + WaitForCaches,
107    C: ConfigureEvm<Primitives = N::Primitives> + 'static,
108{
109    /// Spawns the engine services and returns the [`ChainOrchestrator`] driving them.
110    pub fn build(
111        self,
112    ) -> EngineOrchestrator<N::Payload, N::Primitives, Client, S, PipelineSync<N>> {
113        let Self {
114            engine_kind,
115            consensus,
116            client,
117            incoming_requests,
118            pipeline,
119            pipeline_task_spawner,
120            provider,
121            blockchain_db,
122            pruner,
123            payload_builder,
124            payload_validator,
125            overlay_manager,
126            tree_config,
127            sync_metrics_tx,
128            evm_config,
129            runtime,
130        } = self;
131
132        let downloader = BasicBlockDownloader::new(client, consensus.clone());
133
134        let persistence_handle =
135            PersistenceHandle::<N::Primitives>::spawn_service(provider, pruner, sync_metrics_tx);
136
137        let canonical_in_memory_state = blockchain_db.canonical_in_memory_state();
138
139        let (to_tree_tx, from_tree) = EngineApiTreeHandler::spawn_new(
140            blockchain_db,
141            consensus,
142            payload_validator,
143            persistence_handle,
144            payload_builder,
145            canonical_in_memory_state,
146            overlay_manager,
147            tree_config,
148            engine_kind,
149            evm_config,
150            runtime,
151        );
152
153        let engine_handler = EngineApiRequestHandler::new(to_tree_tx, from_tree);
154        let handler = EngineHandler::new(engine_handler, downloader, incoming_requests);
155
156        let backfill_sync = PipelineSync::new(pipeline, pipeline_task_spawner);
157
158        ChainOrchestrator::new(handler, backfill_sync)
159    }
160}