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
99const CHANGESET_CACHE_RETENTION_BLOCKS: u64 = 64;
104
105#[derive(Debug)]
109pub struct EngineApiTreeState<N: NodePrimitives> {
110 tree_state: TreeState<N>,
112 pending_sparse_trie_prune: bool,
114 forkchoice_state_tracker: ForkchoiceStateTracker,
116 buffer: BlockBuffer<N::Block>,
118 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 pub const fn tree_state(&self) -> &TreeState<N> {
146 &self.tree_state
147 }
148
149 pub const fn pending_sparse_trie_prune(&self) -> bool {
151 self.pending_sparse_trie_prune
152 }
153
154 pub const fn set_pending_sparse_trie_prune(&mut self, pending: bool) {
156 self.pending_sparse_trie_prune = pending;
157 }
158
159 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 pub fn has_invalid_header(&mut self, hash: &B256) -> bool {
184 self.invalid_headers.get(hash).is_some()
185 }
186}
187
188#[derive(Debug)]
190pub struct TreeOutcome<T> {
191 pub outcome: T,
193 pub event: Option<TreeEvent>,
195 pub already_seen: bool,
198}
199
200impl<T> TreeOutcome<T> {
201 pub const fn new(outcome: T) -> Self {
203 Self { outcome, event: None, already_seen: false }
204 }
205
206 pub fn with_event(mut self, event: TreeEvent) -> Self {
208 self.event = Some(event);
209 self
210 }
211
212 pub const fn with_already_seen(mut self, value: bool) -> Self {
214 self.already_seen = value;
215 self
216 }
217}
218
219#[derive(Debug)]
221pub struct TryInsertPayloadResult {
222 pub status: PayloadStatus,
226 pub already_seen: bool,
228}
229
230impl TryInsertPayloadResult {
231 #[inline]
233 pub fn into_outcome(self) -> TreeOutcome<PayloadStatus> {
234 TreeOutcome::new(self.status).with_already_seen(self.already_seen)
235 }
236}
237
238#[derive(Debug)]
240pub enum TreeEvent {
241 TreeAction(TreeAction),
243 BackfillAction(BackfillAction),
245 Download(DownloadRequest),
247}
248
249impl TreeEvent {
250 const fn is_backfill_action(&self) -> bool {
252 matches!(self, Self::BackfillAction(_))
253 }
254}
255
256#[derive(Debug)]
258pub enum TreeAction {
259 MakeCanonical {
261 sync_target_head: B256,
263 },
264}
265
266pub 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 state: EngineApiTreeState<N>,
281 incoming_tx: Sender<FromEngine<EngineApiRequest<T, N>, N::Block>>,
290 incoming: Receiver<FromEngine<EngineApiRequest<T, N>, N::Block>>,
292 outgoing: UnboundedSender<EngineApiEvent<N>>,
294 persistence: PersistenceHandle<N>,
296 persistence_state: PersistenceState,
298 backfill_sync_state: BackfillSyncState,
300 canonical_in_memory_state: CanonicalInMemoryState<N>,
303 payload_builder: PayloadBuilderHandle<T>,
306 config: TreeConfig,
308 metrics: EngineApiMetrics,
310 engine_kind: EngineApiKind,
312 evm_config: C,
314 execution_timing_stats: B256Map<Box<ExecutionTimingStats>>,
318 payload_builds: PayloadBuildTracker,
320 payload_build_finished: Receiver<()>,
322 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 #[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 #[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 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 pub fn sender(&self) -> Sender<FromEngine<EngineApiRequest<T, N>, N::Block>> {
493 self.incoming_tx.clone()
494 }
495
496 fn persistence_gap(&self) -> u64 {
499 self.canonical_in_memory_state.canonical_chain().count() as u64
500 }
501
502 fn persistence_backpressure_gap(&self) -> u64 {
504 self.persistence_gap().saturating_sub(self.config.memory_block_buffer_target())
505 }
506
507 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 pub fn run(mut self) {
521 loop {
522 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 if let Err(err) = self.advance_persistence() {
602 error!(target: "engine::tree", %err, "Advancing persistence failed");
603 return
604 }
605 }
606 }
607
608 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 fn wait_for_event(&mut self) -> LoopEvent<T, N> {
631 let maybe_persistence = self.persistence_state.rx.take();
633
634 if let Some((persistence_rx, start_time, action)) = maybe_persistence {
635 crossbeam_channel::select_biased! {
638 recv(persistence_rx) -> result => {
639 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 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 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 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 fn on_downloaded(
686 &mut self,
687 mut blocks: Vec<SealedBlockWithAccessList<N::Block>>,
688 ) -> Result<Option<TreeEvent>, InsertBlockFatalError> {
689 if blocks.is_empty() {
690 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 return Ok(None)
703 }
704 }
705 }
706
707 if !blocks.is_empty() {
709 let _ = self.incoming_tx.send(FromEngine::DownloadedBlocks(blocks));
710 }
711
712 Ok(None)
713 }
714
715 #[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 let start = Instant::now();
745
746 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 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 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 outcome.outcome.is_valid() && self.is_sync_target_head(block_hash) {
795 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 self.metrics.block_validation.total_duration.record(start.elapsed().as_secs_f64());
805
806 Ok(outcome)
807 }
808
809 #[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 (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 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 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 fn on_new_head(&self, new_head: B256) -> ProviderResult<Option<NewCanonicalChain<N>>> {
896 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 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 return Ok(None)
925 }
926 }
927
928 if current_hash == self.state.tree_state.current_canonical_head.hash {
931 new_chain.reverse();
932
933 return Ok(Some(NewCanonicalChain::Commit { new: new_chain }))
935 }
936
937 let mut old_chain = Vec::new();
939 let mut old_hash = self.state.tree_state.current_canonical_head.hash;
940
941 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 debug_assert_eq!(current_number, current_canonical_number);
952
953 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 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 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 self.state.tree_state.set_canonical_head(canonical_header.num_hash());
998
999 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 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 let old_blocks =
1039 self.collect_blocks_for_canonical_unwind(new_head_number, current_head_number);
1040
1041 self.apply_canonical_ancestor_via_reorg(canonical_header, old_blocks)
1043 }
1044
1045 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 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 let executed_block = self.canonical_block_by_hash(new_head_hash)?;
1087 self.canonical_in_memory_state
1089 .update_chain(NewCanonicalChain::Reorg { new: vec![executed_block], old: old_blocks });
1090
1091 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 fn handle_chain_advance_or_same_height(
1107 &self,
1108 canonical_header: &SealedHeader<N::BlockHeader>,
1109 ) -> ProviderResult<()> {
1110 self.ensure_block_in_memory(canonical_header.number(), canonical_header.hash())?;
1112
1113 self.canonical_in_memory_state.set_canonical_head(canonical_header.clone());
1115
1116 Ok(())
1117 }
1118
1119 fn ensure_block_in_memory(&self, block_number: u64, block_hash: B256) -> ProviderResult<()> {
1121 if self.canonical_in_memory_state.state_by_number(block_number).is_some() {
1123 return Ok(());
1124 }
1125
1126 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 #[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 self.record_forkchoice_metrics();
1159
1160 if let Some(early_result) = self.validate_forkchoice_state(state)? {
1162 return Ok(TreeOutcome::new(early_result));
1163 }
1164
1165 if let Some(result) = self.handle_canonical_head(state, &attrs)? {
1167 return Ok(result);
1168 }
1169
1170 if let Some(result) = self.apply_chain_update(state, &attrs)? {
1173 return Ok(result);
1174 }
1175
1176 self.handle_missing_block(state)
1178 }
1179
1180 fn record_forkchoice_metrics(&self) {
1182 self.canonical_in_memory_state.on_forkchoice_update_received();
1183 }
1184
1185 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 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 trace!(target: "engine::tree", "Pipeline is syncing, skipping forkchoice update");
1208 return Ok(Some(OnForkChoiceUpdated::syncing()));
1209 }
1210
1211 Ok(None)
1212 }
1213
1214 fn handle_canonical_head(
1220 &mut self,
1221 state: ForkchoiceState,
1222 attrs: &Option<T::PayloadAttributes>, ) -> ProviderResult<Option<TreeOutcome<OnForkChoiceUpdated>>> {
1224 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 if let Err(outcome) = self.ensure_consistent_forkchoice_state(state) {
1250 return Ok(Some(TreeOutcome::new(outcome)));
1252 }
1253
1254 self.payload_validator.on_canonical_head_changed(state.head_block_hash, &self.state);
1255
1256 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 ProviderError::HeaderNotFound(state.head_block_hash.into())
1264 })?;
1265 let updated = self.process_payload_attributes(attr.clone(), &tip, state);
1267 return Ok(Some(TreeOutcome::new(updated)));
1268 }
1269
1270 Ok(Some(Self::valid_outcome(state)))
1272 }
1273
1274 fn apply_chain_update(
1286 &mut self,
1287 state: ForkchoiceState,
1288 attrs: &Option<T::PayloadAttributes>,
1289 ) -> ProviderResult<Option<TreeOutcome<OnForkChoiceUpdated>>> {
1290 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 let always_trigger_payload_job = self.engine_kind.is_opstack() ||
1297 self.config.always_process_payload_attributes_on_canonical_head();
1298
1299 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 if always_trigger_payload_job && self.config.unwind_canonical_header() {
1326 self.update_latest_block_to_canonical_ancestor(&canonical_header)?;
1327 }
1328
1329 if let Some(attr) = attrs {
1334 debug!(target: "engine::tree", head = canonical_header.number(), "handling payload attributes for canonical head");
1335 let updated =
1337 self.process_payload_attributes(attr.clone(), &canonical_header, state);
1338 return Ok(Some(TreeOutcome::new(updated)));
1339 }
1340
1341 return Ok(Some(Self::valid_outcome(state)));
1344 }
1345
1346 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 if let Err(outcome) = self.ensure_consistent_forkchoice_state(state) {
1357 return Ok(Some(TreeOutcome::new(outcome)));
1359 }
1360
1361 if let Some(attr) = attrs {
1362 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 fn handle_missing_block(
1378 &self,
1379 state: ForkchoiceState,
1380 ) -> ProviderResult<TreeOutcome<OnForkChoiceUpdated>> {
1381 let target = if self.state.forkchoice_state_tracker.is_empty() &&
1388 !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 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 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 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 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 fn persist_until_complete(&mut self) -> Result<(), AdvancePersistenceError> {
1479 loop {
1480 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 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 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 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 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 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 finalized.number.min(min_threshold)
1557 } else {
1558 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 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 return Ok(ops::ControlFlow::Continue(()))
1608 }
1609
1610 if self.state.tree_state.contains_hash(&block_num_hash.hash) {
1611 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 self.state
1662 .forkchoice_state_tracker
1663 .set_latest(state, res.outcome.forkchoice_status());
1664
1665 self.emit_event(ConsensusEngineEvent::ForkchoiceUpdated(
1667 state,
1668 res.outcome.forkchoice_status(),
1669 ));
1670
1671 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 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 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 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 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 self.state.invalid_headers.insert(**bad_block);
1855
1856 Some(*target)
1858 } else {
1859 ctrl.block_number()
1861 };
1862
1863 let Some(backfill_height) = backfill_height else { return Ok(()) };
1865
1866 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 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 self.state.buffer.remove_old_blocks(backfill_height);
1898 self.purge_timing_stats(backfill_height, None);
1899 self.canonical_in_memory_state.clear_state();
1902
1903 if let Ok(Some(new_head)) = self.provider.sealed_header(backfill_height) {
1904 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 self.canonical_in_memory_state.set_canonical_head(new_head);
1911
1912 if !ctrl.is_unwind() {
1916 self.on_canonicalized_sync_target(backfill_num_hash.hash);
1917 }
1918 }
1919
1920 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 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 let newest_target = self.state.buffer.block(&target_hash).map(|block| block.number());
1936
1937 if let Some(backfill_target) =
1943 ctrl.block_number().zip(newest_target).and_then(|(progress, target_number)| {
1944 self.backfill_sync_target(progress, target_number, None)
1947 })
1948 {
1949 self.emit_event(EngineApiEvent::BackfillAction(BackfillAction::Start(
1951 backfill_target.into(),
1952 )));
1953 return Ok(())
1954 };
1955
1956 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 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 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 self.try_connect_buffered_blocks(self.state.tree_state.current_canonical_head)
1997 }
1998
1999 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 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 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 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 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 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 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 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 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 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 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 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 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 fn get_save_blocks_input(&self, target: PersistTarget) -> Option<SaveBlocksInput<N>> {
2219 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 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(¤t_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 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 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 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 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 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 const fn should_prune_sparse_trie(&self) -> bool {
2354 self.config.use_state_root_task()
2355 }
2356
2357 #[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 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 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 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 fn sealed_header_by_hash(
2439 &self,
2440 hash: B256,
2441 ) -> ProviderResult<Option<SealedHeader<N::BlockHeader>>> {
2442 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 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 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 fn latest_valid_hash_for_invalid_payload(
2489 &mut self,
2490 parent_hash: B256,
2491 ) -> ProviderResult<Option<B256>> {
2492 if self.has_block_by_hash(parent_hash)? {
2494 return Ok(Some(parent_hash))
2495 }
2496
2497 let mut current_hash = parent_hash;
2500 let mut current_block = self.state.invalid_headers.get(¤t_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(¤t_hash);
2504
2505 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 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 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 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 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 fn check_invalid_ancestor_with_head(
2559 &mut self,
2560 check: B256,
2561 head: &SealedBlock<N::Block>,
2562 ) -> ProviderResult<Option<PayloadStatus>> {
2563 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 fn on_invalid_new_payload(
2571 &mut self,
2572 head: SealedBlock<N::Block>,
2573 invalid: BlockWithParent,
2574 ) -> ProviderResult<PayloadStatus> {
2575 let status = self.prepare_invalid_response(invalid.parent)?;
2577
2578 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 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 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 self.state.invalid_headers.get(&lowest_buffered_ancestor)
2617 }
2618
2619 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 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 fn check_invalid_ancestor(&mut self, head: B256) -> ProviderResult<Option<PayloadStatus>> {
2651 let Some(header) = self.state.invalid_headers.get(&head) else { return Ok(None) };
2653
2654 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 Ok(Some(PayloadStatus::from_status(PayloadStatusEnum::Invalid {
2661 validation_error: PayloadValidationError::LinksToRejectedPayload.to_string(),
2662 })))
2663 }
2664 }
2665 }
2666
2667 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 #[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 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 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 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 #[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 #[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 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 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 let exceeds_backfill_threshold = match downloaded_block.as_ref() {
2793 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 Some(buffered_target) => {
2800 self.exceeds_backfill_run_threshold(canonical_tip_num, buffered_target.number())
2801 }
2802 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 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 Ok(None) if !target_hash.is_zero() => Some(target_hash),
2819 Ok(None) => {
2820 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 Ok(Some(_)) => None,
2837 }
2838 }
2839
2840 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 while canonical.number > persisted.number {
2856 canonical = parent_num_hash(canonical)?;
2857 }
2858
2859 if canonical == persisted {
2861 return Ok(None);
2862 }
2863
2864 while persisted.number > canonical.number {
2870 persisted = parent_num_hash(persisted)?;
2871 }
2872
2873 debug_assert_eq!(persisted.number, canonical.number);
2874
2875 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 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 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 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 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 self.metrics.tree.canonical_chain_height.set(tip.number() as f64);
2918
2919 self.canonical_in_memory_state.notify_canon_state(notification);
2921
2922 self.emit_event(ConsensusEngineEvent::CanonicalChainCommitted(
2924 Box::new(tip),
2925 start.elapsed(),
2926 ));
2927 }
2928
2929 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 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 fn on_disconnected_downloaded_block(
2970 &self,
2971 downloaded_block: BlockNumHash,
2972 missing_parent: BlockNumHash,
2973 head: BlockNumHash,
2974 ) -> Option<TreeEvent> {
2975 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 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 DownloadRequest::single_block(missing_parent.hash)
3002 };
3003
3004 Some(TreeEvent::Download(request.with_access_lists(self.should_download_access_lists())))
3005 }
3006
3007 fn on_valid_downloaded_block(
3014 &mut self,
3015 block_num_hash: BlockNumHash,
3016 ) -> Result<Option<TreeEvent>, InsertBlockFatalError> {
3017 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 return Ok(Some(TreeEvent::TreeAction(TreeAction::MakeCanonical {
3027 sync_target_head: block_num_hash.hash,
3028 })))
3029 }
3030
3031 self.make_canonical(block_num_hash.hash)?;
3035 self.try_connect_buffered_blocks(block_num_hash)?;
3036
3037 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 #[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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 let latest_valid_hash =
3364 if error.is_block_hash_mismatch() || error.is_invalid_versioned_hashes() {
3365 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 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 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 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 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 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 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 let _ = self.persistence.save_finalized_block_number(finalized.number());
3477 self.canonical_in_memory_state.set_finalized(finalized.clone());
3478 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 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 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 let _ = self.persistence.save_safe_block_number(safe.number());
3507 self.canonical_in_memory_state.set_safe(safe.clone());
3508 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 fn ensure_consistent_forkchoice_state(
3529 &self,
3530 state: ForkchoiceState,
3531 ) -> Result<(), OnForkChoiceUpdated> {
3532 self.update_finalized_block(state.finalized_block_hash)?;
3538
3539 self.update_safe_block(state.safe_block_hash)
3545 }
3546
3547 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 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 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 OnForkChoiceUpdated::updated_with_pending_payload_id(
3613 PayloadStatus::new(PayloadStatusEnum::Valid, Some(state.head_block_hash)),
3614 pending_payload_id,
3615 )
3616 }
3617
3618 pub(crate) fn remove_before(
3625 &mut self,
3626 upper_bound: BlockNumHash,
3627 finalized_hash: Option<B256>,
3628 ) -> ProviderResult<()> {
3629 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#[derive(Debug)]
3649enum LoopEvent<T, N>
3650where
3651 N: NodePrimitives,
3652 T: PayloadTypes,
3653{
3654 EngineMessage(FromEngine<EngineApiRequest<T, N>, N::Block>),
3656 PersistenceComplete {
3658 result: PersistenceResult,
3660 start_time: Instant,
3662 },
3663 PayloadBuildFinished,
3665 Disconnected,
3667}
3668
3669#[derive(Clone, Debug)]
3671struct PayloadBuildTracker {
3672 active: Arc<AtomicUsize>,
3673 finished_tx: Sender<()>,
3674}
3675
3676impl PayloadBuildTracker {
3677 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 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 fn is_active(&self) -> bool {
3694 self.active.load(Ordering::Acquire) != 0
3695 }
3696}
3697
3698#[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 let _ = self.finished_tx.try_send(());
3714 }
3715 }
3716}
3717
3718#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3724pub enum BlockStatus {
3725 Valid,
3731 Disconnected {
3733 head: BlockNumHash,
3735 missing_ancestor: BlockNumHash,
3737 },
3738}
3739
3740#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3745pub enum InsertPayloadOk {
3746 AlreadySeen(BlockStatus),
3748 Inserted(BlockStatus),
3750}
3751
3752#[derive(Debug, Clone, Copy)]
3754enum PersistTarget {
3755 Threshold,
3757 Head,
3759 Persisted,
3761}
3762
3763#[derive(Debug, Clone, Copy, Default)]
3765pub struct CacheWaitDurations {
3766 pub execution_cache: Duration,
3768 pub sparse_trie: Duration,
3770}
3771
3772pub trait WaitForCaches {
3777 fn wait_for_caches(&self) -> CacheWaitDurations;
3781}