1use crate::{
4 common::{Attached, LaunchContextWith, WithConfigs},
5 hooks::NodeHooks,
6 rpc::{EngineShutdown, EngineValidatorAddOn, EngineValidatorBuilder, RethRpcAddOns, RpcHandle},
7 setup::build_networked_pipeline,
8 AddOns, AddOnsContext, FullNode, LaunchContext, LaunchNode, Node, NodeAdapter,
9 NodeBuilderWithComponents, NodeComponents, NodeComponentsBuilder, NodeHandle, NodeTypesAdapter,
10 RethFullAdapter,
11};
12use alloy_consensus::BlockHeader;
13use futures::{stream::FusedStream, stream_select, FutureExt, StreamExt};
14use reth_chainspec::{EthChainSpec, EthereumHardforks};
15use reth_db::{database_metrics::DatabaseMetrics, Database};
16use reth_engine_tree::{
17 chain::{ChainEvent, FromOrchestrator},
18 engine::{EngineApiKind, EngineApiRequest, EngineRequestHandler},
19 launch::EngineOrchestratorBuilder,
20 tree::TreeConfig,
21};
22use reth_engine_util::EngineMessageStreamExt;
23use reth_exex::ExExManagerHandle;
24use reth_network::{types::BlockRangeUpdate, NetworkSyncUpdater, SyncState};
25use reth_network_api::BlockDownloaderProvider;
26use reth_node_api::{
27 BuiltPayload, ConsensusEngineHandle, FullNodeTypes, NodeTypes, NodeTypesWithDBAdapter,
28};
29use reth_node_core::{
30 args::PruneConfigKind,
31 dirs::{ChainPath, DataDirPath},
32 exit::NodeExitFuture,
33 primitives::Head,
34};
35use reth_node_events::node;
36use reth_provider::{
37 providers::{BlockchainProvider, NodeTypesForProvider},
38 BlockNumReader, StorageSettingsCache,
39};
40use reth_storage_overlay::OverlayManager;
41use reth_tasks::TaskExecutor;
42use reth_tokio_util::EventSender;
43use reth_tracing::tracing::{debug, error, info};
44use std::{future::Future, pin::Pin, sync::Arc};
45use tokio::sync::{mpsc::unbounded_channel, oneshot};
46use tokio_stream::wrappers::UnboundedReceiverStream;
47
48#[derive(Debug)]
50pub struct EngineNodeLauncher {
51 pub ctx: LaunchContext,
53
54 pub engine_tree_config: TreeConfig,
57}
58
59impl EngineNodeLauncher {
60 pub const fn new(
62 task_executor: TaskExecutor,
63 data_dir: ChainPath<DataDirPath>,
64 engine_tree_config: TreeConfig,
65 ) -> Self {
66 Self { ctx: LaunchContext::new(task_executor, data_dir), engine_tree_config }
67 }
68
69 async fn launch_node<N, DB, T, CB, AO>(
70 self,
71 target: NodeBuilderWithComponents<T, CB, AO>,
72 ) -> eyre::Result<NodeHandle<NodeAdapter<T, CB::Components>, AO>>
73 where
74 N: Node<RethFullAdapter<DB, N>> + NodeTypesForProvider,
75 DB: Database + DatabaseMetrics + Clone + Unpin + 'static,
76 T: FullNodeTypes<
77 Types = N,
78 Provider = BlockchainProvider<NodeTypesWithDBAdapter<N, DB>>,
79 DB = DB,
80 >,
81 CB: NodeComponentsBuilder<T>,
82 AO: RethRpcAddOns<NodeAdapter<T, CB::Components>>
83 + EngineValidatorAddOn<NodeAdapter<T, CB::Components>>,
84 {
85 let Self { ctx, engine_tree_config } = self;
86 let NodeBuilderWithComponents {
87 adapter: NodeTypesAdapter { database },
88 rocksdb_provider,
89 components_builder,
90 add_ons: AddOns { hooks, exexs: installed_exex, add_ons },
91 config,
92 } = target;
93 let NodeHooks { on_component_initialized, on_node_started, .. } = hooks;
94
95 let overlay_manager = OverlayManager::<N::Primitives>::new(
97 ctx.task_executor.state_trie_overlay_worker_pool(),
98 );
99 let disabled_stages = N::disabled_stages();
100
101 let ctx = ctx
103 .with_configured_globals(engine_tree_config.reserved_cpu_cores())
104 .with_loaded_toml_config(config)?
106 .with_resolved_peers()?
108 .attach(database.clone())
110 .with_adjusted_configs()
112 .with_provider_factory::<_, <CB::Components as NodeComponents<T>>::Evm>(
114 overlay_manager.clone(),
115 rocksdb_provider,
116 disabled_stages,
117 )
118 .await?
119 .inspect(|_| {
120 info!(target: "reth::cli", "Database opened");
121 })
122 .with_prometheus_server().await?
123 .inspect(|this| {
124 debug!(target: "reth::cli", chain=%this.chain_id(), genesis=?this.genesis_hash(), "Initializing genesis");
125 })
126 .with_genesis()?
127 .inspect(|this: &LaunchContextWith<Attached<WithConfigs<<T::Types as NodeTypes>::ChainSpec>, _>>| {
128 info!(target: "reth::cli", "\n{}", this.chain_spec().display_hardforks());
129 let settings = this.provider_factory().cached_storage_settings();
130 let pruning_mode =
131 PruneConfigKind::from_config(&this.prune_config(), this.chain_spec().as_ref()).as_str();
132 info!(target: "reth::cli", ?settings, ?pruning_mode, "Loaded storage settings");
133 })
134 .with_metrics_task()
135 .with_blockchain_db::<T, _>(move |provider_factory| {
138 Ok(BlockchainProvider::new(provider_factory)?)
139 })?
140 .with_components(components_builder, on_component_initialized).await?;
141
142 let maybe_exex_manager_handle = ctx.launch_exex(installed_exex).await?;
144
145 let network_handle = ctx.components().network().clone();
147 let network_client = network_handle.fetch_client().await?;
148 let (consensus_engine_tx, consensus_engine_rx) = unbounded_channel();
149
150 let node_config = ctx.node_config();
151
152 network_handle.update_sync_state(SyncState::Syncing);
154
155 let max_block = ctx.max_block(network_client.clone()).await?;
156
157 let static_file_producer = ctx.static_file_producer();
158 let static_file_producer_events = static_file_producer.lock().events();
159 info!(target: "reth::cli", "StaticFileProducer initialized");
160
161 let consensus = Arc::new(ctx.components().consensus().clone());
162
163 let pipeline = build_networked_pipeline(
164 &ctx.toml_config().stages,
165 network_client.clone(),
166 consensus.clone(),
167 ctx.provider_factory().clone(),
168 ctx.task_executor(),
169 ctx.sync_metrics_tx(),
170 ctx.prune_config(),
171 max_block,
172 static_file_producer,
173 ctx.components().evm_config().clone(),
174 maybe_exex_manager_handle.clone().unwrap_or_else(ExExManagerHandle::empty),
175 ctx.era_import_source(),
176 disabled_stages,
177 )?;
178
179 pipeline.move_to_static_files()?;
181
182 let pipeline_events = pipeline.events();
183
184 let mut pruner_builder = ctx.pruner_builder();
185 if let Some(exex_manager_handle) = &maybe_exex_manager_handle {
186 pruner_builder =
187 pruner_builder.finished_exex_height(exex_manager_handle.finished_height());
188 }
189 let pruner = pruner_builder.build_with_provider_factory(ctx.provider_factory().clone());
190 let pruner_events = pruner.events();
191 info!(target: "reth::cli", prune_config=?ctx.prune_config(), "Pruner initialized");
192
193 let event_sender = EventSender::default();
194
195 let beacon_engine_handle = ConsensusEngineHandle::new(consensus_engine_tx.clone());
196
197 let jwt_secret = ctx.auth_jwt_secret()?;
199
200 let add_ons_ctx = AddOnsContext {
201 node: ctx.node_adapter().clone(),
202 config: ctx.node_config(),
203 beacon_engine_handle: beacon_engine_handle.clone(),
204 jwt_secret,
205 engine_events: event_sender.clone(),
206 sender_recovery_cache: ctx.sender_recovery_cache().cloned(),
207 };
208 let validator_builder = add_ons.engine_validator_builder();
209
210 let engine_validator = validator_builder
212 .clone()
213 .build_tree_validator(&add_ons_ctx, engine_tree_config.clone(), overlay_manager.clone())
214 .await?;
215
216 let reorg_overlay_manager = overlay_manager.clone();
218 let consensus_engine_stream = UnboundedReceiverStream::from(consensus_engine_rx)
219 .maybe_skip_fcu(node_config.debug.skip_fcu)
220 .maybe_skip_new_payload(node_config.debug.skip_new_payload)
221 .maybe_reorg(
222 ctx.blockchain_db().clone(),
223 ctx.components().evm_config().clone(),
224 || async {
225 validator_builder
226 .build_tree_validator(
227 &add_ons_ctx,
228 engine_tree_config.clone(),
229 reorg_overlay_manager.clone(),
230 )
231 .await
232 },
233 node_config.debug.reorg_frequency,
234 node_config.debug.reorg_depth,
235 )
236 .await?
237 .maybe_store_messages(node_config.debug.engine_api_store.clone());
241
242 let engine_kind = if ctx.chain_spec().is_optimism() {
243 EngineApiKind::OpStack
244 } else {
245 EngineApiKind::Ethereum
246 };
247
248 let mut orchestrator = EngineOrchestratorBuilder {
249 engine_kind,
250 consensus,
251 client: network_client,
252 incoming_requests: Box::pin(consensus_engine_stream),
253 pipeline,
254 pipeline_task_spawner: ctx.task_executor().clone(),
255 provider: ctx.provider_factory().clone(),
256 blockchain_db: ctx.blockchain_db().clone(),
257 pruner,
258 payload_builder: ctx.components().payload_builder_handle().clone(),
259 payload_validator: engine_validator,
260 overlay_manager,
261 tree_config: engine_tree_config,
262 sync_metrics_tx: ctx.sync_metrics_tx(),
263 evm_config: ctx.components().evm_config().clone(),
264 runtime: ctx.task_executor().clone(),
265 }
266 .build();
267
268 info!(target: "reth::cli", "Consensus engine initialized");
269
270 #[expect(clippy::needless_continue)]
271 let events = stream_select!(
272 event_sender.new_listener().map(Into::into),
273 pipeline_events.map(Into::into),
274 ctx.consensus_layer_events(),
275 pruner_events.map(Into::into),
276 static_file_producer_events.map(Into::into),
277 );
278
279 ctx.task_executor().spawn_critical_task(
280 "events task",
281 node::handle_events(
282 Some(Box::new(ctx.components().network().clone())),
283 Some(ctx.head().number),
284 events,
285 ),
286 );
287
288 let RpcHandle {
289 rpc_server_handles,
290 rpc_registry,
291 engine_events,
292 beacon_engine_handle,
293 engine_shutdown: _,
294 } = add_ons.launch_add_ons(add_ons_ctx).await?;
295
296 let (engine_shutdown, shutdown_rx) = EngineShutdown::new();
298
299 let initial_target = ctx.initial_backfill_target(disabled_stages)?;
301 let mut built_payloads = ctx
302 .components()
303 .payload_builder_handle()
304 .subscribe()
305 .await
306 .map_err(|e| eyre::eyre!("Failed to subscribe to payload builder events: {:?}", e))?
307 .into_built_payload_stream()
308 .fuse();
309
310 let chainspec = ctx.chain_spec();
311 let provider = ctx.blockchain_db().clone();
312 let (exit, rx) = oneshot::channel();
313 let terminate_after_backfill = ctx.terminate_after_initial_backfill();
314 let startup_sync_state_idle = ctx.node_config().debug.startup_sync_state_idle;
315
316 info!(target: "reth::cli", "Starting consensus engine");
317 let consensus_engine = move |mut on_graceful_shutdown| async move {
318 if let Some(initial_target) = initial_target {
319 debug!(target: "reth::cli", %initial_target, "start backfill sync");
320 orchestrator.start_backfill_sync(initial_target);
322 } else if startup_sync_state_idle {
323 network_handle.update_sync_state(SyncState::Idle);
324 }
325
326 let mut res = Ok(());
327 let mut shutdown_rx = shutdown_rx.fuse();
328
329 loop {
333 tokio::select! {
334 event = orchestrator.next() => {
335 let Some(event) = event else { break };
336 debug!(target: "reth::cli", "Event: {event}");
337 match event {
338 ChainEvent::BackfillSyncFinished => {
339 if terminate_after_backfill {
340 debug!(target: "reth::cli", "Terminating after initial backfill");
341 break
342 }
343 if startup_sync_state_idle {
344 network_handle.update_sync_state(SyncState::Idle);
345 }
346 }
347 ChainEvent::BackfillSyncStarted => {
348 network_handle.update_sync_state(SyncState::Syncing);
349 }
350 ChainEvent::FatalError => {
351 error!(target: "reth::cli", "Fatal error in consensus engine");
352 res = Err(eyre::eyre!("Fatal error in consensus engine"));
353 break
354 }
355 ChainEvent::Handler(ev) => {
356 if let Some(head) = ev.canonical_header() {
357 network_handle.update_sync_state(SyncState::Idle);
359 let head_block = Head {
360 number: head.number(),
361 hash: head.hash(),
362 difficulty: head.difficulty(),
363 timestamp: head.timestamp(),
364 total_difficulty: chainspec.final_paris_total_difficulty()
365 .filter(|_| chainspec.is_paris_active_at_block(head.number()))
366 .unwrap_or_default(),
367 };
368 network_handle.update_status(head_block);
369
370 let updated = BlockRangeUpdate {
371 earliest: provider.earliest_block_number().unwrap_or_default(),
372 latest: head.number(),
373 latest_hash: head.hash(),
374 };
375 network_handle.update_block_range(updated);
376 }
377 event_sender.notify(ev);
378 }
379 }
380 }
381 Some(payload) = built_payloads.next(), if !built_payloads.is_terminated() => {
382 if let Some(executed_block) = payload.executed_block() {
383 debug!(target: "reth::cli", block=?executed_block.recovered_block.num_hash(), "inserting built payload");
384 orchestrator.handler_mut().handler_mut().on_event(EngineApiRequest::InsertExecutedBlock(executed_block).into());
385 }
386 }
387 shutdown_req = &mut shutdown_rx => {
388 if let Ok(req) = shutdown_req {
389 debug!(target: "reth::cli", "received engine shutdown request");
390 orchestrator.handler_mut().handler_mut().on_event(
391 FromOrchestrator::Terminate { tx: req.done_tx }.into()
392 );
393 }
394 }
395 _guard = &mut on_graceful_shutdown => {
396 debug!(target: "reth::cli", "shutdown signal received, terminating engine");
400 let (done_tx, done_rx) = oneshot::channel();
401 orchestrator.handler_mut().handler_mut().on_event(
402 FromOrchestrator::Terminate { tx: done_tx }.into()
403 );
404 let _ = done_rx.await;
405 break;
406 }
407 }
408 }
409
410 let _ = exit.send(res);
411 };
412 ctx.task_executor()
413 .spawn_critical_with_graceful_shutdown_signal("consensus engine", consensus_engine);
414
415 let engine_events_for_ethstats = engine_events.new_listener();
416
417 let full_node = FullNode {
418 evm_config: ctx.components().evm_config().clone(),
419 pool: ctx.components().pool().clone(),
420 network: ctx.components().network().clone(),
421 provider: ctx.node_adapter().provider.clone(),
422 payload_builder_handle: ctx.components().payload_builder_handle().clone(),
423 task_executor: ctx.task_executor().clone(),
424 config: ctx.node_config().clone(),
425 data_dir: ctx.data_dir().clone(),
426 add_ons_handle: RpcHandle {
427 rpc_server_handles,
428 rpc_registry,
429 engine_events,
430 beacon_engine_handle,
431 engine_shutdown,
432 },
433 };
434 on_node_started.on_event(FullNode::clone(&full_node))?;
436
437 ctx.spawn_ethstats(engine_events_for_ethstats).await?;
438
439 let handle = NodeHandle {
440 node_exit_future: NodeExitFuture::new(async { rx.await? }),
441 node: full_node,
442 };
443
444 Ok(handle)
445 }
446}
447
448impl<N, DB, T, CB, AO> LaunchNode<NodeBuilderWithComponents<T, CB, AO>> for EngineNodeLauncher
449where
450 T: FullNodeTypes<
451 Types = N,
452 DB = DB,
453 Provider = BlockchainProvider<NodeTypesWithDBAdapter<N, DB>>,
454 >,
455 N: Node<RethFullAdapter<DB, N>> + NodeTypesForProvider,
456 DB: Database + DatabaseMetrics + Clone + Unpin + 'static,
457 CB: NodeComponentsBuilder<T> + 'static,
458 AO: RethRpcAddOns<NodeAdapter<T, CB::Components>>
459 + EngineValidatorAddOn<NodeAdapter<T, CB::Components>>
460 + 'static,
461{
462 type Node = NodeHandle<NodeAdapter<T, CB::Components>, AO>;
463 type Future = Pin<Box<dyn Future<Output = eyre::Result<Self::Node>> + Send>>;
464
465 fn launch_node(self, target: NodeBuilderWithComponents<T, CB, AO>) -> Self::Future {
466 Box::pin(self.launch_node(target))
467 }
468}