arb_node/
engine.rs

1//! Custom engine orchestrator builder that exposes the tree sender.
2//!
3//! This is a thin wrapper around reth's `build_engine_orchestrator` pattern
4//! that also returns a clone of the tree sender channel, allowing our
5//! block producer to send `InsertExecutedBlock` and `ForkchoiceUpdated`
6//! directly to the engine tree for persistence.
7
8use 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
35/// The sender type for injecting blocks and FCU into the engine tree.
36pub type TreeSender<T, N> =
37    Sender<FromEngine<EngineApiRequest<T, N>, <N as NodePrimitives>::Block>>;
38
39/// Builds the engine orchestrator AND returns a clone of the tree sender.
40///
41/// This is identical to reth's `build_engine_orchestrator` but clones
42/// `to_tree_tx` before passing it to the request handler, allowing
43/// external code (our block producer) to send `InsertExecutedBlock`
44/// and `ForkchoiceUpdated` directly.
45#[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    // Clone the tree sender BEFORE it's consumed by the request handler.
103    // This allows our block producer to inject ExecutedBlocks directly.
104    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}