Skip to main content

reth_node_builder/launch/
engine.rs

1//! Engine node related functionality.
2
3use 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/// The engine node launcher.
49#[derive(Debug)]
50pub struct EngineNodeLauncher {
51    /// The task executor for the node.
52    pub ctx: LaunchContext,
53
54    /// Temporary configuration for engine tree.
55    /// After engine is stabilized, this should be configured through node builder.
56    pub engine_tree_config: TreeConfig,
57}
58
59impl EngineNodeLauncher {
60    /// Create a new instance of the ethereum node launcher.
61    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        // Create the overlay manager that will be shared across the provider and engine.
96        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        // setup the launch context
102        let ctx = ctx
103            .with_configured_globals(engine_tree_config.reserved_cpu_cores())
104            // load the toml config
105            .with_loaded_toml_config(config)?
106            // add resolved peers
107            .with_resolved_peers()?
108            // attach the database
109            .attach(database.clone())
110            // ensure certain settings take effect
111            .with_adjusted_configs()
112            // Create the provider factory with the shared overlay manager
113            .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            // passing FullNodeTypes as type parameter here so that we can build
136            // later the components.
137            .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        // spawn exexs if any
143        let maybe_exex_manager_handle = ctx.launch_exex(installed_exex).await?;
144
145        // create pipeline
146        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        // We always assume that node is syncing after a restart
153        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        // The new engine writes directly to static files. This ensures that they're up to the tip.
180        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        // extract the jwt secret from the args if possible
198        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        // Build the engine validator with all required components
211        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        // Create the consensus engine stream with optional reorg
217        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            // Store messages _after_ skipping so that `replay-engine` command
238            // would replay only the messages that were observed by the engine
239            // during this run.
240            .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        // Create engine shutdown handle
297        let (engine_shutdown, shutdown_rx) = EngineShutdown::new();
298
299        // Run consensus engine to completion
300        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                // network_handle's sync state is already initialized at Syncing
321                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            // advance the chain and await payloads built locally to add into the engine api
330            // tree handler to prevent re-execution if that block is received as payload from
331            // the CL
332            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                                    // Once we're progressing via live sync, we can consider the node is not syncing anymore
358                                    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                        // Shutdown signal received.
397                        // Send Terminate so the engine OS thread can exit cleanly before we
398                        // drop the orchestrator.
399                        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        // Notify on node started
435        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}