Skip to main content

reth_engine_tree/tree/
mod.rs

1use crate::{
2    backfill::{BackfillAction, BackfillSyncState},
3    chain::FromOrchestrator,
4    engine::{DownloadRequest, EngineApiEvent, EngineApiKind, EngineApiRequest, FromEngine},
5    persistence::PersistenceHandle,
6    tree::{error::InsertPayloadError, payload_validator::TreeCtx},
7};
8use alloy_consensus::BlockHeader;
9use alloy_eips::{eip1898::BlockWithParent, BlockNumHash, NumHash};
10use alloy_primitives::{map::B256Map, B256};
11use alloy_rpc_types_engine::{
12    ForkchoiceState, PayloadStatus, PayloadStatusEnum, PayloadValidationError,
13};
14use error::{
15    InsertBlockError, InsertBlockFatalError, InsertBlockProcessingError, InsertBlockValidationError,
16};
17use reth_chain_state::{
18    CanonicalInMemoryState, ExecutedBlock, ExecutionTimingStats, NewCanonicalChain,
19};
20use reth_consensus::{Consensus, FullConsensus};
21use reth_engine_primitives::{
22    BeaconEngineMessage, ConsensusEngineEvent, ExecutionPayload, ForkchoiceStateTracker,
23    NewPayloadTimings, OnForkChoiceUpdated, SlowBlockInfo,
24};
25use reth_errors::{ConsensusError, ProviderResult};
26use reth_evm::ConfigureEvm;
27use reth_network_p2p::full_block::SealedBlockWithAccessList;
28use reth_payload_builder::{BuildNewPayload, PayloadBuilderHandle, PayloadBuilderLease};
29use reth_payload_primitives::{BuiltPayload, NewPayloadError, PayloadAttributes, PayloadTypes};
30use reth_primitives_traits::{
31    FastInstant as Instant, NodePrimitives, RecoveredBlock, SealedBlock, SealedHeader,
32};
33use reth_provider::{
34    BalProvider, BlockExecutionOutput, BlockExecutionResult, BlockReader, ChangeSetReader,
35    DatabaseProviderFactory, ProviderError, PruneCheckpointReader, SaveBlocksInput,
36    StageCheckpointReader, StateProviderFactory, StateReader, StorageChangeSetReader,
37    StorageSettingsCache, TransactionVariant,
38};
39use reth_revm::database::StateProviderDatabase;
40use reth_stages_api::ControlFlow;
41use reth_storage_overlay::OverlayManager;
42use reth_tasks::{spawn_os_thread, utils::increase_thread_priority};
43use reth_trie::{HashedPostState, KeccakKeyHasher};
44use revm::interpreter::debug_unreachable;
45use state::TreeState;
46use std::{
47    fmt::Debug,
48    ops,
49    sync::{
50        atomic::{AtomicUsize, Ordering},
51        Arc,
52    },
53    time::Duration,
54};
55
56use crossbeam_channel::{Receiver, Sender};
57use tokio::sync::{
58    mpsc::{unbounded_channel, UnboundedReceiver, UnboundedSender},
59    oneshot,
60};
61use tracing::*;
62
63mod block_buffer;
64pub mod error;
65pub mod instrumented_state;
66mod invalid_headers;
67mod metrics;
68pub mod payload_processor;
69pub mod payload_validator;
70mod persistence_state;
71pub mod state_root_strategy;
72#[cfg(test)]
73mod tests;
74mod trie_updates;
75mod txpool_prewarm;
76pub mod types;
77
78use crate::{persistence::PersistenceResult, tree::error::AdvancePersistenceError};
79pub use block_buffer::BlockBuffer;
80pub use invalid_headers::InvalidHeaderCache;
81pub use metrics::EngineApiMetrics;
82pub use payload_processor::*;
83pub use payload_validator::{BasicEngineValidator, EngineValidator};
84pub use persistence_state::PersistenceState;
85pub use reth_engine_primitives::TreeConfig;
86pub use reth_execution_cache::{
87    precompile_cache, CachedStateCacheMetrics, CachedStateMetrics, CachedStateMetricsSource,
88    CachedStateProvider, ExecutionCache, PayloadExecutionCache, SavedCache,
89    TxPoolPrewarmCacheSnapshot,
90};
91pub use txpool_prewarm::{
92    Source as TxPoolPrewarmSource, Transaction as TxPoolPrewarmTransaction,
93    Transactions as TxPoolPrewarmTransactions,
94};
95pub use types::{ExecutionEnv, ValidationOutcome, ValidationOutput};
96
97pub mod state;
98
99/// The minimum number of blocks to retain in the changeset cache after eviction.
100///
101/// This ensures that recent changesets are kept in memory for potential reorgs,
102/// even when the finalized block is not set (e.g., on L2s like Optimism).
103const CHANGESET_CACHE_RETENTION_BLOCKS: u64 = 64;
104
105/// Tracks the state of the engine api internals.
106///
107/// This type is not shareable.
108#[derive(Debug)]
109pub struct EngineApiTreeState<N: NodePrimitives> {
110    /// Tracks the state of the blockchain tree.
111    tree_state: TreeState<N>,
112    /// Whether the next sparse trie task should attempt cache pruning during trie preservation.
113    pending_sparse_trie_prune: bool,
114    /// Tracks the forkchoice state updates received by the CL.
115    forkchoice_state_tracker: ForkchoiceStateTracker,
116    /// Buffer of detached blocks.
117    buffer: BlockBuffer<N::Block>,
118    /// Tracks the header of invalid payloads that were rejected by the engine because they're
119    /// invalid.
120    invalid_headers: InvalidHeaderCache,
121}
122
123impl<N: NodePrimitives> EngineApiTreeState<N> {
124    fn new(
125        block_buffer_limit: u32,
126        max_invalid_header_cache_length: u32,
127        invalid_header_hit_eviction_threshold: u8,
128        canonical_block: BlockNumHash,
129        engine_kind: EngineApiKind,
130        overlay_manager: OverlayManager<N>,
131    ) -> Self {
132        Self {
133            invalid_headers: InvalidHeaderCache::new(
134                max_invalid_header_cache_length,
135                invalid_header_hit_eviction_threshold,
136            ),
137            buffer: BlockBuffer::new(block_buffer_limit),
138            tree_state: TreeState::new(canonical_block, engine_kind, overlay_manager),
139            pending_sparse_trie_prune: false,
140            forkchoice_state_tracker: ForkchoiceStateTracker::default(),
141        }
142    }
143
144    /// Returns a reference to the tree state.
145    pub const fn tree_state(&self) -> &TreeState<N> {
146        &self.tree_state
147    }
148
149    /// Returns whether sparse trie pruning is pending.
150    pub const fn pending_sparse_trie_prune(&self) -> bool {
151        self.pending_sparse_trie_prune
152    }
153
154    /// Sets whether sparse trie pruning is pending for the next sparse trie task.
155    pub const fn set_pending_sparse_trie_prune(&mut self, pending: bool) {
156        self.pending_sparse_trie_prune = pending;
157    }
158
159    /// Takes a pending sparse trie prune request, if any, and snapshots the in-memory parent chain
160    /// ending at `parent_hash`.
161    ///
162    /// `None` means no prune request is pending. `Some(Vec::new())` means a prune was requested,
163    /// but no in-memory parent-chain blocks were found for the parent hash; the sparse trie task
164    /// should still prune nodes cached before the current block's epoch.
165    pub fn take_sparse_trie_prune_blocks(
166        &mut self,
167        parent_hash: B256,
168    ) -> Option<Vec<ExecutedBlock<N>>> {
169        if !self.pending_sparse_trie_prune {
170            return None
171        }
172
173        self.pending_sparse_trie_prune = false;
174        Some(
175            self.tree_state
176                .blocks_by_hash(parent_hash)
177                .map(|(_, blocks)| blocks)
178                .unwrap_or_default(),
179        )
180    }
181
182    /// Returns true if the block has been marked as invalid.
183    pub fn has_invalid_header(&mut self, hash: &B256) -> bool {
184        self.invalid_headers.get(hash).is_some()
185    }
186}
187
188/// The outcome of a tree operation.
189#[derive(Debug)]
190pub struct TreeOutcome<T> {
191    /// The outcome of the operation.
192    pub outcome: T,
193    /// An optional event to tell the caller to do something.
194    pub event: Option<TreeEvent>,
195    /// Whether the block was already seen, meaning no real execution happened during this
196    /// `newPayload` call.
197    pub already_seen: bool,
198}
199
200impl<T> TreeOutcome<T> {
201    /// Create new tree outcome.
202    pub const fn new(outcome: T) -> Self {
203        Self { outcome, event: None, already_seen: false }
204    }
205
206    /// Set event on the outcome.
207    pub fn with_event(mut self, event: TreeEvent) -> Self {
208        self.event = Some(event);
209        self
210    }
211
212    /// Set the `already_seen` flag on the outcome.
213    pub const fn with_already_seen(mut self, value: bool) -> Self {
214        self.already_seen = value;
215        self
216    }
217}
218
219/// Result of trying to insert a new payload in [`EngineApiTreeHandler`].
220#[derive(Debug)]
221pub struct TryInsertPayloadResult {
222    /// - `Valid`: Payload successfully validated and inserted
223    /// - `Syncing`: Parent missing, payload buffered for later
224    /// - Error status: Payload is invalid
225    pub status: PayloadStatus,
226    /// Whether the block was already seen
227    pub already_seen: bool,
228}
229
230impl TryInsertPayloadResult {
231    /// Convert the result into a [`TreeOutcome`].
232    #[inline]
233    pub fn into_outcome(self) -> TreeOutcome<PayloadStatus> {
234        TreeOutcome::new(self.status).with_already_seen(self.already_seen)
235    }
236}
237
238/// Events that are triggered by Tree Chain
239#[derive(Debug)]
240pub enum TreeEvent {
241    /// Tree action is needed.
242    TreeAction(TreeAction),
243    /// Backfill action is needed.
244    BackfillAction(BackfillAction),
245    /// Block download is needed.
246    Download(DownloadRequest),
247}
248
249impl TreeEvent {
250    /// Returns true if the event is a backfill action.
251    const fn is_backfill_action(&self) -> bool {
252        matches!(self, Self::BackfillAction(_))
253    }
254}
255
256/// The actions that can be performed on the tree.
257#[derive(Debug)]
258pub enum TreeAction {
259    /// Make target canonical.
260    MakeCanonical {
261        /// The sync target head hash
262        sync_target_head: B256,
263    },
264}
265
266/// The engine API tree handler implementation.
267///
268/// This type is responsible for processing engine API requests, maintaining the canonical state and
269/// emitting events.
270pub struct EngineApiTreeHandler<N, P, T, V, C>
271where
272    N: NodePrimitives,
273    T: PayloadTypes,
274    C: ConfigureEvm<Primitives = N> + 'static,
275{
276    provider: P,
277    consensus: Arc<dyn FullConsensus<N>>,
278    payload_validator: V,
279    /// Keeps track of internals such as executed and buffered blocks.
280    state: EngineApiTreeState<N>,
281    /// The half for sending messages to the engine.
282    ///
283    /// This is kept so that we can queue in messages to ourself that we can process later, for
284    /// example distributing workload across multiple messages that would otherwise take too long
285    /// to process. E.g. we might receive a range of downloaded blocks and we want to process
286    /// them one by one so that we can handle incoming engine API in between and don't become
287    /// unresponsive. This can happen during live sync transition where we're trying to close the
288    /// gap (up to 3 epochs of blocks in the worst case).
289    incoming_tx: Sender<FromEngine<EngineApiRequest<T, N>, N::Block>>,
290    /// Incoming engine API requests.
291    incoming: Receiver<FromEngine<EngineApiRequest<T, N>, N::Block>>,
292    /// Outgoing events that are emitted to the handler.
293    outgoing: UnboundedSender<EngineApiEvent<N>>,
294    /// Channels to the persistence layer.
295    persistence: PersistenceHandle<N>,
296    /// Tracks the state changes of the persistence task.
297    persistence_state: PersistenceState,
298    /// Flag indicating the state of the node's backfill synchronization process.
299    backfill_sync_state: BackfillSyncState,
300    /// Keeps track of the state of the canonical chain that isn't persisted yet.
301    /// This is intended to be accessed from external sources, such as rpc.
302    canonical_in_memory_state: CanonicalInMemoryState<N>,
303    /// Handle to the payload builder that will receive payload attributes for valid forkchoice
304    /// updates
305    payload_builder: PayloadBuilderHandle<T>,
306    /// Configuration settings.
307    config: TreeConfig,
308    /// Metrics for the engine api.
309    metrics: EngineApiMetrics,
310    /// The engine API variant of this handler
311    engine_kind: EngineApiKind,
312    /// The EVM configuration.
313    evm_config: C,
314    /// Timing statistics for executed blocks, keyed by block hash.
315    /// Stored here (not in `ExecutedBlock`) to avoid leaking observability concerns into the block
316    /// type. Entries are removed when blocks are persisted or invalidated.
317    execution_timing_stats: B256Map<Box<ExecutionTimingStats>>,
318    /// Tracks payload jobs that may still access in-memory overlay state.
319    payload_builds: PayloadBuildTracker,
320    /// Notifies the engine when the final active payload job finishes.
321    payload_build_finished: Receiver<()>,
322    /// Task runtime for spawning blocking work on named, reusable threads.
323    runtime: reth_tasks::Runtime,
324}
325
326impl<N, P: Debug, T: PayloadTypes + Debug, V: Debug, C> std::fmt::Debug
327    for EngineApiTreeHandler<N, P, T, V, C>
328where
329    N: NodePrimitives,
330    C: Debug + ConfigureEvm<Primitives = N>,
331{
332    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
333        f.debug_struct("EngineApiTreeHandler")
334            .field("provider", &self.provider)
335            .field("consensus", &self.consensus)
336            .field("payload_validator", &self.payload_validator)
337            .field("state", &self.state)
338            .field("incoming_tx", &self.incoming_tx)
339            .field("persistence", &self.persistence)
340            .field("persistence_state", &self.persistence_state)
341            .field("backfill_sync_state", &self.backfill_sync_state)
342            .field("canonical_in_memory_state", &self.canonical_in_memory_state)
343            .field("payload_builder", &self.payload_builder)
344            .field("config", &self.config)
345            .field("metrics", &self.metrics)
346            .field("engine_kind", &self.engine_kind)
347            .field("evm_config", &self.evm_config)
348            .field("execution_timing_stats", &self.execution_timing_stats.len())
349            .field("payload_builds_active", &self.payload_builds.is_active())
350            .field("runtime", &self.runtime)
351            .finish()
352    }
353}
354
355impl<N, P, T, V, C> EngineApiTreeHandler<N, P, T, V, C>
356where
357    N: NodePrimitives,
358    P: DatabaseProviderFactory
359        + BlockReader<Block = N::Block, Header = N::BlockHeader>
360        + StateProviderFactory
361        + StateReader<Receipt = N::Receipt>
362        + BalProvider
363        + Clone
364        + 'static,
365    P::Provider: BlockReader<Block = N::Block, Header = N::BlockHeader>
366        + PruneCheckpointReader
367        + StageCheckpointReader
368        + ChangeSetReader
369        + StorageChangeSetReader
370        + StorageSettingsCache
371        + 'static,
372    C: ConfigureEvm<Primitives = N> + 'static,
373    T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>,
374    V: EngineValidator<T> + WaitForCaches,
375{
376    /// Creates a new [`EngineApiTreeHandler`].
377    #[expect(clippy::too_many_arguments)]
378    pub fn new(
379        provider: P,
380        consensus: Arc<dyn FullConsensus<N>>,
381        payload_validator: V,
382        outgoing: UnboundedSender<EngineApiEvent<N>>,
383        state: EngineApiTreeState<N>,
384        canonical_in_memory_state: CanonicalInMemoryState<N>,
385        persistence: PersistenceHandle<N>,
386        persistence_state: PersistenceState,
387        payload_builder: PayloadBuilderHandle<T>,
388        config: TreeConfig,
389        engine_kind: EngineApiKind,
390        evm_config: C,
391        runtime: reth_tasks::Runtime,
392    ) -> Self {
393        let (incoming_tx, incoming) = crossbeam_channel::unbounded();
394
395        let (payload_builds, payload_build_finished) = PayloadBuildTracker::new();
396
397        Self {
398            provider,
399            consensus,
400            payload_validator,
401            incoming,
402            outgoing,
403            persistence,
404            persistence_state,
405            backfill_sync_state: BackfillSyncState::Idle,
406            state,
407            canonical_in_memory_state,
408            payload_builder,
409            config,
410            metrics: Default::default(),
411            incoming_tx,
412            engine_kind,
413            evm_config,
414            execution_timing_stats: B256Map::default(),
415            payload_builds,
416            payload_build_finished,
417            runtime,
418        }
419    }
420
421    /// Creates a new [`EngineApiTreeHandler`] instance and spawns it in its
422    /// own thread.
423    ///
424    /// Returns the sender through which incoming requests can be sent to the task and the receiver
425    /// end of a [`EngineApiEvent`] unbounded channel to receive events from the engine.
426    #[expect(clippy::complexity)]
427    pub fn spawn_new(
428        provider: P,
429        consensus: Arc<dyn FullConsensus<N>>,
430        payload_validator: V,
431        persistence: PersistenceHandle<N>,
432        payload_builder: PayloadBuilderHandle<T>,
433        canonical_in_memory_state: CanonicalInMemoryState<N>,
434        overlay_manager: OverlayManager<N>,
435        config: TreeConfig,
436        kind: EngineApiKind,
437        evm_config: C,
438        runtime: reth_tasks::Runtime,
439    ) -> (Sender<FromEngine<EngineApiRequest<T, N>, N::Block>>, UnboundedReceiver<EngineApiEvent<N>>)
440    {
441        let best_block_number = provider.best_block_number().unwrap_or(0);
442        let header = provider.sealed_header(best_block_number).ok().flatten().unwrap_or_default();
443
444        let persistence_state = PersistenceState {
445            last_persisted_block: BlockNumHash::new(best_block_number, header.hash()),
446            last_state_trie_persisted_block: BlockNumHash::new(best_block_number, header.hash()),
447            rx: None,
448        };
449
450        let (tx, outgoing) = unbounded_channel();
451        let state = EngineApiTreeState::new(
452            config.block_buffer_limit(),
453            config.max_invalid_header_cache_length(),
454            config.invalid_header_hit_eviction_threshold(),
455            header.num_hash(),
456            kind,
457            overlay_manager,
458        );
459
460        let task = Self::new(
461            provider,
462            consensus,
463            payload_validator,
464            tx,
465            state,
466            canonical_in_memory_state,
467            persistence,
468            persistence_state,
469            payload_builder,
470            config,
471            kind,
472            evm_config,
473            runtime,
474        );
475        let incoming = task.incoming_tx.clone();
476        spawn_os_thread("engine", || {
477            increase_thread_priority();
478            task.run()
479        });
480        (incoming, outgoing)
481    }
482
483    /// Returns a [`TreeOutcome`] indicating the forkchoice head is valid and canonical.
484    fn valid_outcome(state: ForkchoiceState) -> TreeOutcome<OnForkChoiceUpdated> {
485        TreeOutcome::new(OnForkChoiceUpdated::valid(PayloadStatus::new(
486            PayloadStatusEnum::Valid,
487            Some(state.head_block_hash),
488        )))
489    }
490
491    /// Returns a new [`Sender`] to send messages to this type.
492    pub fn sender(&self) -> Sender<FromEngine<EngineApiRequest<T, N>, N::Block>> {
493        self.incoming_tx.clone()
494    }
495
496    /// How many canonical blocks are retained in memory. A large count means persistence is
497    /// falling behind execution.
498    fn persistence_gap(&self) -> u64 {
499        self.canonical_in_memory_state.canonical_chain().count() as u64
500    }
501
502    /// How many blocks beyond the configured in-memory buffer are awaiting persistence.
503    fn persistence_backpressure_gap(&self) -> u64 {
504        self.persistence_gap().saturating_sub(self.config.memory_block_buffer_target())
505    }
506
507    /// Returns `true` when the main loop should stop draining the tree input channel.
508    ///
509    /// This is the case when persistence is already running and the number of blocks beyond the
510    /// configured in-memory buffer has reached the configured threshold.
511    fn should_backpressure(&self) -> bool {
512        self.persistence_state.in_progress() &&
513            self.persistence_backpressure_gap() >=
514                self.config.persistence_backpressure_threshold()
515    }
516
517    /// Run the engine API handler.
518    ///
519    /// This will block the current thread and process incoming messages.
520    pub fn run(mut self) {
521        loop {
522            // Each iteration has three phases:
523            //
524            // 1. Non-blocking poll for persistence completion. If the background flush already
525            //    landed, absorb the result now so the gap calculation below is fresh.
526            // 2. Decide how to wait for the next event. When the canonical-to-persisted gap beyond
527            //    the in-memory buffer reaches the backpressure threshold we only block on the
528            //    persistence receiver, leaving new engine requests sitting in the unbounded
529            //    upstream channel.
530            // 3. Handle the event (engine message or persistence completion) and kick off a new
531            //    persistence cycle if the threshold is met again.
532            //
533            // The net effect: when the unbuffered persistence gap reaches the threshold, we stop
534            // processing incoming messages and let them queue in the channel. This is only a soft
535            // form of backpressure: it delays replies and, more importantly, prevents executing
536            // further blocks that would pile up in the persistence queue - where each block
537            // carries heavier state (eg. trie updates) than the raw payload sitting in the engine
538            // channel.
539            //
540            // Standard Ethereum CLs won't truly back off - the engine API has no
541            // backpressure semantics, and CLs typically timeout after ≈8s and resend - so
542            // this cannot prevent the incoming channel from growing under sustained load.
543            // But it shifts the bottleneck to the lighter-weight incoming queue rather than
544            // the costlier persistence pipeline. Other clients that respect reply latency
545            // can treat the delayed responses as a signal to chill out.
546            match self.try_poll_persistence() {
547                Ok(true) => {
548                    if let Err(err) = self.advance_persistence() {
549                        error!(target: "engine::tree", %err, "Advancing persistence failed");
550                        return
551                    }
552                    continue;
553                }
554                Ok(false) => {}
555                Err(err) => {
556                    error!(target: "engine::tree", %err, "Polling persistence failed");
557                    return
558                }
559            }
560
561            let event = if self.should_backpressure() {
562                self.metrics.engine.backpressure_active.set(1.0);
563                let stall_start = Instant::now();
564                let event = self.wait_for_persistence_event();
565                self.metrics.engine.backpressure_stall_duration.record(stall_start.elapsed());
566                event
567            } else {
568                self.metrics.engine.backpressure_active.set(0.0);
569                self.wait_for_event()
570            };
571
572            match event {
573                LoopEvent::EngineMessage(msg) => {
574                    debug!(target: "engine::tree", %msg, "received new engine message");
575                    match self.on_engine_message(msg) {
576                        Ok(ops::ControlFlow::Break(())) => return,
577                        Ok(ops::ControlFlow::Continue(())) => {}
578                        Err(fatal) => {
579                            error!(target: "engine::tree", %fatal, "insert block fatal error");
580                            return
581                        }
582                    }
583                }
584                LoopEvent::PersistenceComplete { result, start_time } => {
585                    if let Err(err) = self.on_persistence_complete(result, start_time) {
586                        error!(target: "engine::tree", %err, "Persistence complete handling failed");
587                        return
588                    }
589                }
590                LoopEvent::PayloadBuildFinished => {}
591                LoopEvent::Disconnected => {
592                    error!(target: "engine::tree", "Channel disconnected");
593                    return
594                }
595            }
596
597            // Always check if we need to trigger new persistence after any event:
598            // - After engine messages: new blocks may have been inserted that exceed the
599            //   persistence threshold
600            // - After persistence completion: we can now persist more blocks if needed
601            if let Err(err) = self.advance_persistence() {
602                error!(target: "engine::tree", %err, "Advancing persistence failed");
603                return
604            }
605        }
606    }
607
608    /// Blocks until the in-flight persistence task completes, used when we are under
609    /// backpressure.
610    ///
611    /// Unlike `wait_for_event`, this deliberately does not read from the tree input channel. Any
612    /// requests sent to the tree remain queued upstream until persistence catches up.
613    fn wait_for_persistence_event(&mut self) -> LoopEvent<T, N> {
614        let maybe_persistence = self.persistence_state.rx.take();
615
616        if let Some((persistence_rx, start_time, _action)) = maybe_persistence {
617            match persistence_rx.recv() {
618                Ok(result) => LoopEvent::PersistenceComplete { result, start_time },
619                Err(_) => LoopEvent::Disconnected,
620            }
621        } else {
622            self.wait_for_event()
623        }
624    }
625
626    /// Blocks until the next event is ready.
627    ///
628    /// Uses biased selection to prioritize persistence completion so in-memory state is updated and
629    /// further writes are unblocked.
630    fn wait_for_event(&mut self) -> LoopEvent<T, N> {
631        // Take ownership of persistence rx if present
632        let maybe_persistence = self.persistence_state.rx.take();
633
634        if let Some((persistence_rx, start_time, action)) = maybe_persistence {
635            // Biased select prioritizes persistence completion to update in memory state and
636            // unblock further writes
637            crossbeam_channel::select_biased! {
638                recv(persistence_rx) -> result => {
639                    // Don't put it back - consumed (oneshot-like behavior)
640                    match result {
641                        Ok(result) => LoopEvent::PersistenceComplete {
642                            result,
643                            start_time,
644                        },
645                        Err(_) => LoopEvent::Disconnected,
646                    }
647                },
648                recv(self.payload_build_finished) -> result => {
649                    // Put the persistence rx back - we didn't consume it.
650                    self.persistence_state.rx = Some((persistence_rx, start_time, action));
651                    match result {
652                        Ok(()) => LoopEvent::PayloadBuildFinished,
653                        Err(_) => LoopEvent::Disconnected,
654                    }
655                },
656                recv(self.incoming) -> msg => {
657                    // Put the persistence rx back - we didn't consume it
658                    self.persistence_state.rx = Some((persistence_rx, start_time, action));
659                    match msg {
660                        Ok(m) => LoopEvent::EngineMessage(m),
661                        Err(_) => LoopEvent::Disconnected,
662                    }
663                },
664            }
665        } else {
666            // No persistence in progress - wait on an incoming message or payload job completion.
667            crossbeam_channel::select_biased! {
668                recv(self.payload_build_finished) -> result => match result {
669                    Ok(()) => LoopEvent::PayloadBuildFinished,
670                    Err(_) => LoopEvent::Disconnected,
671                },
672                recv(self.incoming) -> msg => match msg {
673                    Ok(m) => LoopEvent::EngineMessage(m),
674                    Err(_) => LoopEvent::Disconnected,
675                },
676            }
677        }
678    }
679
680    /// Invoked when previously requested blocks were downloaded.
681    ///
682    /// If the block count exceeds the configured batch size we're allowed to execute at once, this
683    /// will execute the first batch and send the remaining blocks back through the channel so that
684    /// block request processing isn't blocked for a long time.
685    fn on_downloaded(
686        &mut self,
687        mut blocks: Vec<SealedBlockWithAccessList<N::Block>>,
688    ) -> Result<Option<TreeEvent>, InsertBlockFatalError> {
689        if blocks.is_empty() {
690            // nothing to execute
691            return Ok(None)
692        }
693
694        trace!(target: "engine::tree", block_count = %blocks.len(), "received downloaded blocks");
695        let batch = self.config.max_execute_block_batch_size().min(blocks.len());
696        for block in blocks.drain(..batch) {
697            if let Some(event) = self.on_downloaded_block(block)? {
698                let needs_backfill = event.is_backfill_action();
699                self.on_tree_event(event)?;
700                if needs_backfill {
701                    // can exit early if backfill is needed
702                    return Ok(None)
703                }
704            }
705        }
706
707        // if we still have blocks to execute, send them as a followup request
708        if !blocks.is_empty() {
709            let _ = self.incoming_tx.send(FromEngine::DownloadedBlocks(blocks));
710        }
711
712        Ok(None)
713    }
714
715    /// When the Consensus layer receives a new block via the consensus gossip protocol,
716    /// the transactions in the block are sent to the execution layer in the form of a
717    /// [`PayloadTypes::ExecutionData`], for example
718    /// [`ExecutionData`](reth_payload_primitives::PayloadTypes::ExecutionData). The
719    /// Execution layer executes the transactions and validates the state in the block header,
720    /// then passes validation data back to Consensus layer, that adds the block to the head of
721    /// its own blockchain and attests to it. The block is then broadcast over the consensus p2p
722    /// network in the form of a "Beacon block".
723    ///
724    /// These responses should adhere to the [Engine API Spec for
725    /// `engine_newPayload`](https://github.com/ethereum/execution-apis/blob/main/src/engine/paris.md#specification).
726    ///
727    /// This returns a [`PayloadStatus`] that represents the outcome of a processed new payload and
728    /// returns an error if an internal error occurred.
729    #[instrument(
730        level = "debug",
731        target = "engine::tree",
732        skip_all,
733        fields(block_hash = %payload.block_hash(), block_num = %payload.block_number()),
734    )]
735    fn on_new_payload(
736        &mut self,
737        payload: T::ExecutionData,
738    ) -> Result<TreeOutcome<PayloadStatus>, InsertBlockProcessingError> {
739        let _thread_resource_usage =
740            self.metrics.engine.new_payload.measure_thread_resource_usage();
741        trace!(target: "engine::tree", "invoked new payload");
742
743        // start timing for the new payload process
744        let start = Instant::now();
745
746        // Ensures that the given payload does not violate any consensus rules that concern the
747        // block's layout, like:
748        //    - missing or invalid base fee
749        //    - invalid extra data
750        //    - invalid transactions
751        //    - incorrect hash
752        //    - the versioned hashes passed with the payload do not exactly match transaction
753        //      versioned hashes
754        //    - the block does not contain blob transactions if it is pre-cancun
755        //
756        // This validates the following engine API rule:
757        //
758        // 3. Given the expected array of blob versioned hashes client software **MUST** run its
759        //    validation by taking the following steps:
760        //
761        //   1. Obtain the actual array by concatenating blob versioned hashes lists
762        //      (`tx.blob_versioned_hashes`) of each [blob
763        //      transaction](https://eips.ethereum.org/EIPS/eip-4844#new-transaction-type) included
764        //      in the payload, respecting the order of inclusion. If the payload has no blob
765        //      transactions the expected array **MUST** be `[]`.
766        //
767        //   2. Return `{status: INVALID, latestValidHash: null, validationError: errorMessage |
768        //      null}` if the expected and the actual arrays don't match.
769        //
770        // This validation **MUST** be instantly run in all cases even during active sync process.
771
772        let num_hash = payload.num_hash();
773        let engine_event = ConsensusEngineEvent::BlockReceived(num_hash);
774        self.emit_event(EngineApiEvent::BeaconConsensus(engine_event));
775
776        let block_hash = num_hash.hash;
777
778        // Check for invalid ancestors
779        if let Some(invalid) = self.find_invalid_ancestor(&payload) {
780            let status = self.handle_invalid_ancestor_payload(payload, invalid)?;
781            return Ok(TreeOutcome::new(status));
782        }
783
784        // record pre-execution phase duration
785        self.metrics.block_validation.record_payload_validation(start.elapsed().as_secs_f64());
786
787        let mut outcome = if self.backfill_sync_state.is_idle() {
788            self.try_insert_payload(payload)?.into_outcome()
789        } else {
790            TreeOutcome::new(self.try_buffer_payload(payload)?)
791        };
792
793        // if the block is valid and it is the current sync target head, make it canonical
794        if outcome.outcome.is_valid() && self.is_sync_target_head(block_hash) {
795            // Only create the canonical event if this block isn't already the canonical head
796            if self.state.tree_state.canonical_block_hash() != block_hash {
797                outcome = outcome.with_event(TreeEvent::TreeAction(TreeAction::MakeCanonical {
798                    sync_target_head: block_hash,
799                }));
800            }
801        }
802
803        // record total newPayload duration
804        self.metrics.block_validation.total_duration.record(start.elapsed().as_secs_f64());
805
806        Ok(outcome)
807    }
808
809    /// Processes a payload during normal sync operation.
810    #[instrument(level = "debug", target = "engine::tree", skip_all)]
811    fn try_insert_payload(
812        &mut self,
813        payload: T::ExecutionData,
814    ) -> Result<TryInsertPayloadResult, InsertBlockProcessingError> {
815        let block_hash = payload.block_hash();
816        let num_hash = payload.num_hash();
817        let parent_hash = payload.parent_hash();
818        let mut latest_valid_hash = None;
819
820        match self.insert_payload(payload) {
821            Ok(status) => {
822                let (status, already_seen) = match status {
823                    InsertPayloadOk::Inserted(BlockStatus::Valid) => {
824                        latest_valid_hash = Some(block_hash);
825                        self.try_connect_buffered_blocks(num_hash)?;
826                        (PayloadStatusEnum::Valid, false)
827                    }
828                    InsertPayloadOk::AlreadySeen(BlockStatus::Valid) => {
829                        latest_valid_hash = Some(block_hash);
830                        (PayloadStatusEnum::Valid, true)
831                    }
832                    InsertPayloadOk::Inserted(BlockStatus::Disconnected { .. }) => {
833                        (PayloadStatusEnum::Syncing, false)
834                    }
835                    InsertPayloadOk::AlreadySeen(BlockStatus::Disconnected { .. }) => {
836                        // not known to be invalid, but we don't know anything else
837                        (PayloadStatusEnum::Syncing, true)
838                    }
839                };
840
841                Ok(TryInsertPayloadResult {
842                    status: PayloadStatus::new(status, latest_valid_hash),
843                    already_seen,
844                })
845            }
846            Err(error) => {
847                let status = match error {
848                    InsertPayloadError::Block(error) => self.on_insert_block_error(error)?,
849                    InsertPayloadError::Payload(error) => self
850                        .on_new_payload_error(error, num_hash, parent_hash)
851                        .map_err(InsertBlockFatalError::from)?,
852                };
853
854                Ok(TryInsertPayloadResult { status, already_seen: false })
855            }
856        }
857    }
858
859    /// Stores a payload for later processing during backfill sync.
860    ///
861    /// During backfill, the node lacks the state needed to validate payloads,
862    /// so they are buffered (stored in memory) until their parent blocks are synced.
863    ///
864    /// Returns:
865    /// - `Syncing`: Payload successfully buffered
866    /// - Error status: Payload is malformed or invalid
867    fn try_buffer_payload(
868        &mut self,
869        payload: T::ExecutionData,
870    ) -> Result<PayloadStatus, InsertBlockProcessingError> {
871        let parent_hash = payload.parent_hash();
872        let num_hash = payload.num_hash();
873
874        match self.payload_validator.convert_payload_to_block(payload) {
875            // if the block is well-formed, buffer it for later
876            Ok(block) => {
877                if let Err(error) = self.buffer_block(block) {
878                    self.on_insert_block_error(error)
879                } else {
880                    Ok(PayloadStatus::from_status(PayloadStatusEnum::Syncing))
881                }
882            }
883            Err(error) => Ok(self
884                .on_new_payload_error(error, num_hash, parent_hash)
885                .map_err(InsertBlockFatalError::from)?),
886        }
887    }
888
889    /// Returns the new chain for the given head.
890    ///
891    /// This also handles reorgs.
892    ///
893    /// Note: This does not update the tracked state and instead returns the new chain based on the
894    /// given head.
895    fn on_new_head(&self, new_head: B256) -> ProviderResult<Option<NewCanonicalChain<N>>> {
896        // get the executed new head block
897        let Some(new_head_block) = self.state.tree_state.blocks_by_hash.get(&new_head) else {
898            debug!(target: "engine::tree", new_head=?new_head, "New head block not found in inmemory tree state");
899            self.metrics.engine.executed_new_block_cache_miss.increment(1);
900            return Ok(None)
901        };
902
903        let new_head_number = new_head_block.recovered_block().number();
904        let mut current_canonical_number = self.state.tree_state.current_canonical_head.number;
905
906        let mut new_chain = vec![new_head_block.clone()];
907        let mut current_hash = new_head_block.recovered_block().parent_hash();
908        let mut current_number = new_head_number - 1;
909
910        // Walk back the new chain until we reach a block we know about
911        //
912        // This is only done for in-memory blocks, because we should not have persisted any blocks
913        // that are _above_ the current canonical head.
914        while current_number > current_canonical_number {
915            if let Some(block) = self.state.tree_state.executed_block_by_hash(current_hash).cloned()
916            {
917                current_hash = block.recovered_block().parent_hash();
918                current_number -= 1;
919                new_chain.push(block);
920            } else {
921                warn!(target: "engine::tree", current_hash=?current_hash, "Sidechain block not found in TreeState");
922                // This should never happen as we're walking back a chain that should connect to
923                // the canonical chain
924                return Ok(None)
925            }
926        }
927
928        // If we have reached the current canonical head by walking back from the target, then we
929        // know this represents an extension of the canonical chain.
930        if current_hash == self.state.tree_state.current_canonical_head.hash {
931            new_chain.reverse();
932
933            // Simple extension of the current chain
934            return Ok(Some(NewCanonicalChain::Commit { new: new_chain }))
935        }
936
937        // We have a reorg. Walk back both chains to find the fork point.
938        let mut old_chain = Vec::new();
939        let mut old_hash = self.state.tree_state.current_canonical_head.hash;
940
941        // If the canonical chain is ahead of the new chain,
942        // gather all blocks until new head number.
943        while current_canonical_number > current_number {
944            let block = self.canonical_block_by_hash(old_hash)?;
945            old_hash = block.recovered_block().parent_hash();
946            old_chain.push(block);
947            current_canonical_number -= 1;
948        }
949
950        // Both new and old chain pointers are now at the same height.
951        debug_assert_eq!(current_number, current_canonical_number);
952
953        // Walk both chains from specified hashes at same height until
954        // a common ancestor (fork block) is reached.
955        while old_hash != current_hash {
956            let block = self.canonical_block_by_hash(old_hash)?;
957            old_hash = block.recovered_block().parent_hash();
958            old_chain.push(block);
959
960            if let Some(block) = self.state.tree_state.executed_block_by_hash(current_hash).cloned()
961            {
962                current_hash = block.recovered_block().parent_hash();
963                new_chain.push(block);
964            } else {
965                // This shouldn't happen as we've already walked this path
966                warn!(target: "engine::tree", invalid_hash=?current_hash, "New chain block not found in TreeState");
967                return Ok(None)
968            }
969        }
970        new_chain.reverse();
971        old_chain.reverse();
972
973        Ok(Some(NewCanonicalChain::Reorg { new: new_chain, old: old_chain }))
974    }
975
976    /// Updates the latest block state to the specified canonical ancestor.
977    ///
978    /// This method ensures that the latest block tracks the given canonical header by resetting
979    ///
980    /// # Arguments
981    /// * `canonical_header` - The canonical header to set as the new head
982    ///
983    /// # Returns
984    /// * `ProviderResult<()>` - Ok(()) on success, error if state update fails
985    ///
986    /// Caution: This unwinds the canonical chain
987    fn update_latest_block_to_canonical_ancestor(
988        &mut self,
989        canonical_header: &SealedHeader<N::BlockHeader>,
990    ) -> ProviderResult<()> {
991        debug!(target: "engine::tree", head = ?canonical_header.num_hash(), "Update latest block to canonical ancestor");
992        let current_head_number = self.state.tree_state.canonical_block_number();
993        let new_head_number = canonical_header.number();
994        let new_head_hash = canonical_header.hash();
995
996        // Update tree state with the new canonical head
997        self.state.tree_state.set_canonical_head(canonical_header.num_hash());
998
999        // Handle the state update based on whether this is an unwind scenario
1000        if new_head_number < current_head_number {
1001            debug!(
1002                target: "engine::tree",
1003                current_head = current_head_number,
1004                new_head = new_head_number,
1005                new_head_hash = ?new_head_hash,
1006                "FCU unwind detected: reverting to canonical ancestor"
1007            );
1008
1009            self.handle_canonical_chain_unwind(current_head_number, canonical_header)
1010        } else {
1011            debug!(
1012                target: "engine::tree",
1013                previous_head = current_head_number,
1014                new_head = new_head_number,
1015                new_head_hash = ?new_head_hash,
1016                "Advancing latest block to canonical ancestor"
1017            );
1018            self.handle_chain_advance_or_same_height(canonical_header)
1019        }
1020    }
1021
1022    /// Handles chain unwind scenarios by collecting blocks to remove and performing an unwind back
1023    /// to the canonical header
1024    fn handle_canonical_chain_unwind(
1025        &self,
1026        current_head_number: u64,
1027        canonical_header: &SealedHeader<N::BlockHeader>,
1028    ) -> ProviderResult<()> {
1029        let new_head_number = canonical_header.number();
1030        debug!(
1031            target: "engine::tree",
1032            from = current_head_number,
1033            to = new_head_number,
1034            "Handling unwind: collecting blocks to remove from in-memory state"
1035        );
1036
1037        // Collect blocks that need to be removed from memory
1038        let old_blocks =
1039            self.collect_blocks_for_canonical_unwind(new_head_number, current_head_number);
1040
1041        // Load and apply the canonical ancestor block
1042        self.apply_canonical_ancestor_via_reorg(canonical_header, old_blocks)
1043    }
1044
1045    /// Collects blocks from memory that need to be removed during an unwind to a canonical block.
1046    fn collect_blocks_for_canonical_unwind(
1047        &self,
1048        new_head_number: u64,
1049        current_head_number: u64,
1050    ) -> Vec<ExecutedBlock<N>> {
1051        let mut old_blocks =
1052            Vec::with_capacity((current_head_number.saturating_sub(new_head_number)) as usize);
1053
1054        for block_num in (new_head_number + 1)..=current_head_number {
1055            if let Some(block_state) = self.canonical_in_memory_state.state_by_number(block_num) {
1056                let executed_block = block_state.block_ref().clone();
1057                old_blocks.push(executed_block);
1058                debug!(
1059                    target: "engine::tree",
1060                    block_number = block_num,
1061                    "Collected block for removal from in-memory state"
1062                );
1063            }
1064        }
1065
1066        if old_blocks.is_empty() {
1067            debug!(
1068                target: "engine::tree",
1069                "No blocks found in memory to remove, will clear and reset state"
1070            );
1071        }
1072
1073        old_blocks
1074    }
1075
1076    /// Applies the canonical ancestor block via a reorg operation.
1077    fn apply_canonical_ancestor_via_reorg(
1078        &self,
1079        canonical_header: &SealedHeader<N::BlockHeader>,
1080        old_blocks: Vec<ExecutedBlock<N>>,
1081    ) -> ProviderResult<()> {
1082        let new_head_hash = canonical_header.hash();
1083        let new_head_number = canonical_header.number();
1084
1085        // Load the canonical ancestor's block
1086        let executed_block = self.canonical_block_by_hash(new_head_hash)?;
1087        // Perform the reorg to properly handle the unwind
1088        self.canonical_in_memory_state
1089            .update_chain(NewCanonicalChain::Reorg { new: vec![executed_block], old: old_blocks });
1090
1091        // CRITICAL: Update the canonical head after the reorg
1092        // This ensures get_canonical_head() returns the correct block
1093        self.canonical_in_memory_state.set_canonical_head(canonical_header.clone());
1094
1095        debug!(
1096            target: "engine::tree",
1097            block_number = new_head_number,
1098            block_hash = ?new_head_hash,
1099            "Successfully loaded canonical ancestor into memory via reorg"
1100        );
1101
1102        Ok(())
1103    }
1104
1105    /// Handles chain advance or same height scenarios.
1106    fn handle_chain_advance_or_same_height(
1107        &self,
1108        canonical_header: &SealedHeader<N::BlockHeader>,
1109    ) -> ProviderResult<()> {
1110        // Load the block into memory if it's not already present
1111        self.ensure_block_in_memory(canonical_header.number(), canonical_header.hash())?;
1112
1113        // Update the canonical head header
1114        self.canonical_in_memory_state.set_canonical_head(canonical_header.clone());
1115
1116        Ok(())
1117    }
1118
1119    /// Ensures a block is loaded into memory if not already present.
1120    fn ensure_block_in_memory(&self, block_number: u64, block_hash: B256) -> ProviderResult<()> {
1121        // Check if block is already in memory
1122        if self.canonical_in_memory_state.state_by_number(block_number).is_some() {
1123            return Ok(());
1124        }
1125
1126        // Load the block from storage
1127        let executed_block = self.canonical_block_by_hash(block_hash)?;
1128        self.canonical_in_memory_state
1129            .update_chain(NewCanonicalChain::Commit { new: vec![executed_block] });
1130
1131        debug!(
1132            target: "engine::tree",
1133            block_number,
1134            block_hash = ?block_hash,
1135            "Added canonical block to in-memory state"
1136        );
1137
1138        Ok(())
1139    }
1140
1141    /// Invoked when we receive a new forkchoice update message. Calls into the blockchain tree
1142    /// to resolve chain forks and ensure that the Execution Layer is working with the latest valid
1143    /// chain.
1144    ///
1145    /// These responses should adhere to the [Engine API Spec for
1146    /// `engine_forkchoiceUpdated`](https://github.com/ethereum/execution-apis/blob/main/src/engine/paris.md#specification-1).
1147    ///
1148    /// Returns an error if an internal error occurred like a database error.
1149    #[instrument(level = "debug", target = "engine::tree", skip_all, fields(head = % state.head_block_hash, safe = % state.safe_block_hash,finalized = % state.finalized_block_hash))]
1150    fn on_forkchoice_updated(
1151        &mut self,
1152        state: ForkchoiceState,
1153        attrs: Option<T::PayloadAttributes>,
1154    ) -> ProviderResult<TreeOutcome<OnForkChoiceUpdated>> {
1155        trace!(target: "engine::tree", ?attrs, "invoked forkchoice update");
1156
1157        // Record metrics
1158        self.record_forkchoice_metrics();
1159
1160        // Pre-validation of forkchoice state
1161        if let Some(early_result) = self.validate_forkchoice_state(state)? {
1162            return Ok(TreeOutcome::new(early_result));
1163        }
1164
1165        // Return early if we are on the correct fork
1166        if let Some(result) = self.handle_canonical_head(state, &attrs)? {
1167            return Ok(result);
1168        }
1169
1170        // Attempt to apply a chain update when the head differs from our canonical chain.
1171        // This handles reorgs and chain extensions by making the specified head canonical.
1172        if let Some(result) = self.apply_chain_update(state, &attrs)? {
1173            return Ok(result);
1174        }
1175
1176        // Fallback that ensures to catch up to the network's state.
1177        self.handle_missing_block(state)
1178    }
1179
1180    /// Records metrics for forkchoice updated calls
1181    fn record_forkchoice_metrics(&self) {
1182        self.canonical_in_memory_state.on_forkchoice_update_received();
1183    }
1184
1185    /// Pre-validates the forkchoice state and returns early if validation fails.
1186    ///
1187    /// Returns `Some(OnForkChoiceUpdated)` if validation fails and an early response should be
1188    /// returned. Returns `None` if validation passes and processing should continue.
1189    fn validate_forkchoice_state(
1190        &mut self,
1191        state: ForkchoiceState,
1192    ) -> ProviderResult<Option<OnForkChoiceUpdated>> {
1193        if state.head_block_hash.is_zero() {
1194            return Ok(Some(OnForkChoiceUpdated::invalid_state()));
1195        }
1196
1197        // Check if the new head hash is connected to any ancestor that we previously marked as
1198        // invalid
1199        let lowest_buffered_ancestor_fcu = self.lowest_buffered_ancestor_or(state.head_block_hash);
1200        if let Some(status) = self.check_invalid_ancestor(lowest_buffered_ancestor_fcu)? {
1201            return Ok(Some(OnForkChoiceUpdated::with_invalid(status)));
1202        }
1203
1204        if !self.backfill_sync_state.is_idle() {
1205            // We can only process new forkchoice updates if the pipeline is idle, since it requires
1206            // exclusive access to the database
1207            trace!(target: "engine::tree", "Pipeline is syncing, skipping forkchoice update");
1208            return Ok(Some(OnForkChoiceUpdated::syncing()));
1209        }
1210
1211        Ok(None)
1212    }
1213
1214    /// Handles the case where the forkchoice head is already canonical.
1215    ///
1216    /// Returns `Some(TreeOutcome<OnForkChoiceUpdated>)` if the head is already canonical and
1217    /// processing is complete. Returns `None` if the head is not canonical and processing
1218    /// should continue.
1219    fn handle_canonical_head(
1220        &mut self,
1221        state: ForkchoiceState,
1222        attrs: &Option<T::PayloadAttributes>, // Changed to reference
1223    ) -> ProviderResult<Option<TreeOutcome<OnForkChoiceUpdated>>> {
1224        // Process the forkchoice update by trying to make the head block canonical
1225        //
1226        // We can only process this forkchoice update if:
1227        // - we have the `head` block
1228        // - the head block is part of a chain that is connected to the canonical chain. This
1229        //   includes reorgs.
1230        //
1231        // Performing a FCU involves:
1232        // - marking the FCU's head block as canonical
1233        // - updating in memory state to reflect the new canonical chain
1234        // - updating canonical state trackers
1235        // - emitting a canonicalization event for the new chain (including reorg)
1236        // - if we have payload attributes, delegate them to the payload service
1237
1238        if self.state.tree_state.canonical_block_hash() != state.head_block_hash {
1239            return Ok(None);
1240        }
1241
1242        trace!(target: "engine::tree", "fcu head hash is already canonical");
1243
1244        if !self.is_consistent_forkchoice_state(state, None)? {
1245            return Ok(Some(TreeOutcome::new(OnForkChoiceUpdated::invalid_state())));
1246        }
1247
1248        // Update the safe and finalized blocks and ensure their values are valid
1249        if let Err(outcome) = self.ensure_consistent_forkchoice_state(state) {
1250            // safe or finalized hashes are invalid
1251            return Ok(Some(TreeOutcome::new(outcome)));
1252        }
1253
1254        self.payload_validator.on_canonical_head_changed(state.head_block_hash, &self.state);
1255
1256        // Process payload attributes if the head is already canonical
1257        if let Some(attr) = attrs {
1258            let tip = self
1259                .sealed_header_by_hash(self.state.tree_state.canonical_block_hash())?
1260                .ok_or_else(|| {
1261                    // If we can't find the canonical block, then something is wrong and we need
1262                    // to return an error
1263                    ProviderError::HeaderNotFound(state.head_block_hash.into())
1264                })?;
1265            // Clone only when we actually need to process the attributes
1266            let updated = self.process_payload_attributes(attr.clone(), &tip, state);
1267            return Ok(Some(TreeOutcome::new(updated)));
1268        }
1269
1270        // The head block is already canonical
1271        Ok(Some(Self::valid_outcome(state)))
1272    }
1273
1274    /// Applies chain update for the new head block and processes payload attributes.
1275    ///
1276    /// This method handles the case where the forkchoice head differs from our current canonical
1277    /// head. It attempts to make the specified head block canonical by:
1278    /// - Checking if the head is already part of the canonical chain
1279    /// - Applying chain reorganizations (reorgs) if necessary
1280    /// - Processing payload attributes if provided
1281    /// - Returning the appropriate forkchoice update response
1282    ///
1283    /// Returns `Some(TreeOutcome<OnForkChoiceUpdated>)` if a chain update was successfully applied.
1284    /// Returns `None` if no chain update was needed or possible.
1285    fn apply_chain_update(
1286        &mut self,
1287        state: ForkchoiceState,
1288        attrs: &Option<T::PayloadAttributes>,
1289    ) -> ProviderResult<Option<TreeOutcome<OnForkChoiceUpdated>>> {
1290        // Check if the head is already part of the canonical chain
1291        if let Ok(Some(canonical_header)) = self.find_canonical_header(state.head_block_hash) {
1292            debug!(target: "engine::tree", head = canonical_header.number(), "fcu head block is already canonical");
1293
1294            // For OpStack, or if explicitly configured, the proposers are allowed to reorg their
1295            // own chain at will, so we need to always trigger a new payload job if requested.
1296            let always_trigger_payload_job = self.engine_kind.is_opstack() ||
1297                self.config.always_process_payload_attributes_on_canonical_head();
1298
1299            // A canonical ancestor below the latest known finalized block can never become the
1300            // head again, because this would reorg out the finalized block. Such a forkchoice
1301            // update exceeds the supported reorg depth and is rejected regardless of the payload
1302            // attributes:
1303            // <https://github.com/ethereum/execution-apis/blob/bf20b4083284e677db19e7f3871bd669b88354a6/src/engine/paris.md?plain=1#L221>
1304            //
1305            // The stored finalized block is used because a forkchoice update MAY carry a zero
1306            // finalized hash without clearing previously established finality.
1307            if !always_trigger_payload_job &&
1308                self.canonical_in_memory_state
1309                    .get_finalized_num_hash()
1310                    .is_some_and(|finalized| canonical_header.number() < finalized.number)
1311            {
1312                debug!(target: "engine::tree", head = canonical_header.number(), "rejecting canonical ancestor fcu below the finalized block");
1313                return Ok(Some(TreeOutcome::new(OnForkChoiceUpdated::too_deep_reorg())));
1314            }
1315
1316            if !self.is_consistent_forkchoice_state(state, None)? {
1317                return Ok(Some(TreeOutcome::new(OnForkChoiceUpdated::invalid_state())));
1318            }
1319
1320            // We need to effectively unwind the _canonical_ chain to the FCU's head, which is
1321            // part of the canonical chain. We need to update the latest block state to reflect
1322            // the canonical ancestor. This ensures that state providers and the transaction
1323            // pool operate with the correct chain state after forkchoice update processing, and
1324            // new payloads built on the reorg'd head will be added to the tree immediately.
1325            if always_trigger_payload_job && self.config.unwind_canonical_header() {
1326                self.update_latest_block_to_canonical_ancestor(&canonical_header)?;
1327            }
1328
1329            // A canonical ancestor at or above the latest known finalized block can become the
1330            // parent of the next block, e.g. when the CL wants to reorg out the current head.
1331            // The canonical chain remains untouched here; the block built on the ancestor
1332            // triggers the actual reorg once it is inserted via newPayload and FCU'd.
1333            if let Some(attr) = attrs {
1334                debug!(target: "engine::tree", head = canonical_header.number(), "handling payload attributes for canonical head");
1335                // Clone only when we actually need to process the attributes
1336                let updated =
1337                    self.process_payload_attributes(attr.clone(), &canonical_header, state);
1338                return Ok(Some(TreeOutcome::new(updated)));
1339            }
1340
1341            // The head block is already canonical and we're not processing payload attributes,
1342            // so we're not triggering a payload job and can return right away
1343            return Ok(Some(Self::valid_outcome(state)));
1344        }
1345
1346        // Ensure we can apply a new chain update for the head block
1347        if let Some(chain_update) = self.on_new_head(state.head_block_hash)? {
1348            if !self.is_consistent_forkchoice_state(state, Some(&chain_update))? {
1349                return Ok(Some(TreeOutcome::new(OnForkChoiceUpdated::invalid_state())));
1350            }
1351
1352            let tip = chain_update.tip().clone_sealed_header();
1353            self.on_canonical_chain_update(chain_update);
1354
1355            // Update the safe and finalized blocks and ensure their values are valid
1356            if let Err(outcome) = self.ensure_consistent_forkchoice_state(state) {
1357                // safe or finalized hashes are invalid
1358                return Ok(Some(TreeOutcome::new(outcome)));
1359            }
1360
1361            if let Some(attr) = attrs {
1362                // Clone only when we actually need to process the attributes
1363                let updated = self.process_payload_attributes(attr.clone(), &tip, state);
1364                return Ok(Some(TreeOutcome::new(updated)));
1365            }
1366
1367            return Ok(Some(Self::valid_outcome(state)));
1368        }
1369
1370        Ok(None)
1371    }
1372
1373    /// Handles the case where the head block is missing and needs to be downloaded.
1374    ///
1375    /// This is the fallback case when all other forkchoice update scenarios have been exhausted.
1376    /// Returns a `TreeOutcome` with syncing status and download event.
1377    fn handle_missing_block(
1378        &self,
1379        state: ForkchoiceState,
1380    ) -> ProviderResult<TreeOutcome<OnForkChoiceUpdated>> {
1381        // We don't have the block to perform the forkchoice update
1382        // We assume the FCU is valid and at least the head is missing,
1383        // so we need to start syncing to it
1384        //
1385        // find the appropriate target to sync to, if we don't have the safe block hash then we
1386        // start syncing to the safe block via backfill first
1387        let target = if self.state.forkchoice_state_tracker.is_empty() &&
1388        // check that safe block is valid and missing
1389        !state.safe_block_hash.is_zero() &&
1390        self.find_canonical_header(state.safe_block_hash).ok().flatten().is_none()
1391        {
1392            debug!(target: "engine::tree", "missing safe block on initial FCU, downloading safe block");
1393            state.safe_block_hash
1394        } else {
1395            state.head_block_hash
1396        };
1397
1398        let target = self.lowest_buffered_ancestor_or(target);
1399        trace!(target: "engine::tree", %target, "downloading missing block");
1400
1401        Ok(TreeOutcome::new(OnForkChoiceUpdated::valid(PayloadStatus::from_status(
1402            PayloadStatusEnum::Syncing,
1403        )))
1404        .with_event(TreeEvent::Download(
1405            DownloadRequest::single_block(target)
1406                .with_access_lists(self.should_download_access_lists()),
1407        )))
1408    }
1409
1410    /// Helper method to remove blocks and set the persistence state. This ensures we keep track of
1411    /// the current persistence action while we're removing blocks.
1412    fn remove_blocks(&mut self, new_tip_num: u64) {
1413        debug!(target: "engine::tree", ?new_tip_num, last_persisted_block_number=?self.persistence_state.last_persisted_block.number, "Removing blocks using persistence task");
1414        if new_tip_num < self.persistence_state.last_persisted_block.number {
1415            debug!(target: "engine::tree", ?new_tip_num, "Starting remove blocks job");
1416            self.state.set_pending_sparse_trie_prune(false);
1417            let (tx, rx) = crossbeam_channel::bounded(1);
1418            let _ = self.persistence.remove_blocks_above(new_tip_num, tx);
1419            self.persistence_state.start_remove(new_tip_num, rx);
1420        }
1421    }
1422
1423    /// Helper method to save blocks and set the persistence state. This ensures we keep track of
1424    /// the current persistence action while we're saving blocks.
1425    fn persist_blocks(&mut self, input: SaveBlocksInput<N>) {
1426        let highest_num_hash = input.last_block();
1427        debug!(target: "engine::tree", count=input.persist_rest_blocks().len(), blocks = ?input.persist_rest_blocks().iter().map(|block| block.recovered_block().num_hash()).collect::<Vec<_>>(), "Persisting blocks");
1428
1429        let (tx, rx) = crossbeam_channel::bounded(1);
1430        let _ = self.persistence.save_blocks(input, tx);
1431
1432        self.persistence_state.start_save(highest_num_hash, rx);
1433    }
1434
1435    /// Triggers new persistence actions if no persistence task is currently in progress.
1436    ///
1437    /// This checks if we need to remove blocks (disk reorg) or save new blocks to disk.
1438    /// Persistence completion is handled separately via the `wait_for_event` method.
1439    fn advance_persistence(&mut self) -> Result<(), AdvancePersistenceError> {
1440        if !self.persistence_state.in_progress() {
1441            let payload_build_active = self.payload_builds.is_active();
1442            if let Some(new_tip_num) = self.find_disk_reorg()? {
1443                self.remove_blocks(new_tip_num)
1444            } else if self.backfill_sync_state.is_pending_revalidation() &&
1445                !payload_build_active &&
1446                self.persistence_state.last_state_trie_persisted_block !=
1447                    self.persistence_state.last_persisted_block
1448            {
1449                let Some(input) = self.get_save_blocks_input(PersistTarget::Persisted) else {
1450                    return Err(AdvancePersistenceError::StateTrieCatchupUnavailable)
1451                };
1452                self.persist_blocks(input);
1453            } else if self.backfill_sync_state.is_pending_revalidation() && !payload_build_active {
1454                self.revalidate_pending_backfill()?;
1455            } else if let Some(input) = self.get_save_blocks_input(PersistTarget::Threshold) {
1456                self.persist_blocks(input);
1457            }
1458        }
1459
1460        Ok(())
1461    }
1462
1463    /// Finishes termination by persisting all remaining blocks and signaling completion.
1464    ///
1465    /// This blocks until all persistence is complete. Always signals completion,
1466    /// even if an error occurs.
1467    fn finish_termination(
1468        &mut self,
1469        pending_termination: oneshot::Sender<()>,
1470    ) -> Result<(), AdvancePersistenceError> {
1471        trace!(target: "engine::tree", "finishing termination, persisting remaining blocks");
1472        let result = self.persist_until_complete();
1473        let _ = pending_termination.send(());
1474        result
1475    }
1476
1477    /// Persists all remaining blocks until none are left.
1478    fn persist_until_complete(&mut self) -> Result<(), AdvancePersistenceError> {
1479        loop {
1480            // Wait for any in-progress persistence to complete (blocking)
1481            if let Some((rx, start_time, action)) = self.persistence_state.rx.take() {
1482                debug!(target: "engine::tree", ?action, "waiting for in-flight persistence");
1483                let result = rx.recv().map_err(|_| AdvancePersistenceError::ChannelClosed)?;
1484                self.on_persistence_complete(result, start_time)?;
1485                continue
1486            }
1487
1488            // Persistence can finish against a branch that the in-memory canonical chain has
1489            // reorged away from. Unwind that stale disk branch before building a head-targeted
1490            // save.
1491            if let Some(new_tip_num) = self.find_disk_reorg()? {
1492                self.remove_blocks(new_tip_num);
1493                continue
1494            }
1495
1496            let Some(input) = self.get_save_blocks_input(PersistTarget::Head) else {
1497                debug!(target: "engine::tree", "persistence complete, signaling termination");
1498                return Ok(())
1499            };
1500
1501            debug!(target: "engine::tree", count = input.persist_rest_blocks().len(), "persisting remaining blocks before shutdown");
1502            self.persist_blocks(input);
1503        }
1504    }
1505
1506    /// Tries to poll for a completed persistence task (non-blocking).
1507    ///
1508    /// Returns `true` if a persistence task was completed, `false` otherwise.
1509    fn try_poll_persistence(&mut self) -> Result<bool, AdvancePersistenceError> {
1510        let Some((rx, start_time, action)) = self.persistence_state.rx.take() else {
1511            return Ok(false);
1512        };
1513
1514        match rx.try_recv() {
1515            Ok(result) => {
1516                self.on_persistence_complete(result, start_time)?;
1517                Ok(true)
1518            }
1519            Err(crossbeam_channel::TryRecvError::Empty) => {
1520                // Not ready yet, put it back
1521                self.persistence_state.rx = Some((rx, start_time, action));
1522                Ok(false)
1523            }
1524            Err(crossbeam_channel::TryRecvError::Disconnected) => {
1525                Err(AdvancePersistenceError::ChannelClosed)
1526            }
1527        }
1528    }
1529
1530    /// Handles a completed persistence task.
1531    fn on_persistence_complete(
1532        &mut self,
1533        result: PersistenceResult,
1534        start_time: Instant,
1535    ) -> Result<(), AdvancePersistenceError> {
1536        self.metrics.engine.persistence_duration.record(start_time.elapsed());
1537
1538        let PersistenceResult { last_block, last_state_trie_block, commit_duration } = result;
1539        debug_assert!(
1540            last_state_trie_block.number <= last_block.number,
1541            "state/trie frontier cannot exceed the last persisted block"
1542        );
1543
1544        debug!(target: "engine::tree", ?last_block, ?last_state_trie_block, elapsed=?start_time.elapsed(), "Finished persisting, calling finish");
1545        self.persistence_state.finish(last_block, last_state_trie_block);
1546
1547        let last_block_number = last_block.number;
1548
1549        // Evict cached changesets for blocks below the eviction threshold.
1550        // Keep at least CHANGESET_CACHE_RETENTION_BLOCKS from the persisted tip, and also respect
1551        // the finalized block if set.
1552        let min_threshold = last_block_number.saturating_sub(CHANGESET_CACHE_RETENTION_BLOCKS);
1553        let eviction_threshold =
1554            if let Some(finalized) = self.canonical_in_memory_state.get_finalized_num_hash() {
1555                // Use the minimum of finalized block and retention threshold to be conservative
1556                finalized.number.min(min_threshold)
1557            } else {
1558                // When finalized is not set (e.g., on L2s), use the retention threshold
1559                min_threshold
1560            };
1561        debug!(
1562            target: "engine::tree",
1563            last_persisted = last_block_number,
1564            finalized_number = ?self.canonical_in_memory_state.get_finalized_num_hash().map(|f| f.number),
1565            eviction_threshold,
1566            "Evicting changesets below threshold"
1567        );
1568        self.state.tree_state.overlay_manager.evict_cached_changesets(eviction_threshold);
1569
1570        self.on_new_persisted_block()?;
1571
1572        self.purge_timing_stats(last_block_number, commit_duration);
1573
1574        Ok(())
1575    }
1576
1577    /// Handles a message from the engine.
1578    ///
1579    /// Returns `ControlFlow::Break(())` if the engine should terminate.
1580    fn on_engine_message(
1581        &mut self,
1582        msg: FromEngine<EngineApiRequest<T, N>, N::Block>,
1583    ) -> Result<ops::ControlFlow<()>, InsertBlockFatalError> {
1584        match msg {
1585            FromEngine::Event(event) => match event {
1586                FromOrchestrator::BackfillSyncStarted => {
1587                    debug!(target: "engine::tree", "received backfill sync started event");
1588                    self.backfill_sync_state = BackfillSyncState::Active;
1589                }
1590                FromOrchestrator::BackfillSyncFinished(ctrl) => {
1591                    self.on_backfill_sync_finished(ctrl)?;
1592                }
1593                FromOrchestrator::Terminate { tx } => {
1594                    debug!(target: "engine::tree", "received terminate request");
1595                    if let Err(err) = self.finish_termination(tx) {
1596                        error!(target: "engine::tree", %err, "Termination failed");
1597                    }
1598                    return Ok(ops::ControlFlow::Break(()))
1599                }
1600            },
1601            FromEngine::Request(request) => {
1602                match request {
1603                    EngineApiRequest::InsertExecutedBlock(payload) => {
1604                        let block_num_hash = payload.recovered_block.num_hash();
1605                        if block_num_hash.number <= self.state.tree_state.canonical_block_number() {
1606                            // outdated block that can be skipped
1607                            return Ok(ops::ControlFlow::Continue(()))
1608                        }
1609
1610                        if self.state.tree_state.contains_hash(&block_num_hash.hash) {
1611                            // block already known to the tree (e.g. delivered via newPayload first)
1612                            return Ok(ops::ControlFlow::Continue(()))
1613                        }
1614
1615                        debug!(target: "engine::tree", block=?block_num_hash, "inserting already executed block");
1616                        let now = Instant::now();
1617
1618                        let block = match self.payload_validator.on_inserted_executed_block(payload)
1619                        {
1620                            Ok(block) => block,
1621                            Err(err) => {
1622                                warn!(target: "engine::tree", %err, block=?block_num_hash, "Failed to insert already executed block");
1623                                return Ok(ops::ControlFlow::Continue(()))
1624                            }
1625                        };
1626
1627                        let is_pending = self.state.tree_state.canonical_block_hash() ==
1628                            block.recovered_block().parent_hash();
1629                        self.state.tree_state.insert_executed(block.clone());
1630                        self.metrics
1631                            .engine
1632                            .executed_blocks
1633                            .set(self.state.tree_state.block_count() as f64);
1634
1635                        if is_pending {
1636                            debug!(target: "engine::tree", pending=?block_num_hash, "updating pending block");
1637                            self.canonical_in_memory_state.set_pending_block(block.clone());
1638                        }
1639
1640                        self.metrics.engine.inserted_already_executed_blocks.increment(1);
1641                        self.emit_event(EngineApiEvent::BeaconConsensus(
1642                            ConsensusEngineEvent::CanonicalBlockAdded(block, now.elapsed()),
1643                        ));
1644                    }
1645                    EngineApiRequest::Beacon(request) => {
1646                        match request {
1647                            BeaconEngineMessage::ForkchoiceUpdated {
1648                                cause,
1649                                state,
1650                                payload_attrs,
1651                                tx,
1652                            } => {
1653                                let _cause = cause.enter();
1654                                let has_attrs = payload_attrs.is_some();
1655
1656                                let start = Instant::now();
1657                                let mut output = self.on_forkchoice_updated(state, payload_attrs);
1658
1659                                if let Ok(res) = &mut output {
1660                                    // track last received forkchoice state
1661                                    self.state
1662                                        .forkchoice_state_tracker
1663                                        .set_latest(state, res.outcome.forkchoice_status());
1664
1665                                    // emit an event about the handled FCU
1666                                    self.emit_event(ConsensusEngineEvent::ForkchoiceUpdated(
1667                                        state,
1668                                        res.outcome.forkchoice_status(),
1669                                    ));
1670
1671                                    // handle the event if any
1672                                    self.on_maybe_tree_event(res.event.take())?;
1673                                }
1674
1675                                if let Err(ref err) = output {
1676                                    error!(target: "engine::tree", %err, ?state, "Error processing forkchoice update");
1677                                }
1678
1679                                self.metrics.engine.forkchoice_updated.update_response_metrics(
1680                                    start,
1681                                    &mut self.metrics.engine.new_payload.latest_finish_at,
1682                                    has_attrs,
1683                                    &output,
1684                                );
1685
1686                                if let Err(err) =
1687                                    tx.send(output.map(|o| o.outcome).map_err(Into::into))
1688                                {
1689                                    self.metrics
1690                                        .engine
1691                                        .failed_forkchoice_updated_response_deliveries
1692                                        .increment(1);
1693                                    warn!(target: "engine::tree", ?state, elapsed=?start.elapsed(), "Failed to deliver forkchoiceUpdated response, receiver dropped (request cancelled): {err:?}");
1694                                }
1695                            }
1696                            BeaconEngineMessage::NewPayload { cause, payload, tx } => {
1697                                let _cause = cause.enter();
1698                                let start = Instant::now();
1699                                let gas_used = payload.gas_used();
1700                                let num_hash = payload.num_hash();
1701                                let mut output = self.on_new_payload(payload);
1702                                self.metrics.engine.new_payload.update_response_metrics(
1703                                    start,
1704                                    &mut self.metrics.engine.forkchoice_updated.latest_finish_at,
1705                                    &output,
1706                                    gas_used,
1707                                );
1708
1709                                let maybe_event =
1710                                    output.as_mut().ok().and_then(|out| out.event.take());
1711
1712                                // emit response
1713                                if let Err(err) =
1714                                    tx.send(output.map(|o| o.outcome).map_err(Into::into))
1715                                {
1716                                    warn!(target: "engine::tree", payload=?num_hash, elapsed=?start.elapsed(), "Failed to deliver newPayload response, receiver dropped (request cancelled): {err:?}");
1717                                    self.metrics
1718                                        .engine
1719                                        .failed_new_payload_response_deliveries
1720                                        .increment(1);
1721                                }
1722
1723                                // handle the event if any
1724                                self.on_maybe_tree_event(maybe_event)?;
1725                            }
1726                            BeaconEngineMessage::RethNewPayload {
1727                                cause,
1728                                payload,
1729                                wait_for_persistence,
1730                                wait_for_caches,
1731                                tx,
1732                                enqueued_at,
1733                            } => {
1734                                let _cause = cause.enter();
1735                                debug!(
1736                                    target: "engine::tree",
1737                                    wait_for_persistence,
1738                                    wait_for_caches,
1739                                    "Processing reth_newPayload"
1740                                );
1741
1742                                let backpressure_wait = enqueued_at.elapsed();
1743
1744                                let explicit_persistence_wait = if wait_for_persistence {
1745                                    let pending_persistence = self.persistence_state.rx.take();
1746                                    if let Some((rx, start_time, _action)) = pending_persistence {
1747                                        let (persistence_tx, persistence_rx) =
1748                                            std::sync::mpsc::channel();
1749                                        self.runtime.spawn_blocking_named(
1750                                            "wait-persist",
1751                                            move || {
1752                                                let start = Instant::now();
1753                                                let result = rx
1754                                                    .recv()
1755                                                    .expect("persistence state channel closed");
1756                                                let _ = persistence_tx.send((
1757                                                    result,
1758                                                    start_time,
1759                                                    start.elapsed(),
1760                                                ));
1761                                            },
1762                                        );
1763                                        let (result, start_time, wait_duration) = persistence_rx
1764                                            .recv()
1765                                            .expect("persistence result channel closed");
1766                                        let _ = self.on_persistence_complete(result, start_time);
1767                                        wait_duration
1768                                    } else {
1769                                        Duration::ZERO
1770                                    }
1771                                } else {
1772                                    Duration::ZERO
1773                                };
1774
1775                                let cache_wait = wait_for_caches
1776                                    .then(|| self.payload_validator.wait_for_caches());
1777
1778                                let start = Instant::now();
1779                                let gas_used = payload.gas_used();
1780                                let num_hash = payload.num_hash();
1781                                let mut output = self.on_new_payload(payload);
1782                                let latency = start.elapsed();
1783                                self.metrics.engine.new_payload.update_response_metrics(
1784                                    start,
1785                                    &mut self.metrics.engine.forkchoice_updated.latest_finish_at,
1786                                    &output,
1787                                    gas_used,
1788                                );
1789
1790                                let maybe_event =
1791                                    output.as_mut().ok().and_then(|out| out.event.take());
1792
1793                                let timings = NewPayloadTimings {
1794                                    latency,
1795                                    persistence_wait: backpressure_wait + explicit_persistence_wait,
1796                                    execution_cache_wait: cache_wait
1797                                        .map(|wait| wait.execution_cache),
1798                                    sparse_trie_wait: cache_wait.map(|wait| wait.sparse_trie),
1799                                };
1800                                if let Err(err) = tx
1801                                    .send(output.map(|o| (o.outcome, timings)).map_err(Into::into))
1802                                {
1803                                    error!(
1804                                        target: "engine::tree",
1805                                        payload=?num_hash,
1806                                        elapsed=?latency,
1807                                        "Failed to send event: {err:?}"
1808                                    );
1809                                    self.metrics
1810                                        .engine
1811                                        .failed_new_payload_response_deliveries
1812                                        .increment(1);
1813                                }
1814
1815                                self.on_maybe_tree_event(maybe_event)?;
1816                            }
1817                        }
1818                    }
1819                }
1820            }
1821            FromEngine::DownloadedBlocks(blocks) => {
1822                if let Some(event) = self.on_downloaded(blocks)? {
1823                    self.on_tree_event(event)?;
1824                }
1825            }
1826        }
1827        Ok(ops::ControlFlow::Continue(()))
1828    }
1829
1830    /// Invoked if the backfill sync has finished to target.
1831    ///
1832    /// At this point we consider the block synced to the backfill target.
1833    ///
1834    /// Checks the tracked finalized block against the block on disk and requests another backfill
1835    /// run if the distance to the tip exceeds the threshold for another backfill run.
1836    ///
1837    /// This will also do the necessary housekeeping of the tree state, this includes:
1838    ///  - removing all blocks below the backfill height
1839    ///  - resetting the canonical in-memory state
1840    ///
1841    /// In case backfill resulted in an unwind, this will clear the tree state above the unwind
1842    /// target block.
1843    fn on_backfill_sync_finished(
1844        &mut self,
1845        ctrl: ControlFlow,
1846    ) -> Result<(), InsertBlockFatalError> {
1847        debug!(target: "engine::tree", "received backfill sync finished event");
1848        self.backfill_sync_state = BackfillSyncState::Idle;
1849
1850        // Pipeline unwound, memorize the invalid block and wait for CL for next sync target.
1851        let backfill_height = if let ControlFlow::Unwind { bad_block, target } = &ctrl {
1852            warn!(target: "engine::tree", invalid_block=?bad_block, "Bad block detected in unwind");
1853            // update the `invalid_headers` cache with the new invalid header
1854            self.state.invalid_headers.insert(**bad_block);
1855
1856            // if this was an unwind then the target is the new height
1857            Some(*target)
1858        } else {
1859            // backfill height is the block number that the backfill finished at
1860            ctrl.block_number()
1861        };
1862
1863        // backfill height is the block number that the backfill finished at
1864        let Some(backfill_height) = backfill_height else { return Ok(()) };
1865
1866        // state house keeping after backfill sync
1867        // remove all executed blocks below the backfill height
1868        //
1869        // We set the `finalized_num` to `Some(backfill_height)` to ensure we remove all state
1870        // before that
1871        let Some(backfill_num_hash) = self
1872            .provider
1873            .block_hash(backfill_height)?
1874            .map(|hash| BlockNumHash { hash, number: backfill_height })
1875        else {
1876            debug!(target: "engine::tree", ?ctrl, "Backfill block not found");
1877            return Ok(())
1878        };
1879
1880        if ctrl.is_unwind() {
1881            // the node reset so we need to clear everything above that height so that backfill
1882            // height is the new canonical block.
1883            self.state.set_pending_sparse_trie_prune(false);
1884            self.state.tree_state.reset(backfill_num_hash)
1885        } else {
1886            self.state.tree_state.remove_until(
1887                backfill_num_hash,
1888                self.persistence_state.last_persisted_block.hash,
1889                Some(backfill_num_hash),
1890            );
1891        }
1892
1893        self.metrics.engine.executed_blocks.set(self.state.tree_state.block_count() as f64);
1894        self.metrics.tree.canonical_chain_height.set(backfill_height as f64);
1895
1896        // remove all buffered blocks below the backfill height
1897        self.state.buffer.remove_old_blocks(backfill_height);
1898        self.purge_timing_stats(backfill_height, None);
1899        // we remove all entries because now we're synced to the backfill target and consider this
1900        // the canonical chain
1901        self.canonical_in_memory_state.clear_state();
1902
1903        if let Ok(Some(new_head)) = self.provider.sealed_header(backfill_height) {
1904            // update the tracked chain height, after backfill sync both the canonical height and
1905            // persisted height are the same
1906            self.state.tree_state.set_canonical_head(new_head.num_hash());
1907            self.persistence_state.finish(new_head.num_hash(), new_head.num_hash());
1908
1909            // update the tracked canonical head
1910            self.canonical_in_memory_state.set_canonical_head(new_head);
1911
1912            // If the pipeline reached the head of the syncing FCU, apply its safe and finalized
1913            // blocks now that the head is canonical. An unwind must wait for a new FCU because the
1914            // supplied sync target may have been invalidated.
1915            if !ctrl.is_unwind() {
1916                self.on_canonicalized_sync_target(backfill_num_hash.hash);
1917            }
1918        }
1919
1920        // check if we need to run backfill again by comparing the most recent backfill target
1921        // height to the backfill height
1922        let Some(sync_target_state) = self.state.forkchoice_state_tracker.sync_target_state()
1923        else {
1924            return Ok(())
1925        };
1926        if !self.engine_kind.is_opstack() && sync_target_state.finalized_block_hash.is_zero() {
1927            // no finalized block, can't check distance on non-OP Stack chains
1928            return Ok(())
1929        }
1930        let target_hash = self.backfill_target_hash(sync_target_state);
1931        if target_hash.is_zero() {
1932            return Ok(())
1933        }
1934        // get the block number of the backfill target block, if we have it buffered
1935        let newest_target = self.state.buffer.block(&target_hash).map(|block| block.number());
1936
1937        // The block number that the backfill finished at - if the progress or newest target is
1938        // None then we can't check the distance anyways.
1939        //
1940        // If both are Some, we perform another distance check and return the desired backfill
1941        // target
1942        if let Some(backfill_target) =
1943            ctrl.block_number().zip(newest_target).and_then(|(progress, target_number)| {
1944                // Determines whether or not we should run backfill again, in case
1945                // the new gap is still large enough and requires running backfill again
1946                self.backfill_sync_target(progress, target_number, None)
1947            })
1948        {
1949            // request another backfill run
1950            self.emit_event(EngineApiEvent::BackfillAction(BackfillAction::Start(
1951                backfill_target.into(),
1952            )));
1953            return Ok(())
1954        };
1955
1956        // Check if there are more blocks to sync between current head and FCU target
1957        if let Some(lowest_buffered) =
1958            self.state.buffer.lowest_ancestor(&sync_target_state.head_block_hash)
1959        {
1960            let current_head_num = self.state.tree_state.current_canonical_head.number;
1961            let target_head_num = lowest_buffered.number();
1962
1963            if let Some(distance) = self.distance_from_local_tip(current_head_num, target_head_num)
1964            {
1965                // There are blocks between current head and FCU target, download them
1966                debug!(
1967                    target: "engine::tree",
1968                    %current_head_num,
1969                    %target_head_num,
1970                    %distance,
1971                    "Backfill complete, downloading remaining blocks to reach FCU target"
1972                );
1973
1974                self.emit_event(EngineApiEvent::Download(
1975                    DownloadRequest::block_range(lowest_buffered.parent_hash(), distance)
1976                        .with_access_lists(self.should_download_access_lists()),
1977                ));
1978                return Ok(());
1979            }
1980        } else {
1981            // We don't have the head block or any of its ancestors buffered. Request
1982            // a download for the head block which will then trigger further sync.
1983            debug!(
1984                target: "engine::tree",
1985                head_hash = %sync_target_state.head_block_hash,
1986                "Backfill complete but head block not buffered, requesting download"
1987            );
1988            self.emit_event(EngineApiEvent::Download(
1989                DownloadRequest::single_block(sync_target_state.head_block_hash)
1990                    .with_access_lists(self.should_download_access_lists()),
1991            ));
1992            return Ok(());
1993        }
1994
1995        // try to close the gap by executing buffered blocks that are child blocks of the new head
1996        self.try_connect_buffered_blocks(self.state.tree_state.current_canonical_head)
1997    }
1998
1999    /// Attempts to make the given target canonical.
2000    ///
2001    /// This will update the tracked canonical in memory state and do the necessary housekeeping.
2002    fn make_canonical(&mut self, target: B256) -> ProviderResult<()> {
2003        if let Some(chain_update) = self.on_new_head(target)? {
2004            self.on_canonical_chain_update(chain_update);
2005        }
2006
2007        self.on_canonicalized_sync_target(target);
2008
2009        Ok(())
2010    }
2011
2012    /// Applies the tracked forkchoice state once its sync target head becomes canonical.
2013    fn on_canonicalized_sync_target(&mut self, target: B256) {
2014        let Some(sync_target_state) = self
2015            .state
2016            .forkchoice_state_tracker
2017            .sync_target_state()
2018            .filter(|state| state.head_block_hash == target)
2019        else {
2020            return;
2021        };
2022
2023        if let Err(outcome) = self.ensure_consistent_forkchoice_state(sync_target_state) {
2024            debug!(
2025                target: "engine::tree",
2026                head = %sync_target_state.head_block_hash,
2027                safe = %sync_target_state.safe_block_hash,
2028                finalized = %sync_target_state.finalized_block_hash,
2029                ?outcome,
2030                "Canonicalized sync target head before safe/finalized could be applied"
2031            );
2032            return;
2033        }
2034
2035        self.state.forkchoice_state_tracker.promote_sync_target_to_valid(sync_target_state);
2036    }
2037
2038    /// Convenience function to handle an optional tree event.
2039    fn on_maybe_tree_event(&mut self, event: Option<TreeEvent>) -> ProviderResult<()> {
2040        if let Some(event) = event {
2041            self.on_tree_event(event)?;
2042        }
2043
2044        Ok(())
2045    }
2046
2047    /// Handles a tree event.
2048    ///
2049    /// Returns an error if a [`TreeAction::MakeCanonical`] results in a fatal error.
2050    fn on_tree_event(&mut self, event: TreeEvent) -> ProviderResult<()> {
2051        match event {
2052            TreeEvent::TreeAction(action) => match action {
2053                TreeAction::MakeCanonical { sync_target_head } => {
2054                    self.make_canonical(sync_target_head)?;
2055                }
2056            },
2057            TreeEvent::BackfillAction(action) => {
2058                self.emit_event(EngineApiEvent::BackfillAction(action));
2059            }
2060            TreeEvent::Download(action) => {
2061                self.emit_event(EngineApiEvent::Download(action));
2062            }
2063        }
2064
2065        Ok(())
2066    }
2067
2068    /// Removes timing stats for blocks at or below `below_number`.
2069    ///
2070    /// No-op when detailed block logging is disabled (no stats are recorded in that case).
2071    /// When `commit_duration` is provided and a slow block threshold is configured, checks
2072    /// each removed block against the threshold and emits a [`ConsensusEngineEvent::SlowBlock`]
2073    /// event for blocks that exceed it.
2074    fn purge_timing_stats(&mut self, below_number: u64, commit_duration: Option<Duration>) {
2075        let threshold = self.config.slow_block_threshold();
2076        let check_slow = commit_duration.is_some() && threshold.is_some();
2077
2078        // Two-pass: collect keys first because emit_event borrows &mut self.
2079        let keys_to_remove: Vec<B256> = self
2080            .execution_timing_stats
2081            .iter()
2082            .filter(|(_, stats)| stats.block_number <= below_number)
2083            .map(|(k, _)| *k)
2084            .collect();
2085
2086        for key in keys_to_remove {
2087            let stats = self.execution_timing_stats.remove(&key).expect("key just found");
2088            if check_slow {
2089                let commit_dur = commit_duration.expect("checked above");
2090                // state_read_duration is already included in execution_duration
2091                let total_duration =
2092                    stats.execution_duration + stats.state_hash_duration + commit_dur;
2093
2094                if total_duration > threshold.expect("checked above") {
2095                    self.emit_event(ConsensusEngineEvent::SlowBlock(SlowBlockInfo {
2096                        stats,
2097                        commit_duration: Some(commit_dur),
2098                        total_duration,
2099                    }));
2100                }
2101            }
2102        }
2103    }
2104
2105    /// Re-evaluates whether a deferred backfill is still required after persistence catches up.
2106    fn revalidate_pending_backfill(&mut self) -> ProviderResult<()> {
2107        debug_assert!(self.backfill_sync_state.is_pending_revalidation());
2108
2109        let sync_target_state = self.state.forkchoice_state_tracker.sync_target_state();
2110        let backfill_target = if let Some(state) = sync_target_state {
2111            let configured_target = self.backfill_target_hash(state);
2112            let target_hash =
2113                if configured_target.is_zero() { state.head_block_hash } else { configured_target };
2114            let target_number = if let Some(block) = self.state.buffer.block(&target_hash) {
2115                Some(block.number())
2116            } else {
2117                self.sealed_header_by_hash(target_hash)?.map(|header| header.number())
2118            };
2119
2120            target_number.and_then(|target_number| {
2121                self.backfill_sync_target(
2122                    self.state.tree_state.canonical_block_number(),
2123                    target_number,
2124                    None,
2125                )
2126            })
2127        } else {
2128            None
2129        };
2130
2131        if let Some(target) = backfill_target {
2132            self.dispatch_backfill_action(BackfillAction::Start(target.into()));
2133            return Ok(())
2134        }
2135
2136        self.backfill_sync_state = BackfillSyncState::Idle;
2137        debug!(target: "engine::tree", "dropping deferred backfill after re-evaluation");
2138
2139        // The target may have changed while persistence was draining. Resume the live-sync
2140        // download flow so a newer target can produce a fresh backfill decision.
2141        if let Some(state) = sync_target_state &&
2142            state.head_block_hash != self.state.tree_state.canonical_block_hash()
2143        {
2144            let target = self.lowest_buffered_ancestor_or(state.head_block_hash);
2145            self.send_event(EngineApiEvent::Download(DownloadRequest::single_block(target)));
2146        }
2147
2148        Ok(())
2149    }
2150
2151    /// Emits an outgoing event to the engine.
2152    fn emit_event(&mut self, event: impl Into<EngineApiEvent<N>>) {
2153        let event = event.into();
2154
2155        if let EngineApiEvent::BackfillAction(action) = event {
2156            debug_assert_eq!(
2157                self.backfill_sync_state,
2158                BackfillSyncState::Idle,
2159                "backfill action should only be emitted when backfill is idle"
2160            );
2161
2162            let persistence_in_progress = self.persistence_state.in_progress();
2163            let state_trie_needs_catchup = self.persistence_state.last_state_trie_persisted_block !=
2164                self.persistence_state.last_persisted_block;
2165            if self.payload_builds.is_active() ||
2166                persistence_in_progress ||
2167                state_trie_needs_catchup
2168            {
2169                // Backfill can remove the same in-memory blocks as an active payload job or a
2170                // persistence task. Enter pending mode to prevent new payload jobs and
2171                // re-evaluate the current sync target after all readers and writes drain.
2172                debug!(
2173                    target: "engine::tree",
2174                    last_persisted_block = self.persistence_state.last_persisted_block.number,
2175                    last_state_trie_persisted_block = self
2176                        .persistence_state
2177                        .last_state_trie_persisted_block
2178                        .number,
2179                    "deferring backfill until persistence and payload jobs drain"
2180                );
2181                self.backfill_sync_state = BackfillSyncState::PendingRevalidation;
2182                return
2183            }
2184
2185            self.dispatch_backfill_action(action);
2186            return
2187        }
2188
2189        self.send_event(event);
2190    }
2191
2192    /// Dispatches a validated backfill action to the orchestrator.
2193    fn dispatch_backfill_action(&mut self, action: BackfillAction) {
2194        debug_assert!(
2195            self.backfill_sync_state.is_idle() ||
2196                self.backfill_sync_state.is_pending_revalidation(),
2197            "backfill action can only be dispatched while idle or pending revalidation"
2198        );
2199        self.backfill_sync_state = BackfillSyncState::Pending;
2200        self.metrics.engine.pipeline_runs.increment(1);
2201        debug!(target: "engine::tree", "emitting backfill action event");
2202        self.send_event(EngineApiEvent::BackfillAction(action));
2203    }
2204
2205    /// Sends an event to the orchestrator.
2206    fn send_event(&self, event: EngineApiEvent<N>) {
2207        let _ = self.outgoing.send(event).inspect_err(
2208            |err| error!(target: "engine::tree", "Failed to send internal event: {err:?}"),
2209        );
2210    }
2211
2212    /// Returns the blocks and frontiers for the next persistence cycle, if one should start.
2213    ///
2214    /// Threshold persistence honors the normal scheduling gates and retains the configured
2215    /// in-memory block buffer. Persisted-target persistence catches the state/trie frontier up to
2216    /// the existing database tip. Head persistence bypasses those gates during shutdown and
2217    /// returns `None` once both persistence frontiers have reached the canonical head.
2218    fn get_save_blocks_input(&self, target: PersistTarget) -> Option<SaveBlocksInput<N>> {
2219        // We will calculate the state root using the database, so we need to be sure there are no
2220        // changes
2221        debug_assert!(!self.persistence_state.in_progress());
2222
2223        let prev_partial_state_trie = self.persistence_state.last_state_trie_persisted_block.number;
2224        let prev_db_tip = self.persistence_state.last_persisted_block.number;
2225        let canonical_head_number = self.state.tree_state.canonical_block_number();
2226
2227        let (new_db_tip, new_partial_state_trie) = match target {
2228            PersistTarget::Head => (canonical_head_number, canonical_head_number),
2229            PersistTarget::Persisted => {
2230                // Catch-up persistence is the transition into pipeline sync, so it deliberately
2231                // runs while backfill is pending and bypasses the normal threshold gates.
2232                debug_assert!(self.backfill_sync_state.is_pending_revalidation());
2233                debug_assert!(!self.payload_builds.is_active());
2234                (prev_db_tip, prev_db_tip)
2235            }
2236            PersistTarget::Threshold => {
2237                if (self.config.suppress_persistence_during_build() &&
2238                    self.payload_builds.is_active()) ||
2239                    !self.backfill_sync_state.is_idle()
2240                {
2241                    return None
2242                }
2243
2244                let persistence_threshold =
2245                    usize::try_from(self.config.persistence_threshold()).unwrap_or(usize::MAX);
2246                if self.canonical_in_memory_state.canonical_chain().count() <= persistence_threshold
2247                {
2248                    return None
2249                }
2250
2251                let new_db_tip =
2252                    canonical_head_number.saturating_sub(self.config.memory_block_buffer_target());
2253                if new_db_tip <= prev_db_tip {
2254                    return None
2255                }
2256
2257                let new_partial_state_trie = new_db_tip
2258                    .saturating_sub(self.config.num_state_masking_blocks())
2259                    .max(prev_partial_state_trie);
2260                (new_db_tip, new_partial_state_trie)
2261            }
2262        };
2263
2264        debug_assert!(
2265            new_db_tip >= prev_db_tip,
2266            "disk reorg must be resolved before saving blocks"
2267        );
2268        debug_assert!(
2269            new_partial_state_trie >= prev_partial_state_trie,
2270            "disk reorg must be resolved before saving state/trie"
2271        );
2272
2273        if new_db_tip == prev_db_tip && new_partial_state_trie == prev_partial_state_trie {
2274            return None
2275        }
2276
2277        let mut blocks = Vec::new();
2278        let mut current_hash = self.state.tree_state.canonical_block_hash();
2279
2280        debug!(
2281            target: "engine::tree",
2282            ?current_hash,
2283            ?prev_partial_state_trie,
2284            ?prev_db_tip,
2285            ?canonical_head_number,
2286            ?new_partial_state_trie,
2287            ?new_db_tip,
2288            target = ?target,
2289            "Returning save input"
2290        );
2291        while let Some(block) = self.state.tree_state.blocks_by_hash.get(&current_hash) {
2292            if block.recovered_block().number() <= prev_partial_state_trie {
2293                break;
2294            }
2295
2296            if block.recovered_block().number() <= new_db_tip {
2297                blocks.push(block.clone());
2298            }
2299
2300            current_hash = block.recovered_block().parent_hash();
2301        }
2302
2303        // Reverse the order so that the oldest block comes first
2304        blocks.reverse();
2305
2306        Some(SaveBlocksInput::new(
2307            blocks,
2308            prev_db_tip,
2309            prev_partial_state_trie,
2310            new_db_tip,
2311            new_partial_state_trie,
2312        ))
2313    }
2314
2315    /// This clears the blocks from the in-memory tree state that have been persisted to the
2316    /// database.
2317    ///
2318    /// This also updates the canonical in-memory state to reflect the newest persisted block
2319    /// height.
2320    ///
2321    /// Assumes that `finish` has been called on the `persistence_state` at least once
2322    fn on_new_persisted_block(&mut self) -> ProviderResult<()> {
2323        let in_memory_persisted_block = self.persistence_state.last_state_trie_persisted_block;
2324
2325        // If we have an on-disk reorg, we need to handle it first before touching the in-memory
2326        // state.
2327        if let Some(remove_above) = self.find_disk_reorg()? {
2328            self.remove_blocks(remove_above);
2329            return Ok(())
2330        }
2331
2332        let finalized = self.state.forkchoice_state_tracker.last_valid_finalized();
2333        // Trim the canonical in-memory state first: state providers build their overlays from the
2334        // canonical chain, so it must never reference blocks whose overlays the manager has
2335        // already pruned. `remove_before` does not read the canonical in-memory state, so the
2336        // order between the two trims is free to choose.
2337        self.canonical_in_memory_state.remove_persisted_blocks_until(
2338            self.persistence_state.last_persisted_block,
2339            in_memory_persisted_block.number,
2340        );
2341        self.remove_before(in_memory_persisted_block, finalized)?;
2342        // Persistence changes the overlay anchor. Prepare the remaining canonical range before
2343        // the next payload needs to read execution state against the new durable frontier.
2344        self.state.tree_state.overlay_manager.precompute_execution_overlay(
2345            self.state.tree_state.canonical_block_hash(),
2346            in_memory_persisted_block.hash,
2347        );
2348        self.state.set_pending_sparse_trie_prune(self.should_prune_sparse_trie());
2349        Ok(())
2350    }
2351
2352    /// Returns whether sparse trie pruning should be attempted by the next sparse trie task.
2353    const fn should_prune_sparse_trie(&self) -> bool {
2354        self.config.use_state_root_task()
2355    }
2356
2357    /// Return an [`ExecutedBlock`] from database or in-memory state by hash.
2358    ///
2359    /// Note: This function attempts to fetch the `ExecutedBlock` from either in-memory state
2360    /// or the database. If the required historical data (such as trie change sets) has been
2361    /// pruned for a given block, this operation will return an error. On archive nodes, it
2362    /// can retrieve any block.
2363    #[instrument(level = "debug", target = "engine::tree", skip(self))]
2364    fn canonical_block_by_hash(&self, hash: B256) -> ProviderResult<ExecutedBlock<N>> {
2365        trace!(target: "engine::tree", ?hash, "Fetching executed block by hash");
2366        // check memory first
2367        if let Some(block) = self.state.tree_state.executed_block_by_hash(hash) {
2368            return Ok(block.clone())
2369        }
2370
2371        let (block, senders) = self
2372            .provider
2373            .sealed_block_with_senders(hash.into(), TransactionVariant::WithHash)?
2374            .ok_or_else(|| ProviderError::HeaderNotFound(hash.into()))?
2375            .split_sealed();
2376        let mut execution_output = self
2377            .provider
2378            .get_state(block.header().number())?
2379            .ok_or_else(|| ProviderError::StateForNumberNotFound(block.header().number()))?;
2380        let bundle_state = execution_output.state();
2381        // `get_state` can return an in-memory execution outcome that retains destruction statuses.
2382        // Hashing it requires the parent provider to expand a pre-existing destroyed account's
2383        // storage into zero-valued slots. Genesis has no parent state, and no account existed
2384        // before it to be destroyed.
2385        let hashed_state = if block.parent_hash().is_zero() {
2386            HashedPostState::from_bundle_state::<KeccakKeyHasher>(bundle_state.state())
2387        } else {
2388            self.provider
2389                .state_by_block_hash(block.parent_hash())?
2390                .hashed_post_state(bundle_state)?
2391        };
2392
2393        debug!(
2394            target: "engine::tree",
2395            number = ?block.number(),
2396            "computing block trie updates",
2397        );
2398        let db_provider = self.provider.database_provider_ro()?;
2399        let trie_updates = self
2400            .state
2401            .tree_state
2402            .overlay_manager
2403            .compute_block_trie_updates(&db_provider, block.number())?;
2404
2405        let sorted_hashed_state = Arc::new(hashed_state.into_sorted());
2406        let sorted_trie_updates = Arc::new(trie_updates);
2407
2408        let execution_output = Arc::new(BlockExecutionOutput {
2409            state: execution_output.bundle,
2410            result: BlockExecutionResult {
2411                receipts: execution_output.receipts.pop().unwrap_or_default(),
2412                requests: execution_output.requests.pop().unwrap_or_default(),
2413                gas_used: block.gas_used(),
2414                blob_gas_used: block.blob_gas_used().unwrap_or_default(),
2415            },
2416        });
2417
2418        Ok(ExecutedBlock::new(
2419            Arc::new(RecoveredBlock::new_sealed(block, senders)),
2420            execution_output,
2421            sorted_hashed_state,
2422            sorted_trie_updates,
2423        ))
2424    }
2425
2426    /// Returns `true` if a block with the given hash is known, either in memory or in the
2427    /// database. This is a lightweight existence check that avoids constructing a full
2428    /// [`SealedHeader`].
2429    fn has_block_by_hash(&self, hash: B256) -> ProviderResult<bool> {
2430        if self.state.tree_state.contains_hash(&hash) {
2431            Ok(true)
2432        } else {
2433            self.provider.is_known(hash)
2434        }
2435    }
2436
2437    /// Return sealed block header from in-memory state or database by hash.
2438    fn sealed_header_by_hash(
2439        &self,
2440        hash: B256,
2441    ) -> ProviderResult<Option<SealedHeader<N::BlockHeader>>> {
2442        // check memory first
2443        let header = self.state.tree_state.sealed_header_by_hash(&hash);
2444
2445        if header.is_some() {
2446            Ok(header)
2447        } else {
2448            self.provider.sealed_header_by_hash(hash)
2449        }
2450    }
2451
2452    /// Return the parent hash of the lowest buffered ancestor for the requested block, if there
2453    /// are any buffered ancestors. If there are no buffered ancestors, and the block itself does
2454    /// not exist in the buffer, this returns the hash that is passed in.
2455    ///
2456    /// Returns the parent hash of the block itself if the block is buffered and has no other
2457    /// buffered ancestors.
2458    fn lowest_buffered_ancestor_or(&self, hash: B256) -> B256 {
2459        self.state
2460            .buffer
2461            .lowest_ancestor(&hash)
2462            .map(|block| block.parent_hash())
2463            .unwrap_or_else(|| hash)
2464    }
2465
2466    /// Returns whether block downloads should also attempt to fetch the blocks' access lists.
2467    ///
2468    /// This is the case once Amsterdam is active at the current canonical head, indicated by the
2469    /// head header carrying a block access list hash.
2470    ///
2471    /// This is a coarse gate: a download issued while the head is still pre-Amsterdam won't ask
2472    /// for access lists of blocks past the fork. Access list downloads are best-effort anyway, so
2473    /// those blocks simply take the regular execution path.
2474    fn should_download_access_lists(&self) -> bool {
2475        self.canonical_in_memory_state.get_canonical_head().block_access_list_hash().is_some()
2476    }
2477
2478    /// If validation fails, the response MUST contain the latest valid hash:
2479    ///
2480    ///   - The block hash of the ancestor of the invalid payload satisfying the following two
2481    ///     conditions:
2482    ///     - It is fully validated and deemed VALID
2483    ///     - Any other ancestor of the invalid payload with a higher blockNumber is INVALID
2484    ///   - 0x0000000000000000000000000000000000000000000000000000000000000000 if the above
2485    ///     conditions are satisfied by a `PoW` block.
2486    ///   - null if client software cannot determine the ancestor of the invalid payload satisfying
2487    ///     the above conditions.
2488    fn latest_valid_hash_for_invalid_payload(
2489        &mut self,
2490        parent_hash: B256,
2491    ) -> ProviderResult<Option<B256>> {
2492        // Check if parent exists in side chain or in canonical chain.
2493        if self.has_block_by_hash(parent_hash)? {
2494            return Ok(Some(parent_hash))
2495        }
2496
2497        // iterate over ancestors in the invalid cache
2498        // until we encounter the first valid ancestor
2499        let mut current_hash = parent_hash;
2500        let mut current_block = self.state.invalid_headers.get(&current_hash);
2501        while let Some(block_with_parent) = current_block {
2502            current_hash = block_with_parent.parent;
2503            current_block = self.state.invalid_headers.get(&current_hash);
2504
2505            // If current_header is None, then the current_hash does not have an invalid
2506            // ancestor in the cache, check its presence in blockchain tree
2507            if current_block.is_none() && self.has_block_by_hash(current_hash)? {
2508                return Ok(Some(current_hash))
2509            }
2510        }
2511        Ok(None)
2512    }
2513
2514    /// Prepares the invalid payload response for the given hash, checking the
2515    /// database for the parent hash and populating the payload status with the latest valid hash
2516    /// according to the engine api spec.
2517    fn prepare_invalid_response(&mut self, parent_hash: B256) -> ProviderResult<PayloadStatus> {
2518        let valid_parent_hash = match self.sealed_header_by_hash(parent_hash)? {
2519            // Edge case: the `latestValid` field is the zero hash if the parent block is the
2520            // terminal PoW block, which we need to identify by looking at the parent's block
2521            // difficulty
2522            Some(parent) if !parent.difficulty().is_zero() => Some(B256::ZERO),
2523            Some(_) => Some(parent_hash),
2524            None => self.latest_valid_hash_for_invalid_payload(parent_hash)?,
2525        };
2526
2527        Ok(PayloadStatus::from_status(PayloadStatusEnum::Invalid {
2528            validation_error: PayloadValidationError::LinksToRejectedPayload.to_string(),
2529        })
2530        .with_latest_valid_hash(valid_parent_hash.unwrap_or_default()))
2531    }
2532
2533    /// Returns true if the given hash is the last received sync target block.
2534    ///
2535    /// See [`ForkchoiceStateTracker::sync_target_state`]
2536    fn is_sync_target_head(&self, block_hash: B256) -> bool {
2537        if let Some(target) = self.state.forkchoice_state_tracker.sync_target_state() {
2538            return target.head_block_hash == block_hash
2539        }
2540        false
2541    }
2542
2543    /// Returns true if the given hash is part of the last received sync target fork choice update.
2544    ///
2545    /// See [`ForkchoiceStateTracker::sync_target_state`]
2546    fn is_any_sync_target(&self, block_hash: B256) -> bool {
2547        if let Some(target) = self.state.forkchoice_state_tracker.sync_target_state() {
2548            return target.contains(block_hash)
2549        }
2550        false
2551    }
2552
2553    /// Checks if the given `check` hash points to an invalid header, inserting the given `head`
2554    /// block into the invalid header cache if the `check` hash has a known invalid ancestor.
2555    ///
2556    /// Returns a payload status response according to the engine API spec if the block is known to
2557    /// be invalid.
2558    fn check_invalid_ancestor_with_head(
2559        &mut self,
2560        check: B256,
2561        head: &SealedBlock<N::Block>,
2562    ) -> ProviderResult<Option<PayloadStatus>> {
2563        // check if the check hash was previously marked as invalid
2564        let Some(header) = self.state.invalid_headers.get(&check) else { return Ok(None) };
2565
2566        Ok(Some(self.on_invalid_new_payload(head.clone(), header)?))
2567    }
2568
2569    /// Invoked when a new payload received is invalid.
2570    fn on_invalid_new_payload(
2571        &mut self,
2572        head: SealedBlock<N::Block>,
2573        invalid: BlockWithParent,
2574    ) -> ProviderResult<PayloadStatus> {
2575        // populate the latest valid hash field
2576        let status = self.prepare_invalid_response(invalid.parent)?;
2577
2578        // insert the head block into the invalid header cache
2579        self.state.invalid_headers.insert_with_invalid_ancestor(head.hash(), invalid);
2580        self.emit_event(ConsensusEngineEvent::InvalidBlock {
2581            block: Box::new(head),
2582            error: PayloadValidationError::LinksToRejectedPayload.to_string(),
2583        });
2584
2585        Ok(status)
2586    }
2587
2588    /// Finds any invalid ancestor for the given payload.
2589    ///
2590    /// This function first checks if the block itself is in the invalid headers cache (to
2591    /// avoid re-executing a known-invalid block). Then it walks up the chain of buffered
2592    /// ancestors and checks if any ancestor is marked as invalid.
2593    ///
2594    /// The check works by:
2595    /// 1. Checking if the block hash itself is in the `invalid_headers` map
2596    /// 2. Finding the lowest buffered ancestor for the given block hash
2597    /// 3. If the ancestor is the same as the block hash itself, using the parent hash instead
2598    /// 4. Checking if this ancestor is in the `invalid_headers` map
2599    ///
2600    /// Returns the invalid ancestor block info if found, or None if no invalid ancestor exists.
2601    fn find_invalid_ancestor(&mut self, payload: &T::ExecutionData) -> Option<BlockWithParent> {
2602        let parent_hash = payload.parent_hash();
2603        let block_hash = payload.block_hash();
2604
2605        // Check if the block itself is already known to be invalid, avoiding re-execution
2606        if let Some(entry) = self.state.invalid_headers.get(&block_hash) {
2607            return Some(entry);
2608        }
2609
2610        let mut lowest_buffered_ancestor = self.lowest_buffered_ancestor_or(block_hash);
2611        if lowest_buffered_ancestor == block_hash {
2612            lowest_buffered_ancestor = parent_hash;
2613        }
2614
2615        // Check if the block has an invalid ancestor
2616        self.state.invalid_headers.get(&lowest_buffered_ancestor)
2617    }
2618
2619    /// Handles a payload that has an invalid ancestor.
2620    ///
2621    /// This function validates the payload and processes it according to whether it's
2622    /// well-formed or malformed:
2623    /// 1. **Well-formed payload**: The payload is marked as invalid since it descends from a
2624    ///    known-bad block, which violates consensus rules
2625    /// 2. **Malformed payload**: Returns an appropriate error status since the payload cannot be
2626    ///    validated due to its own structural issues
2627    fn handle_invalid_ancestor_payload(
2628        &mut self,
2629        payload: T::ExecutionData,
2630        invalid: BlockWithParent,
2631    ) -> Result<PayloadStatus, InsertBlockFatalError> {
2632        let parent_hash = payload.parent_hash();
2633        let num_hash = payload.num_hash();
2634
2635        // Here we might have 2 cases
2636        // 1. the block is well formed and indeed links to an invalid header, meaning we should
2637        //    remember it as invalid
2638        // 2. the block is not well formed (i.e block hash is incorrect), and we should just return
2639        //    an error and forget it
2640        let block = match self.payload_validator.convert_payload_to_block(payload) {
2641            Ok(block) => block,
2642            Err(error) => return Ok(self.on_new_payload_error(error, num_hash, parent_hash)?),
2643        };
2644
2645        Ok(self.on_invalid_new_payload(block, invalid)?)
2646    }
2647
2648    /// Checks if the given `head` points to an invalid header, which requires a specific response
2649    /// to a forkchoice update.
2650    fn check_invalid_ancestor(&mut self, head: B256) -> ProviderResult<Option<PayloadStatus>> {
2651        // check if the head was previously marked as invalid
2652        let Some(header) = self.state.invalid_headers.get(&head) else { return Ok(None) };
2653
2654        // Try to prepare invalid response, but handle errors gracefully
2655        match self.prepare_invalid_response(header.parent) {
2656            Ok(status) => Ok(Some(status)),
2657            Err(err) => {
2658                debug!(target: "engine::tree", %err, "Failed to prepare invalid response for ancestor check");
2659                // Return a basic invalid status without latest valid hash
2660                Ok(Some(PayloadStatus::from_status(PayloadStatusEnum::Invalid {
2661                    validation_error: PayloadValidationError::LinksToRejectedPayload.to_string(),
2662                })))
2663            }
2664        }
2665    }
2666
2667    /// Validate if block is correct and satisfies all the consensus rules that concern the header
2668    /// and block body itself.
2669    fn validate_block(&self, block: &SealedBlock<N::Block>) -> Result<(), ConsensusError> {
2670        if let Err(e) = self.consensus.validate_header(block.sealed_header()) {
2671            error!(target: "engine::tree", ?block, "Failed to validate header {}: {e}", block.hash());
2672            return Err(e)
2673        }
2674
2675        if let Err(e) = self.consensus.validate_block_pre_execution(block) {
2676            error!(target: "engine::tree", ?block, "Failed to validate block {}: {e}", block.hash());
2677            return Err(e)
2678        }
2679
2680        Ok(())
2681    }
2682
2683    /// Attempts to connect any buffered blocks that are connected to the given parent hash.
2684    #[instrument(level = "debug", target = "engine::tree", skip(self))]
2685    fn try_connect_buffered_blocks(
2686        &mut self,
2687        parent: BlockNumHash,
2688    ) -> Result<(), InsertBlockFatalError> {
2689        let blocks = self.state.buffer.remove_block_with_children(&parent.hash);
2690
2691        if blocks.is_empty() {
2692            // nothing to append
2693            return Ok(())
2694        }
2695
2696        let now = Instant::now();
2697        let block_count = blocks.len();
2698        for child in blocks {
2699            let child_num_hash = child.num_hash();
2700            match self.insert_block(child) {
2701                Ok(res) => {
2702                    debug!(target: "engine::tree", child =?child_num_hash, ?res, "connected buffered block");
2703                    if self.is_any_sync_target(child_num_hash.hash) &&
2704                        matches!(res, InsertPayloadOk::Inserted(BlockStatus::Valid))
2705                    {
2706                        debug!(target: "engine::tree", child =?child_num_hash, "connected sync target block");
2707                        // we just inserted a block that we know is part of the canonical chain, so
2708                        // we can make it canonical
2709                        self.make_canonical(child_num_hash.hash)?;
2710                    }
2711                }
2712                Err(err) => {
2713                    if let InsertPayloadError::Block(err) = err {
2714                        debug!(target: "engine::tree", ?err, "failed to connect buffered block to tree");
2715                        if let Err(InsertBlockProcessingError::Fatal(fatal)) =
2716                            self.on_insert_block_error(err)
2717                        {
2718                            warn!(target: "engine::tree", %fatal, "fatal error occurred while connecting buffered blocks");
2719                        }
2720                    }
2721                }
2722            }
2723        }
2724
2725        debug!(target: "engine::tree", elapsed = ?now.elapsed(), %block_count, "connected buffered blocks");
2726        Ok(())
2727    }
2728
2729    /// Pre-validates the block and inserts it into the buffer.
2730    fn buffer_block(
2731        &mut self,
2732        block: SealedBlock<N::Block>,
2733    ) -> Result<(), InsertBlockError<N::Block>> {
2734        if let Err(err) = self.validate_block(&block) {
2735            return Err(InsertBlockError::consensus_error(err, block))
2736        }
2737        self.state.buffer.insert_block(block.into());
2738        Ok(())
2739    }
2740
2741    /// Returns true if the distance from the local tip to the block is greater than the configured
2742    /// threshold.
2743    ///
2744    /// If the `local_tip` is greater than the `block`, then this will return false.
2745    #[inline]
2746    const fn exceeds_backfill_run_threshold(&self, local_tip: u64, block: u64) -> bool {
2747        block > local_tip && block - local_tip > self.config.backfill_run_threshold()
2748    }
2749
2750    /// Returns how far the local tip is from the given block. If the local tip is at the same
2751    /// height or its block number is greater than the given block, this returns None.
2752    #[inline]
2753    const fn distance_from_local_tip(&self, local_tip: u64, block: u64) -> Option<u64> {
2754        if block > local_tip {
2755            Some(block - local_tip)
2756        } else {
2757            None
2758        }
2759    }
2760
2761    /// Returns the block hash that backfill should target.
2762    ///
2763    /// Defaults to the finalized block hash. On OP Stack, the CL finalizes in large batches and the
2764    /// finalized hash can lag the canonical tip by a wide margin, so backfill targets the head.
2765    ///
2766    /// The zero-finalized optimistic-sync fallback for non-OP Stack chains is handled by
2767    /// [`Self::backfill_sync_target`].
2768    const fn backfill_target_hash(&self, state: ForkchoiceState) -> B256 {
2769        if self.engine_kind.is_opstack() {
2770            state.head_block_hash
2771        } else {
2772            state.finalized_block_hash
2773        }
2774    }
2775
2776    /// Returns the target hash to sync to if the distance from the local tip is greater than the
2777    /// threshold and we're not yet synced to the backfill target (see
2778    /// [`Self::backfill_target_hash`]).
2779    ///
2780    /// If this is invoked after a new block has been downloaded, the downloaded block could be
2781    /// the (missing) target block.
2782    fn backfill_sync_target(
2783        &self,
2784        canonical_tip_num: u64,
2785        target_block_number: u64,
2786        downloaded_block: Option<BlockNumHash>,
2787    ) -> Option<B256> {
2788        let state = self.state.forkchoice_state_tracker.sync_target_state()?;
2789        let target_hash = self.backfill_target_hash(state);
2790
2791        // check if the downloaded block is the tracked backfill target
2792        let exceeds_backfill_threshold = match downloaded_block.as_ref() {
2793            // if we downloaded the target block we can now check how far we're off
2794            Some(downloaded_block) if downloaded_block.hash == target_hash => {
2795                self.exceeds_backfill_run_threshold(canonical_tip_num, downloaded_block.number)
2796            }
2797            _ => match self.state.buffer.block(&target_hash) {
2798                // if we have buffered the target block, we should check how far we're off
2799                Some(buffered_target) => {
2800                    self.exceeds_backfill_run_threshold(canonical_tip_num, buffered_target.number())
2801                }
2802                // check if the distance exceeds the threshold for backfill sync
2803                None => self.exceeds_backfill_run_threshold(canonical_tip_num, target_block_number),
2804            },
2805        };
2806
2807        if !exceeds_backfill_threshold {
2808            return None
2809        }
2810
2811        // if we have already canonicalized the target block, we should skip backfill
2812        match self.provider.header_by_hash_or_number(target_hash.into()) {
2813            Err(err) => {
2814                warn!(target: "engine::tree", %err, "Failed to get backfill target block header");
2815                None
2816            }
2817            // we don't have the block yet and the distance exceeds the allowed threshold
2818            Ok(None) if !target_hash.is_zero() => Some(target_hash),
2819            Ok(None) => {
2820                // OPTIMISTIC SYNCING
2821                //
2822                // It can happen when the node is doing an
2823                // optimistic sync, where the CL has no knowledge of the finalized hash,
2824                // but is expecting the EL to sync as high
2825                // as possible before finalizing.
2826                //
2827                // This usually doesn't happen on ETH mainnet since CLs use the more
2828                // secure checkpoint syncing.
2829                //
2830                // However, optimism chains will do this. The risk of a reorg is however
2831                // low.
2832                debug!(target: "engine::tree", hash=?state.head_block_hash, "Setting head hash as an optimistic backfill target.");
2833                Some(state.head_block_hash)
2834            }
2835            // we're fully synced to the target block
2836            Ok(Some(_)) => None,
2837        }
2838    }
2839
2840    /// This method tries to detect whether on-disk and in-memory states have diverged. It might
2841    /// happen if a reorg is happening while we are persisting a block.
2842    fn find_disk_reorg(&self) -> ProviderResult<Option<u64>> {
2843        let mut canonical = self.state.tree_state.current_canonical_head;
2844        let mut persisted = self.persistence_state.last_persisted_block;
2845
2846        let parent_num_hash = |num_hash: NumHash| -> ProviderResult<NumHash> {
2847            Ok(self
2848                .sealed_header_by_hash(num_hash.hash)?
2849                .ok_or(ProviderError::BlockHashNotFound(num_hash.hash))?
2850                .parent_num_hash())
2851        };
2852
2853        // Happy path, canonical chain is ahead or equal to persisted chain.
2854        // Walk canonical chain back to make sure that it connects to persisted chain.
2855        while canonical.number > persisted.number {
2856            canonical = parent_num_hash(canonical)?;
2857        }
2858
2859        // If we've reached persisted tip by walking the canonical chain back, everything is fine.
2860        if canonical == persisted {
2861            return Ok(None);
2862        }
2863
2864        // At this point, we know that `persisted` block can't be reached by walking the canonical
2865        // chain back. In this case we need to truncate it to the first canonical block it connects
2866        // to.
2867
2868        // Firstly, walk back until we reach the same height as `canonical`.
2869        while persisted.number > canonical.number {
2870            persisted = parent_num_hash(persisted)?;
2871        }
2872
2873        debug_assert_eq!(persisted.number, canonical.number);
2874
2875        // Now walk both chains back until we find a common ancestor.
2876        while persisted.hash != canonical.hash {
2877            canonical = parent_num_hash(canonical)?;
2878            persisted = parent_num_hash(persisted)?;
2879        }
2880
2881        debug!(target: "engine::tree", remove_above=persisted.number, "on-disk reorg detected");
2882
2883        Ok(Some(persisted.number))
2884    }
2885
2886    /// Invoked when we the canonical chain has been updated.
2887    ///
2888    /// This is invoked on a valid forkchoice update, or if we can make the target block canonical.
2889    fn on_canonical_chain_update(&mut self, chain_update: NewCanonicalChain<N>) {
2890        trace!(target: "engine::tree", new_blocks = %chain_update.new_block_count(), reorged_blocks =  %chain_update.reorged_block_count(), "applying new chain update");
2891        let start = Instant::now();
2892
2893        // update the tracked canonical head
2894        self.state.tree_state.set_canonical_head(chain_update.tip().num_hash());
2895
2896        let tip = chain_update.tip().clone_sealed_header();
2897        let notification = chain_update.to_chain_notification();
2898
2899        // reinsert any missing reorged blocks
2900        if let NewCanonicalChain::Reorg { new, old } = &chain_update {
2901            let new_first = new.first().map(|first| first.recovered_block().num_hash());
2902            let old_first = old.first().map(|first| first.recovered_block().num_hash());
2903            trace!(target: "engine::tree", ?new_first, ?old_first, "Reorg detected, new and old first blocks");
2904
2905            self.state.set_pending_sparse_trie_prune(false);
2906            self.update_reorg_metrics(old.len(), old_first);
2907            self.reinsert_reorged_blocks(new.clone());
2908            self.reinsert_reorged_blocks(old.clone());
2909        }
2910
2911        // update the tracked in-memory state with the new chain
2912        self.canonical_in_memory_state.update_chain(chain_update);
2913        self.canonical_in_memory_state.set_canonical_head(tip.clone());
2914        self.payload_validator.on_canonical_head_changed(tip.hash(), &self.state);
2915
2916        // Update metrics based on new tip
2917        self.metrics.tree.canonical_chain_height.set(tip.number() as f64);
2918
2919        // sends an event to all active listeners about the new canonical chain
2920        self.canonical_in_memory_state.notify_canon_state(notification);
2921
2922        // emit event
2923        self.emit_event(ConsensusEngineEvent::CanonicalChainCommitted(
2924            Box::new(tip),
2925            start.elapsed(),
2926        ));
2927    }
2928
2929    /// This updates metrics based on the given reorg length and first reorged block number.
2930    fn update_reorg_metrics(&self, old_chain_length: usize, first_reorged_block: Option<NumHash>) {
2931        if let Some(first_reorged_block) = first_reorged_block.map(|block| block.number) {
2932            if let Some(finalized) = self.canonical_in_memory_state.get_finalized_num_hash() &&
2933                first_reorged_block <= finalized.number
2934            {
2935                self.metrics.tree.reorgs.finalized.increment(1);
2936            } else if let Some(safe) = self.canonical_in_memory_state.get_safe_num_hash() &&
2937                first_reorged_block <= safe.number
2938            {
2939                self.metrics.tree.reorgs.safe.increment(1);
2940            } else {
2941                self.metrics.tree.reorgs.head.increment(1);
2942            }
2943        } else {
2944            debug_unreachable!("Reorged chain doesn't have any blocks");
2945        }
2946        self.metrics.tree.latest_reorg_depth.set(old_chain_length as f64);
2947    }
2948
2949    /// This reinserts any blocks in the new chain that do not already exist in the tree
2950    fn reinsert_reorged_blocks(&mut self, new_chain: Vec<ExecutedBlock<N>>) {
2951        for block in new_chain {
2952            if self
2953                .state
2954                .tree_state
2955                .executed_block_by_hash(block.recovered_block().hash())
2956                .is_none()
2957            {
2958                trace!(target: "engine::tree", num=?block.recovered_block().number(), hash=?block.recovered_block().hash(), "Reinserting block into tree state");
2959                self.state.tree_state.insert_executed(block);
2960            }
2961        }
2962        self.metrics.engine.executed_blocks.set(self.state.tree_state.block_count() as f64);
2963    }
2964
2965    /// This handles downloaded blocks that are shown to be disconnected from the canonical chain.
2966    ///
2967    /// This mainly compares the missing parent of the downloaded block with the current canonical
2968    /// tip, and decides whether or not backfill sync should be triggered.
2969    fn on_disconnected_downloaded_block(
2970        &self,
2971        downloaded_block: BlockNumHash,
2972        missing_parent: BlockNumHash,
2973        head: BlockNumHash,
2974    ) -> Option<TreeEvent> {
2975        // compare the missing parent with the canonical tip
2976        if let Some(target) =
2977            self.backfill_sync_target(head.number, missing_parent.number, Some(downloaded_block))
2978        {
2979            trace!(target: "engine::tree", %target, "triggering backfill on downloaded block");
2980            return Some(TreeEvent::BackfillAction(BackfillAction::Start(target.into())));
2981        }
2982
2983        // continue downloading the missing parent
2984        //
2985        // this happens if either:
2986        //  * the missing parent block num < canonical tip num
2987        //    * this case represents a missing block on a fork that is shorter than the canonical
2988        //      chain
2989        //  * the missing parent block num >= canonical tip num, but the number of missing blocks is
2990        //    less than the backfill threshold
2991        //    * this case represents a potentially long range of blocks to download and execute
2992        let request = if let Some(distance) =
2993            self.distance_from_local_tip(head.number, missing_parent.number)
2994        {
2995            trace!(target: "engine::tree", %distance, missing=?missing_parent, "downloading missing parent block range");
2996            DownloadRequest::block_range(missing_parent.hash, distance)
2997        } else {
2998            trace!(target: "engine::tree", missing=?missing_parent, "downloading missing parent block");
2999            // This happens when the missing parent is on an outdated
3000            // sidechain and we can only download the missing block itself
3001            DownloadRequest::single_block(missing_parent.hash)
3002        };
3003
3004        Some(TreeEvent::Download(request.with_access_lists(self.should_download_access_lists())))
3005    }
3006
3007    /// Handles a downloaded block that was successfully inserted as valid.
3008    ///
3009    /// If the block matches the sync target head, returns [`TreeAction::MakeCanonical`].
3010    /// If it matches a non-head sync target (safe or finalized), makes it canonical inline
3011    /// and triggers a download for the remaining blocks towards the actual head.
3012    /// Otherwise, tries to connect buffered blocks.
3013    fn on_valid_downloaded_block(
3014        &mut self,
3015        block_num_hash: BlockNumHash,
3016    ) -> Result<Option<TreeEvent>, InsertBlockFatalError> {
3017        // check if we just inserted a block that's part of sync targets,
3018        // i.e. head, safe, or finalized
3019        if let Some(sync_target) = self.state.forkchoice_state_tracker.sync_target_state() &&
3020            sync_target.contains(block_num_hash.hash)
3021        {
3022            debug!(target: "engine::tree", ?sync_target, "appended downloaded sync target block");
3023
3024            if sync_target.head_block_hash == block_num_hash.hash {
3025                // we just inserted the sync target head block, make it canonical
3026                return Ok(Some(TreeEvent::TreeAction(TreeAction::MakeCanonical {
3027                    sync_target_head: block_num_hash.hash,
3028                })))
3029            }
3030
3031            // This block is part of the sync target (safe or finalized) but not the
3032            // head. Make it canonical and try to connect any buffered children, then
3033            // continue downloading towards the actual head if needed.
3034            self.make_canonical(block_num_hash.hash)?;
3035            self.try_connect_buffered_blocks(block_num_hash)?;
3036
3037            // Check if we've reached the sync target head after connecting buffered
3038            // blocks (e.g. the head block may have already been buffered).
3039            if self.state.tree_state.canonical_block_hash() != sync_target.head_block_hash {
3040                let target = self.lowest_buffered_ancestor_or(sync_target.head_block_hash);
3041                trace!(target: "engine::tree", %target, "sync target head not yet reached, downloading head block");
3042                return Ok(Some(TreeEvent::Download(
3043                    DownloadRequest::single_block(target)
3044                        .with_access_lists(self.should_download_access_lists()),
3045                )))
3046            }
3047
3048            return Ok(None)
3049        }
3050        trace!(target: "engine::tree", "appended downloaded block");
3051        self.try_connect_buffered_blocks(block_num_hash)?;
3052        Ok(None)
3053    }
3054
3055    /// Invoked with a block downloaded from the network
3056    ///
3057    /// Returns an event with the appropriate action to take, such as:
3058    ///  - download more missing blocks
3059    ///  - try to canonicalize the target if the `block` is the tracked target (head) block.
3060    #[instrument(level = "debug", target = "engine::tree", skip_all, fields(block_hash = %block.hash(), block_num = %block.number()))]
3061    fn on_downloaded_block(
3062        &mut self,
3063        block: SealedBlockWithAccessList<N::Block>,
3064    ) -> Result<Option<TreeEvent>, InsertBlockFatalError> {
3065        let block_num_hash = block.num_hash();
3066        let lowest_buffered_ancestor = self.lowest_buffered_ancestor_or(block_num_hash.hash);
3067        if self.check_invalid_ancestor_with_head(lowest_buffered_ancestor, &block)?.is_some() {
3068            return Ok(None)
3069        }
3070
3071        if !self.backfill_sync_state.is_idle() {
3072            return Ok(None)
3073        }
3074
3075        // try to append the block
3076        match self.insert_block(block) {
3077            Ok(InsertPayloadOk::Inserted(BlockStatus::Valid)) => {
3078                return self.on_valid_downloaded_block(block_num_hash);
3079            }
3080            Ok(InsertPayloadOk::Inserted(BlockStatus::Disconnected { head, missing_ancestor })) => {
3081                // block is not connected to the canonical head, we need to download
3082                // its missing branch first
3083                return Ok(self.on_disconnected_downloaded_block(
3084                    block_num_hash,
3085                    missing_ancestor,
3086                    head,
3087                ))
3088            }
3089            Ok(InsertPayloadOk::AlreadySeen(_)) => {
3090                trace!(target: "engine::tree", "downloaded block already executed");
3091            }
3092            Err(err) => {
3093                if let InsertPayloadError::Block(err) = err {
3094                    debug!(target: "engine::tree", err=%err.kind(), "failed to insert downloaded block");
3095                    if let Err(InsertBlockProcessingError::Fatal(fatal)) =
3096                        self.on_insert_block_error(err)
3097                    {
3098                        warn!(target: "engine::tree", %fatal, "fatal error occurred while inserting downloaded block");
3099                    }
3100                }
3101            }
3102        }
3103        Ok(None)
3104    }
3105
3106    /// Inserts a payload into the tree and executes it.
3107    ///
3108    /// This function validates the payload's basic structure, then executes it using the
3109    /// payload validator. The execution includes running all transactions in the payload
3110    /// and validating the resulting state transitions.
3111    ///
3112    /// Returns `InsertPayloadOk` if the payload was successfully inserted and executed,
3113    /// or `InsertPayloadError` if validation or execution failed.
3114    fn insert_payload(
3115        &mut self,
3116        payload: T::ExecutionData,
3117    ) -> Result<InsertPayloadOk, InsertPayloadError<N::Block>> {
3118        self.insert_block_or_payload(
3119            payload.block_with_parent(),
3120            payload,
3121            |validator, payload, ctx| validator.validate_payload(payload, ctx),
3122            |this, payload| Ok(this.payload_validator.convert_payload_to_block(payload)?.into()),
3123        )
3124    }
3125
3126    fn insert_block(
3127        &mut self,
3128        block: SealedBlockWithAccessList<N::Block>,
3129    ) -> Result<InsertPayloadOk, InsertPayloadError<N::Block>> {
3130        self.insert_block_or_payload(
3131            block.block_with_parent(),
3132            block,
3133            |validator, block, ctx| validator.validate_block(block, ctx),
3134            |_, block| Ok(block),
3135        )
3136    }
3137
3138    /// Inserts a block or payload into the blockchain tree with full execution.
3139    ///
3140    /// This is a generic function that handles both blocks and payloads by accepting
3141    /// a block identifier, input data, and execution/validation functions. It performs
3142    /// comprehensive checks and execution:
3143    ///
3144    /// - Validates that the block doesn't already exist in the tree
3145    /// - Ensures parent state is available, buffering if necessary
3146    /// - Executes the block/payload using the provided execute function
3147    /// - Handles both canonical and fork chain insertions
3148    /// - Updates pending block state when appropriate
3149    /// - Emits consensus engine events and records metrics
3150    ///
3151    /// Returns `InsertPayloadOk::Inserted(BlockStatus::Valid)` on successful execution,
3152    /// `InsertPayloadOk::AlreadySeen` if the block already exists, or
3153    /// `InsertPayloadOk::Inserted(BlockStatus::Disconnected)` if parent state is missing.
3154    #[instrument(level = "debug", target = "engine::tree", skip_all, fields(?block_id))]
3155    fn insert_block_or_payload<Input, Err>(
3156        &mut self,
3157        block_id: BlockWithParent,
3158        input: Input,
3159        execute: impl FnOnce(&mut V, Input, TreeCtx<'_, N>) -> Result<ValidationOutput<N>, Err>,
3160        convert_to_block: impl FnOnce(
3161            &mut Self,
3162            Input,
3163        ) -> Result<SealedBlockWithAccessList<N::Block>, Err>,
3164    ) -> Result<InsertPayloadOk, Err>
3165    where
3166        Err: From<InsertBlockError<N::Block>>,
3167    {
3168        let block_insert_start = Instant::now();
3169        let block_num_hash = block_id.block;
3170        debug!(target: "engine::tree", block=?block_num_hash, parent = ?block_id.parent, "Inserting new block into tree");
3171
3172        // Check if block already exists - first in memory, then DB only if it could be persisted
3173        if self.state.tree_state.contains_hash(&block_num_hash.hash) {
3174            convert_to_block(self, input)?;
3175            return Ok(InsertPayloadOk::AlreadySeen(BlockStatus::Valid));
3176        }
3177
3178        // Only query DB if block could be persisted (number <= last persisted block).
3179        // New blocks from CL always have number > last persisted, so skip DB lookup for them.
3180        if block_num_hash.number <= self.persistence_state.last_persisted_block.number {
3181            match self.provider.sealed_header_by_hash(block_num_hash.hash) {
3182                Err(err) => {
3183                    let block = convert_to_block(self, input)?;
3184                    return Err(InsertBlockError::new(block.split().0, err.into()).into());
3185                }
3186                Ok(Some(_)) => {
3187                    convert_to_block(self, input)?;
3188                    return Ok(InsertPayloadOk::AlreadySeen(BlockStatus::Valid));
3189                }
3190                Ok(None) => {}
3191            }
3192        }
3193
3194        // Ensure that the parent state is available.
3195        if !self.state.tree_state.contains_hash(&block_id.parent) {
3196            let parent_exists = match self.provider.header(block_id.parent) {
3197                Ok(header) => header.is_some(),
3198                Err(err) => {
3199                    let block = convert_to_block(self, input)?;
3200                    return Err(InsertBlockError::new(block.split().0, err.into()).into());
3201                }
3202            };
3203
3204            if !parent_exists {
3205                let block = convert_to_block(self, input)?;
3206                // we don't have the state required to execute this block, buffering it and find the
3207                // missing parent block
3208                let missing_ancestor = self
3209                    .state
3210                    .buffer
3211                    .lowest_ancestor(&block.parent_hash())
3212                    .map(|block| block.parent_num_hash())
3213                    .unwrap_or_else(|| block.parent_num_hash());
3214
3215                self.state.buffer.insert_block(block);
3216
3217                return Ok(InsertPayloadOk::Inserted(BlockStatus::Disconnected {
3218                    head: self.state.tree_state.current_canonical_head,
3219                    missing_ancestor,
3220                }))
3221            }
3222        }
3223
3224        // determine whether we are on a fork chain by comparing the block number with the
3225        // canonical head. This is a simple check that is sufficient for the event emission below.
3226        // A block is considered a fork if its number is less than or equal to the canonical head,
3227        // as this indicates there's already a canonical block at that height.
3228        let is_fork = block_id.block.number <= self.state.tree_state.current_canonical_head.number;
3229
3230        let ctx = TreeCtx::new(&mut self.state, &self.canonical_in_memory_state);
3231
3232        let start = Instant::now();
3233
3234        let ValidationOutput { executed_block: executed, execution_timing_stats: timing_stats } =
3235            execute(&mut self.payload_validator, input, ctx)?;
3236
3237        if let Some(raw_bal) = executed.bal().map(|bal| bal.as_raw_bal().clone()) {
3238            let num_hash = executed.recovered_block().num_hash();
3239            if let Err(err) = self.provider.bal_store().insert(num_hash, raw_bal) {
3240                warn!(
3241                    target: "engine::tree",
3242                    ?num_hash,
3243                    %err,
3244                    "Failed to store validated block access list"
3245                );
3246            }
3247        }
3248
3249        // Emit slow block event immediately after execution so it appears even when
3250        // persistence hasn't completed yet (e.g. blocks arriving faster than persistence).
3251        if let Some(stats) = timing_stats {
3252            if let Some(threshold) = self.config.slow_block_threshold() {
3253                let total_duration = stats.execution_duration + stats.state_hash_duration;
3254                if total_duration > threshold {
3255                    self.emit_event(ConsensusEngineEvent::SlowBlock(SlowBlockInfo {
3256                        stats: stats.clone(),
3257                        commit_duration: None,
3258                        total_duration,
3259                    }));
3260                }
3261            }
3262            self.execution_timing_stats.insert(executed.recovered_block().hash(), stats);
3263        }
3264
3265        let is_pending = self.state.tree_state.canonical_block_hash() ==
3266            executed.recovered_block().parent_hash();
3267        self.state.tree_state.insert_executed(executed.clone());
3268
3269        if is_pending {
3270            debug!(target: "engine::tree", pending=?block_num_hash, "updating pending block");
3271            self.canonical_in_memory_state.set_pending_block(executed.clone());
3272        }
3273
3274        self.metrics.engine.executed_blocks.set(self.state.tree_state.block_count() as f64);
3275
3276        // emit insert event
3277        let elapsed = start.elapsed();
3278        let engine_event = if is_fork {
3279            ConsensusEngineEvent::ForkBlockAdded(executed, elapsed)
3280        } else {
3281            ConsensusEngineEvent::CanonicalBlockAdded(executed, elapsed)
3282        };
3283        self.emit_event(EngineApiEvent::BeaconConsensus(engine_event));
3284
3285        self.metrics
3286            .engine
3287            .block_insert_total_duration
3288            .record(block_insert_start.elapsed().as_secs_f64());
3289        debug!(target: "engine::tree", block=?block_num_hash, "Finished inserting block");
3290        Ok(InsertPayloadOk::Inserted(BlockStatus::Valid))
3291    }
3292
3293    /// Handles an error that occurred while inserting a block.
3294    ///
3295    /// If this is a validation error this will mark the block as invalid.
3296    ///
3297    /// Returns the proper payload status response if the block is invalid.
3298    fn on_insert_block_error(
3299        &mut self,
3300        error: InsertBlockError<N::Block>,
3301    ) -> Result<PayloadStatus, InsertBlockProcessingError> {
3302        let (block, error) = error.split();
3303
3304        let validation_err = error.ensure_validation_error()?;
3305
3306        // If the error was due to an invalid payload, the payload is added to the
3307        // invalid headers cache and `Ok` with [PayloadStatusEnum::Invalid] is
3308        // returned.
3309        warn!(
3310            target: "engine::tree",
3311            invalid_hash=%block.hash(),
3312            invalid_number=block.number(),
3313            %validation_err,
3314            "Invalid block error on new payload",
3315        );
3316        // The Amsterdam Engine API requires `latestValidHash: null` for an undecodable BAL.
3317        // <https://github.com/ethereum/execution-apis/blob/df75e230befef0de56ee8833322ed714bacb479c/src/engine/amsterdam.md?plain=1#L129>
3318        let latest_valid_hash =
3319            if matches!(&validation_err, InsertBlockValidationError::BlockAccessListDecode(_)) {
3320                None
3321            } else {
3322                self.latest_valid_hash_for_invalid_payload(block.parent_hash())
3323                    .map_err(InsertBlockFatalError::from)?
3324            };
3325
3326        // keep track of the invalid header unless the consensus impl considers it transient
3327        let is_transient = match &validation_err {
3328            InsertBlockValidationError::Consensus(err) => self.consensus.is_transient_error(err),
3329            _ => false,
3330        };
3331        if is_transient {
3332            warn!(
3333                target: "engine::tree",
3334                invalid_hash=%block.hash(),
3335                invalid_number=block.number(),
3336                %validation_err,
3337                "Skipping invalid header cache insert for transient validation error",
3338            );
3339        } else {
3340            self.state.invalid_headers.insert(block.block_with_parent());
3341        }
3342        self.emit_event(EngineApiEvent::BeaconConsensus(ConsensusEngineEvent::InvalidBlock {
3343            block: Box::new(block),
3344            error: validation_err.to_string(),
3345        }));
3346
3347        Ok(PayloadStatus::new(
3348            PayloadStatusEnum::Invalid { validation_error: validation_err.to_string() },
3349            latest_valid_hash,
3350        ))
3351    }
3352
3353    /// Handles a [`NewPayloadError`] by converting it to a [`PayloadStatus`].
3354    fn on_new_payload_error(
3355        &mut self,
3356        error: NewPayloadError,
3357        payload_num_hash: NumHash,
3358        parent_hash: B256,
3359    ) -> ProviderResult<PayloadStatus> {
3360        error!(target: "engine::tree", payload=?payload_num_hash, %error, "Invalid payload");
3361        // we need to convert the error to a payload status (response to the CL)
3362
3363        let latest_valid_hash =
3364            if error.is_block_hash_mismatch() || error.is_invalid_versioned_hashes() {
3365                // Engine-API rules:
3366                // > `latestValidHash: null` if the blockHash validation has failed (<https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/shanghai.md?plain=1#L113>)
3367                // > `latestValidHash: null` if the expected and the actual arrays don't match (<https://github.com/ethereum/execution-apis/blob/fe8e13c288c592ec154ce25c534e26cb7ce0530d/src/engine/cancun.md?plain=1#L103>)
3368                None
3369            } else {
3370                self.latest_valid_hash_for_invalid_payload(parent_hash)?
3371            };
3372
3373        let status = PayloadStatusEnum::from(error);
3374        Ok(PayloadStatus::new(status, latest_valid_hash))
3375    }
3376
3377    /// Attempts to find the header for the given block hash if it is canonical.
3378    ///
3379    /// A persisted header can still be found by hash while its disk reorg is pending, so it is
3380    /// only canonical if it is not above the canonical head and matches the canonical hash at its
3381    /// height. The in-memory canonical hash is preferred over the persisted one.
3382    pub fn find_canonical_header(
3383        &self,
3384        hash: B256,
3385    ) -> Result<Option<SealedHeader<N::BlockHeader>>, ProviderError> {
3386        if let Some(header) = self.canonical_in_memory_state.header_by_hash(hash) {
3387            return Ok(Some(header))
3388        }
3389
3390        let Some(header) = self.provider.sealed_header_by_hash(hash)? else { return Ok(None) };
3391        // Reorged-out headers remain on disk until persistence catches up, so reject old tips above
3392        // the in-memory head and use in-memory hashes to reject old blocks at shared heights.
3393        let number = header.number();
3394        if number > self.canonical_in_memory_state.get_canonical_block_number() {
3395            return Ok(None)
3396        }
3397        let canonical_hash =
3398            if let Some(hash) = self.canonical_in_memory_state.hash_by_number(number) {
3399                Some(hash)
3400            } else {
3401                self.provider.block_hash(number)?
3402            };
3403
3404        Ok((canonical_hash == Some(hash)).then_some(header))
3405    }
3406
3407    /// Checks that nonzero safe and finalized hashes belong to the chain defined by the FCU head.
3408    ///
3409    /// The [Engine API forkchoiceUpdated specification] requires `-38002` when a `VALID` head's
3410    /// safe or finalized hash is outside that chain, and requires all forkchoice state updates to
3411    /// be atomic. In direct FCU processing, this check must therefore run before canonicalizing
3412    /// the head or updating either marker; checking only after canonicalization would leave an
3413    /// invalid reorg applied.
3414    ///
3415    /// `chain_update` describes a proposed commit or reorg that has not yet been applied. Its new
3416    /// blocks and the canonical prefix below its first block form the proposed chain. Without a
3417    /// chain update, the proposed head is already canonical, so only canonical blocks through its
3418    /// height are eligible. Canonical-prefix hashes are resolved via
3419    /// [`Self::find_canonical_header`], which rejects stale persisted headers whose disk reorg
3420    /// cleanup is pending. A zero safe or finalized hash leaves that marker unchanged.
3421    /// Returns `Ok(false)` for an unknown or off-chain hash and propagates provider errors.
3422    ///
3423    /// [Engine API forkchoiceUpdated specification]: https://github.com/ethereum/execution-apis/blob/main/src/engine/paris.md#specification-1
3424    fn is_consistent_forkchoice_state(
3425        &self,
3426        state: ForkchoiceState,
3427        chain_update: Option<&NewCanonicalChain<N>>,
3428    ) -> ProviderResult<bool> {
3429        let canonical_head_number = match chain_update {
3430            Some(chain_update) => {
3431                let new = chain_update.new_blocks();
3432                // Only the canonical prefix below the new branch remains on the proposed chain.
3433                new.first().expect("non empty chain").block_number() - 1
3434            }
3435            None => {
3436                let Some(head) = self.find_canonical_header(state.head_block_hash)? else {
3437                    return Ok(false)
3438                };
3439                head.number()
3440            }
3441        };
3442
3443        for hash in [state.finalized_block_hash, state.safe_block_hash] {
3444            if hash.is_zero() || chain_update.is_some_and(|update| update.contains(hash)) {
3445                continue
3446            }
3447            let Some(header) = self.find_canonical_header(hash)? else { return Ok(false) };
3448            if header.number() > canonical_head_number {
3449                return Ok(false)
3450            }
3451        }
3452        Ok(true)
3453    }
3454
3455    /// Updates the tracked finalized block if we have it.
3456    fn update_finalized_block(
3457        &self,
3458        finalized_block_hash: B256,
3459    ) -> Result<(), OnForkChoiceUpdated> {
3460        if finalized_block_hash.is_zero() {
3461            return Ok(())
3462        }
3463
3464        match self.find_canonical_header(finalized_block_hash) {
3465            Ok(None) => {
3466                debug!(target: "engine::tree", "Finalized block not found in canonical chain");
3467                // if the finalized block is not known, we can't update the finalized block
3468                return Err(OnForkChoiceUpdated::invalid_state())
3469            }
3470            Ok(Some(finalized)) => {
3471                if Some(finalized.num_hash()) !=
3472                    self.canonical_in_memory_state.get_finalized_num_hash()
3473                {
3474                    // we're also persisting the finalized block on disk so we can reload it on
3475                    // restart this is required by optimism which queries the finalized block: <https://github.com/ethereum-optimism/optimism/blob/c383eb880f307caa3ca41010ec10f30f08396b2e/op-node/rollup/sync/start.go#L65-L65>
3476                    let _ = self.persistence.save_finalized_block_number(finalized.number());
3477                    self.canonical_in_memory_state.set_finalized(finalized.clone());
3478                    // Update finalized block height metric
3479                    self.metrics.tree.finalized_block_height.set(finalized.number() as f64);
3480                }
3481            }
3482            Err(err) => {
3483                error!(target: "engine::tree", %err, "Failed to fetch finalized block header");
3484            }
3485        }
3486
3487        Ok(())
3488    }
3489
3490    /// Updates the tracked safe block if we have it
3491    fn update_safe_block(&self, safe_block_hash: B256) -> Result<(), OnForkChoiceUpdated> {
3492        if safe_block_hash.is_zero() {
3493            return Ok(())
3494        }
3495
3496        match self.find_canonical_header(safe_block_hash) {
3497            Ok(None) => {
3498                debug!(target: "engine::tree", "Safe block not found in canonical chain");
3499                // if the safe block is not known, we can't update the safe block
3500                return Err(OnForkChoiceUpdated::invalid_state())
3501            }
3502            Ok(Some(safe)) => {
3503                if Some(safe.num_hash()) != self.canonical_in_memory_state.get_safe_num_hash() {
3504                    // we're also persisting the safe block on disk so we can reload it on
3505                    // restart this is required by optimism which queries the safe block: <https://github.com/ethereum-optimism/optimism/blob/c383eb880f307caa3ca41010ec10f30f08396b2e/op-node/rollup/sync/start.go#L65-L65>
3506                    let _ = self.persistence.save_safe_block_number(safe.number());
3507                    self.canonical_in_memory_state.set_safe(safe.clone());
3508                    // Update safe block height metric
3509                    self.metrics.tree.safe_block_height.set(safe.number() as f64);
3510                }
3511            }
3512            Err(err) => {
3513                error!(target: "engine::tree", %err, "Failed to fetch safe block header");
3514            }
3515        }
3516
3517        Ok(())
3518    }
3519
3520    /// Ensures that the given forkchoice state is consistent, assuming the head block has been
3521    /// made canonical.
3522    ///
3523    /// If the forkchoice state is consistent, this will return Ok(()). Otherwise, this will
3524    /// return an instance of [`OnForkChoiceUpdated`] that is INVALID.
3525    ///
3526    /// This also updates the safe and finalized blocks in the [`CanonicalInMemoryState`], if they
3527    /// are consistent with the head block.
3528    fn ensure_consistent_forkchoice_state(
3529        &self,
3530        state: ForkchoiceState,
3531    ) -> Result<(), OnForkChoiceUpdated> {
3532        // Ensure that the finalized block, if not zero, is known and in the canonical chain
3533        // after the head block is canonicalized.
3534        //
3535        // This ensures that the finalized block is consistent with the head block, i.e. the
3536        // finalized block is an ancestor of the head block.
3537        self.update_finalized_block(state.finalized_block_hash)?;
3538
3539        // Also ensure that the safe block, if not zero, is known and in the canonical chain
3540        // after the head block is canonicalized.
3541        //
3542        // This ensures that the safe block is consistent with the head block, i.e. the safe
3543        // block is an ancestor of the head block.
3544        self.update_safe_block(state.safe_block_hash)
3545    }
3546
3547    /// Validates the payload attributes with respect to the header and fork choice state.
3548    ///
3549    /// This is called during `engine_forkchoiceUpdated` when the CL provides payload attributes,
3550    /// indicating it wants the EL to start building a new block.
3551    ///
3552    /// Runs [`PayloadValidator::validate_payload_attributes_against_header`](reth_engine_primitives::PayloadValidator::validate_payload_attributes_against_header) to ensure
3553    /// `payloadAttributes.timestamp > headBlock.timestamp` per the Engine API spec.
3554    ///
3555    /// If validation passes, sends the attributes to the payload builder to start a new
3556    /// payload job. If it fails, returns `INVALID_PAYLOAD_ATTRIBUTES` without rolling back
3557    /// the forkchoice update.
3558    ///
3559    /// Note: At this point, the fork choice update is considered to be VALID, however, we can still
3560    /// return an error if the payload attributes are invalid.
3561    fn process_payload_attributes(
3562        &mut self,
3563        attributes: T::PayloadAttributes,
3564        head: &N::BlockHeader,
3565        state: ForkchoiceState,
3566    ) -> OnForkChoiceUpdated {
3567        if let Err(err) =
3568            self.payload_validator.validate_payload_attributes_against_header(&attributes, head)
3569        {
3570            warn!(target: "engine::tree", %err, ?head, "Invalid payload attributes");
3571            return OnForkChoiceUpdated::invalid_payload_attributes()
3572        }
3573
3574        // 8. Client software MUST begin a payload build process building on top of
3575        //    forkchoiceState.headBlockHash and identified via buildProcessId value if
3576        //    payloadAttributes is not null and the forkchoice state has been updated successfully.
3577        //    The build process is specified in the Payload building section.
3578
3579        // Acquire this before preparing resources because state-root setup can already start
3580        // workers that need the current in-memory overlay.
3581        let payload_build = self.payload_builds.acquire();
3582
3583        let resources = self
3584            .payload_validator
3585            .payload_builder_resources(
3586                state.head_block_hash,
3587                head,
3588                attributes.timestamp(),
3589                &mut self.state,
3590            )
3591            .with_lease(PayloadBuilderLease::new(payload_build));
3592
3593        // send the payload to the builder and return the receiver for the pending payload
3594        // id, initiating payload job is handled asynchronously
3595        let pending_payload_id = self.payload_builder.send_new_payload(BuildNewPayload {
3596            parent_hash: state.head_block_hash,
3597            attributes,
3598            resources,
3599        });
3600
3601        // Client software MUST respond to this method call in the following way:
3602        // {
3603        //      payloadStatus: {
3604        //          status: VALID,
3605        //          latestValidHash: forkchoiceState.headBlockHash,
3606        //          validationError: null
3607        //      },
3608        //      payloadId: buildProcessId
3609        // }
3610        //
3611        // if the payload is deemed VALID and the build process has begun.
3612        OnForkChoiceUpdated::updated_with_pending_payload_id(
3613            PayloadStatus::new(PayloadStatusEnum::Valid, Some(state.head_block_hash)),
3614            pending_payload_id,
3615        )
3616    }
3617
3618    /// Remove all blocks up to __and including__ the given block number.
3619    ///
3620    /// If a finalized hash is provided, the only non-canonical blocks which will be removed are
3621    /// those which have a fork point at or below the finalized hash.
3622    ///
3623    /// Canonical blocks below the upper bound will still be removed.
3624    pub(crate) fn remove_before(
3625        &mut self,
3626        upper_bound: BlockNumHash,
3627        finalized_hash: Option<B256>,
3628    ) -> ProviderResult<()> {
3629        // first fetch the finalized block number and then call the remove_before method on
3630        // tree_state
3631        let num = if let Some(hash) = finalized_hash {
3632            self.provider.block_number(hash)?.map(|number| BlockNumHash { number, hash })
3633        } else {
3634            None
3635        };
3636
3637        self.state.tree_state.remove_until(
3638            upper_bound,
3639            self.persistence_state.last_persisted_block.hash,
3640            num,
3641        );
3642        self.metrics.engine.executed_blocks.set(self.state.tree_state.block_count() as f64);
3643        Ok(())
3644    }
3645}
3646
3647/// Events received in the main engine loop.
3648#[derive(Debug)]
3649enum LoopEvent<T, N>
3650where
3651    N: NodePrimitives,
3652    T: PayloadTypes,
3653{
3654    /// An engine API message was received.
3655    EngineMessage(FromEngine<EngineApiRequest<T, N>, N::Block>),
3656    /// A persistence task completed.
3657    PersistenceComplete {
3658        /// The unified result of the persistence operation.
3659        result: PersistenceResult,
3660        /// When the persistence operation started.
3661        start_time: Instant,
3662    },
3663    /// The last active payload job has finished, so suppressed persistence may resume.
3664    PayloadBuildFinished,
3665    /// A channel was disconnected.
3666    Disconnected,
3667}
3668
3669/// Tracks payload jobs that may access the current in-memory overlay.
3670#[derive(Clone, Debug)]
3671struct PayloadBuildTracker {
3672    active: Arc<AtomicUsize>,
3673    finished_tx: Sender<()>,
3674}
3675
3676impl PayloadBuildTracker {
3677    /// Creates a tracker and a receiver notified when its active count reaches zero.
3678    fn new() -> (Self, Receiver<()>) {
3679        let (finished_tx, finished_rx) = crossbeam_channel::bounded(1);
3680        (Self { active: Arc::new(AtomicUsize::new(0)), finished_tx }, finished_rx)
3681    }
3682
3683    /// Acquires a lease for one payload job.
3684    fn acquire(&self) -> PayloadBuildLease {
3685        self.active.fetch_add(1, Ordering::AcqRel);
3686        PayloadBuildLease {
3687            active: Arc::clone(&self.active),
3688            finished_tx: self.finished_tx.clone(),
3689        }
3690    }
3691
3692    /// Returns whether at least one payload job is active.
3693    fn is_active(&self) -> bool {
3694        self.active.load(Ordering::Acquire) != 0
3695    }
3696}
3697
3698/// A lease held for the lifetime of a payload job.
3699#[derive(Debug)]
3700struct PayloadBuildLease {
3701    active: Arc<AtomicUsize>,
3702    finished_tx: Sender<()>,
3703}
3704
3705impl Drop for PayloadBuildLease {
3706    fn drop(&mut self) {
3707        let previous = self.active.fetch_sub(1, Ordering::AcqRel);
3708        debug_assert!(previous > 0, "payload build lease count underflow");
3709
3710        if previous == 1 {
3711            // The bounded channel coalesces completion notifications. The engine always checks
3712            // the counter again before applying a pending handoff.
3713            let _ = self.finished_tx.try_send(());
3714        }
3715    }
3716}
3717
3718/// Block inclusion can be valid, accepted, or invalid. Invalid blocks are returned as an error
3719/// variant.
3720///
3721/// If we don't know the block's parent, we return `Disconnected`, as we can't claim that the block
3722/// is valid or not.
3723#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3724pub enum BlockStatus {
3725    /// The block is valid: its parent state was available, so it was executed and inserted into
3726    /// the tree.
3727    ///
3728    /// Note: this does not imply the block extends the canonical chain. Blocks on a fork are
3729    /// executed and inserted the same way and report this status as well.
3730    Valid,
3731    /// The block may be valid and has an unknown missing ancestor.
3732    Disconnected {
3733        /// Current canonical head.
3734        head: BlockNumHash,
3735        /// The lowest ancestor block that is not connected to the canonical chain.
3736        missing_ancestor: BlockNumHash,
3737    },
3738}
3739
3740/// How a payload was inserted if it was valid.
3741///
3742/// If the payload was valid, but has already been seen, [`InsertPayloadOk::AlreadySeen`] is
3743/// returned, otherwise [`InsertPayloadOk::Inserted`] is returned.
3744#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3745pub enum InsertPayloadOk {
3746    /// The payload was valid, but we have already seen it.
3747    AlreadySeen(BlockStatus),
3748    /// The payload was valid and inserted into the tree.
3749    Inserted(BlockStatus),
3750}
3751
3752/// Target for block persistence.
3753#[derive(Debug, Clone, Copy)]
3754enum PersistTarget {
3755    /// Persist up to `canonical_head - memory_block_buffer_target`.
3756    Threshold,
3757    /// Persist all blocks up to and including the canonical head.
3758    Head,
3759    /// Persist state/trie updates through the persisted block frontier.
3760    Persisted,
3761}
3762
3763/// Result of waiting for caches to become available.
3764#[derive(Debug, Clone, Copy, Default)]
3765pub struct CacheWaitDurations {
3766    /// Time spent waiting for the execution cache lock, excluding post-unlock cache destruction.
3767    pub execution_cache: Duration,
3768    /// Time spent waiting for the sparse trie lock.
3769    pub sparse_trie: Duration,
3770}
3771
3772/// Trait for types that can wait for caches to become available.
3773///
3774/// Used by `reth_newPayload` to wait for cache updates before starting payload processing.
3775/// Removed execution-cache allocations may still be destroyed concurrently after unlocking.
3776pub trait WaitForCaches {
3777    /// Waits for cache updates to complete.
3778    ///
3779    /// Returns the time spent waiting for each cache separately.
3780    fn wait_for_caches(&self) -> CacheWaitDurations;
3781}