1use crate::tree::{
99 error::{
100 BlockAccessListDecodeError, InsertBlockError, InsertBlockErrorKind, InsertPayloadError,
101 },
102 instrumented_state::{InstrumentedStateProvider, StateProviderMetrics, StateProviderStats},
103 payload_processor::PayloadProcessor,
104 precompile_cache::{CachedPrecompile, CachedPrecompileMetrics, PrecompileCacheMap},
105 txpool_prewarm,
106 types::{InsertPayloadResult, ValidationOutput},
107 CacheWaitDurations, CachedStateProvider, EngineApiMetrics, EngineApiTreeState, ExecutionEnv,
108 PayloadHandle, StateProviderDatabase, TreeConfig, WaitForCaches,
109};
110use alloy_consensus::transaction::{Either, TxHashRef};
111use alloy_eip7928::{
112 bal::{Bal, DecodedBal},
113 BlockAccessList,
114};
115use alloy_eips::{eip1898::BlockWithParent, eip4895::Withdrawal, NumHash};
116use alloy_evm::Evm;
117use alloy_primitives::{
118 map::{AddressMap, B256Set},
119 B256,
120};
121use reth_tasks::LazyHandle;
122
123use crate::tree::{
124 payload_processor::receipt_root_task::{IndexedReceipt, ReceiptRootTaskHandle},
125 state_root_strategy::{
126 DefaultStateRootStrategy, LazyHashedPostState, PayloadStateRootHandle,
127 PayloadStateRootJobContext, StateRootHintStream, StateRootJobContext, StateRootStrategy,
128 StateRootUpdateStream,
129 },
130};
131use alloy_primitives::Address;
132use reth_chain_state::{CanonicalInMemoryState, ExecutedBlock, ExecutionTimingStats};
133use reth_consensus::{ConsensusError, FullConsensus, ReceiptRootBloom};
134use reth_engine_primitives::{
135 ConfigureEngineEvm, ExecutableTxIterator, ExecutionPayload, InvalidBlockHook, PayloadValidator,
136};
137use reth_errors::{BlockExecutionError, BlockValidationError, ProviderResult};
138use reth_evm::{
139 block::BlockExecutor, execute::ExecutableTxFor, ConfigureEvm, EvmEnvFor, ExecutionCtxFor,
140 OnStateHook, SpecFor,
141};
142use reth_execution_cache::{CacheFillMode, CacheStats};
143use reth_execution_types::DecodedRevmBal;
144use reth_network_p2p::full_block::SealedBlockWithAccessList;
145use reth_payload_builder::{PayloadBuilderLease, PayloadBuilderResources};
146use reth_payload_primitives::{
147 BuiltPayload, BuiltPayloadExecutedBlock, InvalidPayloadAttributesError, NewPayloadError,
148 PayloadTypes,
149};
150use reth_primitives_traits::{
151 AlloyBlockHeader, BlockBody, BlockTy, FastInstant as Instant, GotExpected, NodePrimitives,
152 RecoveredBlock, SealedBlock, SealedHeader, SignerRecoverable,
153};
154use reth_provider::{
155 BlockExecutionOutput, BlockHashReader, BlockReader, ChangeSetReader, DatabaseProviderFactory,
156 DatabaseProviderROFactory, EvmStateProvider, EvmStateProviderBox, HashedPostStateProvider,
157 HistoryReader, ProviderError, PruneCheckpointReader, StageCheckpointReader, StateProvider,
158 StateProviderFactory, StateReader, StateRootProvider, StorageChangeSetReader,
159 StorageSettingsCache,
160};
161use reth_revm::db::{states::bundle_state::BundleRetention, BundleAccount, State};
162use reth_storage_overlay::{OverlayManager, OverlayStateProviderFactory};
163use reth_trie::{
164 hashed_cursor::HashedCursorFactory,
165 trie_cursor::TrieCursorFactory,
166 updates::{TrieUpdates, TrieUpdatesSorted},
167 HashedPostState, KeccakKeyHasher, LazyHashedPostStateSorted,
168};
169use revm::state::{bal::Bal as RevmBal, AccountInfo};
170use std::{
171 sync::{
172 atomic::{AtomicUsize, Ordering},
173 Arc,
174 },
175 time::Duration,
176};
177use tracing::{debug, debug_span, error, info, instrument, trace, warn, Level, Span};
178
179pub use crate::tree::types::ValidationOutcome;
180
181const MAX_EXPECTED_GAS_LIMIT_MULTIPLIER: u64 = 2;
184
185const DEFERRED_TRIE_WORKER_NAME: &str = "deferred-trie";
187
188type ReceiptRootSender<N> =
189 crossbeam_channel::Sender<IndexedReceipt<<N as NodePrimitives>::Receipt>>;
190type ReceiptRootReceiver = tokio::sync::oneshot::Receiver<(B256, alloy_primitives::Bloom)>;
191
192pub struct TreeCtx<'a, N: NodePrimitives> {
197 state: &'a mut EngineApiTreeState<N>,
199 canonical_in_memory_state: &'a CanonicalInMemoryState<N>,
201}
202
203impl<'a, N: NodePrimitives> std::fmt::Debug for TreeCtx<'a, N> {
204 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
205 f.debug_struct("TreeCtx")
206 .field("state", &"EngineApiTreeState")
207 .field("canonical_in_memory_state", &self.canonical_in_memory_state)
208 .finish()
209 }
210}
211
212impl<'a, N: NodePrimitives> TreeCtx<'a, N> {
213 pub const fn new(
215 state: &'a mut EngineApiTreeState<N>,
216 canonical_in_memory_state: &'a CanonicalInMemoryState<N>,
217 ) -> Self {
218 Self { state, canonical_in_memory_state }
219 }
220}
221
222impl<'a, N: NodePrimitives> TreeCtx<'a, N> {
223 pub const fn state(&self) -> &EngineApiTreeState<N> {
225 &*self.state
226 }
227
228 pub const fn state_mut(&mut self) -> &mut EngineApiTreeState<N> {
230 self.state
231 }
232
233 pub const fn canonical_in_memory_state(&self) -> &'a CanonicalInMemoryState<N> {
235 self.canonical_in_memory_state
236 }
237}
238
239struct JitPauseGuard<Evm: ConfigureEvm>(Evm);
245
246impl<Evm: ConfigureEvm> JitPauseGuard<Evm> {
247 fn new(evm_config: &Evm) -> Self {
248 if let Some(jit_backend) = evm_config.jit_backend() {
249 jit_backend.pause();
250 }
251 Self(evm_config.clone())
252 }
253}
254
255impl<Evm: ConfigureEvm> Drop for JitPauseGuard<Evm> {
256 fn drop(&mut self) {
257 if let Some(jit_backend) = self.0.jit_backend() {
258 jit_backend.resume();
259 }
260 }
261}
262
263#[derive(derive_more::Debug)]
271pub struct BasicEngineValidator<P, Evm, V>
272where
273 Evm: ConfigureEvm,
274{
275 provider: P,
277 consensus: Arc<dyn FullConsensus<Evm::Primitives>>,
279 evm_config: Evm,
281 config: TreeConfig,
283 payload_processor: PayloadProcessor<Evm>,
285 precompile_cache_map: PrecompileCacheMap<SpecFor<Evm>>,
287 precompile_cache_metrics: AddressMap<CachedPrecompileMetrics>,
289 #[debug(skip)]
291 invalid_block_hook: Box<dyn InvalidBlockHook<Evm::Primitives>>,
292 metrics: EngineApiMetrics,
294 validator: V,
296 runtime: reth_tasks::Runtime,
298 overlay_manager: OverlayManager<Evm::Primitives>,
300 #[debug(skip)]
302 state_root_strategy: Arc<dyn StateRootStrategy<Evm::Primitives, P, Evm>>,
303 #[debug(skip)]
307 txpool_prewarm: Option<txpool_prewarm::Handle<Evm::Primitives, P, Evm>>,
308 bal_hash_buf: Vec<u8>,
310}
311
312impl<N, P, Evm, V> BasicEngineValidator<P, Evm, V>
313where
314 N: NodePrimitives,
315 P: DatabaseProviderFactory<
316 Provider: BlockReader
317 + BlockHashReader
318 + StageCheckpointReader
319 + PruneCheckpointReader
320 + ChangeSetReader
321 + StorageChangeSetReader
322 + StorageSettingsCache
323 + HistoryReader
324 + 'static,
325 > + BlockReader<Header = N::BlockHeader>
326 + ChangeSetReader
327 + StateProviderFactory
328 + StateReader
329 + Clone
330 + 'static,
331 OverlayStateProviderFactory<P, N>: DatabaseProviderROFactory<
332 Provider: TrieCursorFactory
333 + HashedCursorFactory
334 + HashedPostStateProvider
335 + StateRootProvider
336 + StateProvider
337 + Send,
338 > + Clone
339 + 'static,
340 Evm: ConfigureEvm<Primitives = N> + 'static,
341{
342 #[expect(clippy::too_many_arguments)]
344 pub fn new(
345 provider: P,
346 consensus: Arc<dyn FullConsensus<N>>,
347 evm_config: Evm,
348 validator: V,
349 config: TreeConfig,
350 invalid_block_hook: Box<dyn InvalidBlockHook<N>>,
351 overlay_manager: OverlayManager<N>,
352 runtime: reth_tasks::Runtime,
353 ) -> Self {
354 let precompile_cache_map = PrecompileCacheMap::default();
355 let payload_processor = PayloadProcessor::new(
356 runtime.clone(),
357 evm_config.clone(),
358 &config,
359 precompile_cache_map.clone(),
360 );
361 Self {
362 provider,
363 consensus,
364 evm_config,
365 payload_processor,
366 precompile_cache_map,
367 precompile_cache_metrics: AddressMap::default(),
368 config,
369 invalid_block_hook,
370 metrics: EngineApiMetrics::default(),
371 validator,
372 runtime,
373 overlay_manager,
374 state_root_strategy: Arc::new(DefaultStateRootStrategy::default()),
375 txpool_prewarm: None,
376 bal_hash_buf: Vec::new(),
377 }
378 }
379
380 pub fn with_state_root_strategy(
382 mut self,
383 state_root_strategy: Arc<dyn StateRootStrategy<N, P, Evm>>,
384 ) -> Self {
385 self.state_root_strategy = state_root_strategy;
386 self
387 }
388
389 pub fn with_txpool_prewarming(
391 mut self,
392 source: impl crate::tree::TxPoolPrewarmSource<N> + 'static,
393 ) -> Self {
394 self.txpool_prewarm = Some(txpool_prewarm::Handle::spawn(
395 &self.runtime,
396 Arc::new(source),
397 self.evm_config.clone(),
398 ));
399 self
400 }
401
402 #[instrument(level = "debug", target = "engine::tree::payload_validator", skip_all)]
404 pub fn convert_to_block<T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
405 &self,
406 input: BlockOrPayload<T>,
407 ) -> Result<SealedBlock<N::Block>, NewPayloadError>
408 where
409 V: PayloadValidator<T, Block = N::Block>,
410 {
411 match input {
412 BlockOrPayload::Payload(payload) => self.validator.convert_payload_to_block(payload),
413 BlockOrPayload::Block(block) => Ok(block.split().0),
414 }
415 }
416
417 pub fn evm_env_for<T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
419 &self,
420 input: &BlockOrPayload<T>,
421 ) -> Result<EvmEnvFor<Evm>, Evm::Error>
422 where
423 V: PayloadValidator<T, Block = N::Block>,
424 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
425 {
426 match input {
427 BlockOrPayload::Payload(payload) => Ok(self.evm_config.evm_env_for_payload(payload)?),
428 BlockOrPayload::Block(block) => Ok(self.evm_config.evm_env(block.header())?),
429 }
430 }
431
432 pub fn tx_iterator_for<'a, T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
434 &'a self,
435 input: &'a BlockOrPayload<T>,
436 ) -> Result<impl ExecutableTxIterator<Evm>, NewPayloadError>
437 where
438 V: PayloadValidator<T, Block = N::Block>,
439 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
440 {
441 Ok(match input {
442 BlockOrPayload::Payload(payload) => {
443 let iter = self
444 .evm_config
445 .tx_iterator_for_payload(payload)
446 .map_err(NewPayloadError::other)?;
447 Either::Left(iter)
448 }
449 BlockOrPayload::Block(block) => {
450 let txs = block.body().clone_transactions();
451 let convert = |tx: N::SignedTx| tx.try_into_recovered();
452 Either::Right((txs, convert))
453 }
454 })
455 }
456
457 pub fn execution_ctx_for<'a, T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
459 &self,
460 input: &'a BlockOrPayload<T>,
461 ) -> Result<ExecutionCtxFor<'a, Evm>, Evm::Error>
462 where
463 V: PayloadValidator<T, Block = N::Block>,
464 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
465 {
466 match input {
467 BlockOrPayload::Payload(payload) => Ok(self.evm_config.context_for_payload(payload)?),
468 BlockOrPayload::Block(block) => Ok(self.evm_config.context_for_block(block)?),
469 }
470 }
471
472 #[instrument(
480 level = "debug",
481 target = "engine::tree::payload_validator",
482 skip_all,
483 fields(
484 parent = ?input.parent_hash(),
485 type_name = ?input.type_name(),
486 )
487 )]
488 pub fn validate_block_with_state<T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
489 &mut self,
490 input: BlockOrPayload<T>,
491 mut ctx: TreeCtx<'_, N>,
492 ) -> InsertPayloadResult<N>
493 where
494 V: PayloadValidator<T, Block = N::Block> + Clone,
495 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
496 {
497 let parent_hash = input.parent_hash();
498 let _txpool_pause = self.txpool_prewarm.as_ref().map(txpool_prewarm::Handle::pause);
499 let txpool_snapshot =
500 self.txpool_prewarm.as_ref().and_then(|prewarmer| prewarmer.snapshot(parent_hash));
501 let _jit_pause = JitPauseGuard::new(&self.evm_config);
502
503 let parent_block = match self.sealed_header_by_hash(parent_hash, ctx.state()) {
506 Ok(Some(parent_block)) => parent_block,
507 Ok(None) => {
508 return Err(InsertBlockError::new(
509 self.convert_to_block(input)?,
510 ProviderError::HeaderNotFound(parent_hash.into()).into(),
511 )
512 .into())
513 }
514 Err(e) => {
515 return Err(InsertBlockError::new(self.convert_to_block(input)?, e.into()).into())
516 }
517 };
518
519 let validated_block = self.spawn_convert_and_validate(&input, parent_block.clone());
523
524 macro_rules! ensure_ok {
527 ($expr:expr) => {
528 match $expr {
529 Ok(val) => val,
530 Err(e) => {
531 let block = validated_block.try_into_inner().expect("sole handle")?;
532 return Err(InsertBlockError::new(block, e.into()).into())
533 }
534 }
535 };
536 }
537
538 macro_rules! ensure_ok_post_block {
540 ($expr:expr, $block:expr) => {
541 match $expr {
542 Ok(val) => val,
543 Err(e) => {
544 return Err(
545 InsertBlockError::new($block.into_sealed_block(), e.into()).into()
546 )
547 }
548 }
549 };
550 }
551
552 if input.gas_limit() >
555 parent_block.gas_limit().saturating_mul(MAX_EXPECTED_GAS_LIMIT_MULTIPLIER)
556 {
557 if validated_block.get().is_err() {
559 return Err(validated_block
560 .try_into_inner()
561 .expect("sole handle")
562 .expect_err("Err result checked"))
563 }
564 }
565
566 trace!(target: "engine::tree::payload_validator", "Fetching block state provider");
567 let _enter =
568 debug_span!(target: "engine::tree::payload_validator", "state_provider").entered();
569 let Some(state_provider_factory) =
570 ensure_ok!(self.overlay_state_provider_factory(parent_hash, ctx.state()))
571 else {
572 return Err(InsertBlockError::new(
574 validated_block.try_into_inner().expect("sole handle")?,
575 ProviderError::HeaderNotFound(parent_hash.into()).into(),
576 )
577 .into())
578 };
579 drop(_enter);
580
581 let evm_env = debug_span!(target: "engine::tree::payload_validator", "evm_env")
582 .in_scope(|| self.evm_env_for(&input))
583 .map_err(NewPayloadError::other)?;
584
585 let decoded_bal =
588 ensure_ok!(input.try_decoded_access_list().map_err(BlockAccessListDecodeError::new))
589 .map(Arc::new);
590
591 if let Some(decoded_bal) = decoded_bal.as_deref() {
592 ensure_ok!(decoded_bal
594 .as_bal()
595 .validate_gas_limit(input.gas_limit())
596 .map_err(ConsensusError::from));
597 }
598
599 let env = ExecutionEnv {
600 evm_env,
601 hash: input.hash(),
602 parent_hash: input.parent_hash(),
603 parent_state_root: parent_block.state_root(),
604 transaction_count: input.transaction_count(),
605 gas_used: input.gas_used(),
606 withdrawals: input.withdrawals().map(|w| w.to_vec()),
607 decoded_bal: decoded_bal.as_ref().map(Arc::clone),
608 txpool_snapshot: txpool_snapshot.clone(),
609 };
610
611 let txs = self.tx_iterator_for(&input)?;
613
614 let parallel_bal_execution = ensure_ok!(self.bal_path_eligible(env.decoded_bal.as_deref()));
615
616 let mut state_root_job =
618 ensure_ok!(self.state_root_strategy.prepare(StateRootJobContext::new(
619 &self.runtime,
620 &self.overlay_manager,
621 &env,
622 &parent_block,
623 state_provider_factory.clone(),
624 &self.config,
625 parallel_bal_execution,
626 ctx.state_mut(),
627 )));
628 let state_root_job_name = state_root_job.name();
629
630 debug!(
631 target: "engine::tree::payload_validator",
632 strategy = state_root_job_name,
633 "Prepared state root job"
634 );
635
636 let execution_state_hook = state_root_job.take_execution_hook();
639 let hint_stream = state_root_job.take_hint_stream();
642 let hashed_update_stream = state_root_job.take_hashed_update_stream();
643
644 let mut handle = ensure_ok!(self.spawn_payload_processor(
646 env.clone(),
647 txs,
648 state_provider_factory.clone(),
649 hint_stream,
650 hashed_update_stream,
651 parallel_bal_execution,
652 ));
653
654 let slow_block_enabled = self.config.slow_block_threshold().is_some();
656 let cache_stats = slow_block_enabled.then(|| Arc::new(CacheStats::default()));
657 let instrument_state_provider = slow_block_enabled || self.config.state_provider_metrics();
658 let state_provider_metrics =
659 instrument_state_provider.then(|| StateProviderMetrics::with_source("engine"));
660 let state_provider_stats =
661 instrument_state_provider.then(|| Arc::new(StateProviderStats::default()));
662 let execution_cache = handle.caches().map(|caches| (caches, handle.cache_metrics()));
663
664 let make_state_provider = |fill_on_miss: bool| -> ProviderResult<EvmStateProviderBox> {
684 let provider = state_provider_factory.database_provider_ro()?.into_evm_state_provider();
685 let provider: EvmStateProviderBox =
686 if let Some((caches, cache_metrics)) = &execution_cache {
687 let fill_mode = if fill_on_miss {
688 CacheFillMode::FillOnMiss
689 } else {
690 CacheFillMode::LookupOnly
691 };
692 Box::new(
693 CachedStateProvider::new_with_mode(
694 provider,
695 caches.clone(),
696 fill_mode,
697 cache_metrics.clone(),
698 cache_stats.clone(),
699 )
700 .with_txpool_snapshot(txpool_snapshot.clone()),
701 )
702 } else {
703 Box::new(provider)
704 };
705
706 let provider: EvmStateProviderBox = if instrument_state_provider {
707 let stats = state_provider_stats
708 .as_ref()
709 .expect("instrumented state provider requires shared stats");
710 let metrics = state_provider_metrics
711 .as_ref()
712 .expect("instrumented state provider requires metrics");
713 Box::new(InstrumentedStateProvider::with_stats(
714 provider,
715 metrics.clone(),
716 Arc::clone(stats),
717 ))
718 } else {
719 provider
720 };
721
722 Ok(provider)
723 };
724
725 let execute_block_start = Instant::now();
729 let execution_result = if parallel_bal_execution {
730 self.execute_block_bal(env, &input, &handle, &make_state_provider)
731 } else {
732 let state_provider = make_state_provider(false);
733 match state_provider {
734 Ok(state_provider) => self.execute_block(
735 state_provider,
736 env,
737 &input,
738 &mut handle,
739 execution_state_hook,
740 ),
741 Err(err) => Err(err.into()),
742 }
743 };
744 let execution_duration = execute_block_start.elapsed();
745 if let (Some(metrics), Some(stats)) = (&state_provider_metrics, &state_provider_stats) {
746 metrics.record_totals(stats);
747 }
748 let (output, senders, receipt_root_rx, executed_bal) = ensure_ok!(execution_result);
749 let (built_bal, revm_bal) =
750 executed_bal.map(|ExecutedBal { alloy, revm }| (alloy, revm)).unzip();
751
752 handle.stop_prewarming_execution();
754
755 let output = Arc::new(output);
759
760 let valid_block_tx = handle.terminate_caching(Some(output.clone()));
763
764 let hashed_state_output = output.clone();
768 let mut hashed_state_rx = state_root_job.take_hashed_state_rx();
769 let parent_span = Span::current();
770 let mut hashed_state: LazyHashedPostState =
771 self.runtime.spawn_blocking_named("hash-post-state", move || {
772 let _span = debug_span!(
773 target: "engine::tree::payload_validator",
774 parent: parent_span,
775 "hashed_post_state",
776 )
777 .entered();
778 if let Some(Ok(state)) = hashed_state_rx.as_mut().map(|rx| rx.recv()) {
779 state
780 } else {
781 Arc::new(HashedPostState::from_bundle_state::<KeccakKeyHasher>(
782 hashed_state_output.state.state(),
783 ))
784 }
785 });
786
787 let block = validated_block.try_into_inner().expect("sole handle")?;
788 let block = block.with_senders(senders);
789
790 let receipt_root_bloom = {
792 let _enter = debug_span!(
793 target: "engine::tree::payload_validator",
794 "wait_receipt_root",
795 )
796 .entered();
797
798 receipt_root_rx
799 .blocking_recv()
800 .inspect_err(|_| {
801 tracing::error!(
802 target: "engine::tree::payload_validator",
803 "Receipt root task dropped sender without result, receipt root calculation likely aborted"
804 );
805 })
806 .ok()
807 };
808
809 ensure_ok_post_block!(
810 self.validate_post_execution(
811 &block,
812 &parent_block,
813 &output,
814 &mut ctx,
815 receipt_root_bloom,
816 built_bal
817 ),
818 block
819 );
820
821 let mut hashed_state_validate_result = debug_span!(
822 target: "engine::tree::payload_validator",
823 "validate_block_post_execution_with_hashed_state"
824 )
825 .in_scope(|| {
826 self.validator.validate_block_post_execution_with_hashed_state(
827 || hashed_state.get(),
828 &block,
829 &parent_block,
830 || {
831 state_provider_factory
832 .database_provider_ro()
833 .map(|provider| Box::new(provider) as _)
834 },
835 )
836 });
837
838 let root_start = Instant::now();
839 let root_outcome = ensure_ok_post_block!(
840 state_root_job.finish(&block, output.clone(), &hashed_state),
841 block
842 );
843 let root_elapsed = root_start.elapsed();
844
845 info!(
846 target: "engine::tree::payload_validator",
847 strategy = state_root_job_name,
848 state_root = ?root_outcome.state_root,
849 elapsed = ?root_elapsed,
850 "State root job finished"
851 );
852
853 let state_root = root_outcome.state_root;
854 let trie_output = root_outcome.trie_updates;
855
856 if let Some(refreshed) = root_outcome.hashed_state {
860 hashed_state = LazyHandle::ready(refreshed);
861 hashed_state_validate_result = debug_span!(
862 target: "engine::tree::payload_validator",
863 "validate_block_post_execution_with_hashed_state"
864 )
865 .in_scope(|| {
866 self.validator.validate_block_post_execution_with_hashed_state(
867 || hashed_state.get(),
868 &block,
869 &parent_block,
870 || {
871 state_provider_factory
872 .database_provider_ro()
873 .map(|provider| Box::new(provider) as _)
874 },
875 )
876 });
877 }
878
879 if let Err(err) = hashed_state_validate_result {
880 if err.is_validation_error() {
881 self.on_invalid_block(&parent_block, &block, &output, None, ctx.state_mut());
882 }
883 return Err(InsertBlockError::new(block.into_sealed_block(), err).into())
884 }
885
886 self.metrics.block_validation.record_state_root(&trie_output, root_elapsed.as_secs_f64());
887 self.metrics
888 .record_state_root_gas_bucket(block.header().gas_used(), root_elapsed.as_secs_f64());
889 debug!(target: "engine::tree::payload_validator", ?root_elapsed, "Calculated state root");
890
891 if state_root != block.header().state_root() {
893 self.on_invalid_block(
895 &parent_block,
896 &block,
897 &output,
898 Some((&trie_output.as_ref().clone().into(), state_root)),
899 ctx.state_mut(),
900 );
901 let block_state_root = block.header().state_root();
902 return Err(InsertBlockError::new(
903 block.into_sealed_block(),
904 ConsensusError::BodyStateRootDiff(
905 GotExpected { got: state_root, expected: block_state_root }.into(),
906 )
907 .into(),
908 )
909 .into())
910 }
911
912 let timing_stats = state_provider_stats.filter(|_| slow_block_enabled).map(|stats| {
913 self.calculate_timing_stats(
914 &block,
915 stats,
916 cache_stats,
917 &output,
918 execution_duration,
919 root_elapsed,
920 )
921 });
922
923 if let Some(valid_block_tx) = valid_block_tx {
924 let _ = valid_block_tx.send(());
925 }
926
927 let bal = revm_bal.zip(decoded_bal).map(|(revm_bal, decoded_bal)| {
932 Arc::new(DecodedRevmBal::with_raw_bal(revm_bal, decoded_bal.as_raw_bal().clone()))
933 });
934 let executed_block = self
935 .spawn_deferred_hashed_state_task(Arc::new(block), output, hashed_state, trie_output)
936 .with_bal(bal);
937 Ok(ValidationOutput::new(executed_block, timing_stats))
938 }
939
940 #[expect(clippy::type_complexity)]
943 pub fn spawn_convert_and_validate<T>(
944 &self,
945 input: &BlockOrPayload<T>,
946 parent: SealedHeader<N::BlockHeader>,
947 ) -> LazyHandle<Result<SealedBlock<N::Block>, InsertPayloadError<N::Block>>>
948 where
949 T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>,
950 V: PayloadValidator<T, Block = N::Block> + Clone,
951 {
952 let input = input.clone();
953 let validator = self.validator.clone();
954 let consensus = self.consensus.clone();
955 let parent_span = Span::current();
956 self.runtime.spawn_blocking_named("payload-convert", move || {
957 let _span = debug_span!(
958 target: "engine::tree::payload_validator",
959 parent: parent_span,
960 "convert_and_validate",
961 )
962 .entered();
963 let block = match input {
964 BlockOrPayload::Block(block) => block.split().0,
965 BlockOrPayload::Payload(payload) => {
966 validator.convert_payload_to_block(payload)?
967 }
968 };
969
970 if let Err(e) = consensus.validate_header(block.sealed_header()) {
971 error!(target: "engine::tree::payload_validator", ?block, "Failed to validate header {}: {e}", block.hash());
972 return Err(InsertBlockError::consensus_error(e, block).into())
973 }
974
975 let _enter = debug_span!(target: "engine::tree::payload_validator", "validate_header_against_parent").entered();
977 if let Err(e) = consensus.validate_header_against_parent(block.sealed_header(), &parent)
978 {
979 warn!(target: "engine::tree::payload_validator", ?block, "Failed to validate header {} against parent: {e}", block.hash());
980 return Err(InsertBlockError::consensus_error(e, block).into())
981 }
982 drop(_enter);
983
984 if let Err(e) =
985 consensus.validate_block_pre_execution_with_tx_root(&block, None)
986 {
987 error!(target: "engine::tree::payload_validator", ?block, "Failed to validate block {}: {e}", block.hash());
988 return Err(InsertBlockError::consensus_error(e, block).into())
989 }
990
991 Ok(block)
992 })
993 }
994
995 fn sealed_header_by_hash(
997 &self,
998 hash: B256,
999 state: &EngineApiTreeState<N>,
1000 ) -> ProviderResult<Option<SealedHeader<N::BlockHeader>>> {
1001 let header = state.tree_state.sealed_header_by_hash(&hash);
1003
1004 if header.is_some() {
1005 Ok(header)
1006 } else {
1007 self.provider.sealed_header_by_hash(hash)
1008 }
1009 }
1010
1011 #[instrument(level = "debug", target = "engine::tree::payload_validator", skip_all)]
1019 #[expect(clippy::type_complexity)]
1020 fn execute_block<S, Err, T>(
1021 &mut self,
1022 state_provider: S,
1023 env: ExecutionEnv<Evm>,
1024 input: &BlockOrPayload<T>,
1025 handle: &mut PayloadHandle<impl ExecutableTxFor<Evm>, Err, N::Receipt>,
1026 state_hook: Option<Box<dyn OnStateHook + 'static>>,
1027 ) -> Result<
1028 (BlockExecutionOutput<N::Receipt>, Vec<Address>, ReceiptRootReceiver, Option<ExecutedBal>),
1029 InsertBlockErrorKind,
1030 >
1031 where
1032 S: EvmStateProvider + Send,
1033 Err: core::error::Error + Send + Sync + 'static,
1034 V: PayloadValidator<T, Block = N::Block>,
1035 T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>,
1036 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
1037 {
1038 debug!(target: "engine::tree::payload_validator", "Executing block");
1039
1040 let has_bal = input.has_block_access_list();
1041 let mut db = debug_span!(target: "engine::tree", "build_state_db").in_scope(|| {
1042 State::builder()
1043 .with_database(StateProviderDatabase::new(state_provider))
1044 .with_bundle_update()
1045 .with_bal_builder_if(has_bal)
1046 .build()
1047 });
1048
1049 let (spec_id, mut executor) = {
1050 let _span = debug_span!(target: "engine::tree", "create_evm").entered();
1051 let spec_id = *env.evm_env.spec_id();
1052 let evm_config = self.evm_config.clone().with_jit_support();
1053 let evm = evm_config.evm_with_env(&mut db, env.evm_env);
1054 let ctx = self
1055 .execution_ctx_for(input)
1056 .map_err(|e| InsertBlockErrorKind::Other(Box::new(e)))?;
1057 let executor = self.evm_config.create_executor(evm, ctx);
1058 (spec_id, executor)
1059 };
1060
1061 if !self.config.precompile_cache_disabled() {
1062 let _span = debug_span!(target: "engine::tree", "setup_precompile_cache").entered();
1063 executor.evm_mut().precompiles_mut().map_cacheable_precompiles(
1064 |address, precompile| {
1065 let metrics = self
1066 .precompile_cache_metrics
1067 .entry(*address)
1068 .or_insert_with(|| CachedPrecompileMetrics::new_with_address(*address))
1069 .clone();
1070 CachedPrecompile::wrap(
1071 precompile,
1072 self.precompile_cache_map.cache_for_address(*address),
1073 spec_id,
1074 Some(metrics),
1075 )
1076 },
1077 );
1078 }
1079
1080 let transaction_count = input.transaction_count();
1081 let (receipt_tx, result_rx) = self.spawn_receipt_root_task(transaction_count);
1082 let executed_tx_index = Arc::clone(handle.executed_tx_index());
1083 executor.evm_mut().db_mut().set_state_hook(state_hook);
1084
1085 let execution_start = Instant::now();
1086
1087 let (executor, senders) = self.execute_transactions(
1089 executor,
1090 transaction_count,
1091 handle.iter_transactions(),
1092 &receipt_tx,
1093 &executed_tx_index,
1094 has_bal,
1095 )?;
1096 drop(receipt_tx);
1097
1098 let post_exec_start = Instant::now();
1100 let (_evm, result) = debug_span!(target: "engine::tree", "BlockExecutor::finish")
1101 .in_scope(|| executor.finish())
1102 .map(|(evm, result)| (evm.into_db(), result))?;
1103 self.metrics.record_post_execution(post_exec_start.elapsed());
1104
1105 debug_span!(target: "engine::tree", "merge_transitions")
1107 .in_scope(|| db.merge_transitions(BundleRetention::Reverts));
1108
1109 let built_bal = db.take_built_bal().map(|revm_bal| ExecutedBal {
1113 alloy: revm_bal.clone().into_alloy_bal(),
1114 revm: Arc::new(revm_bal),
1115 });
1116 let output = BlockExecutionOutput { result, state: db.take_bundle() };
1117
1118 let execution_duration = execution_start.elapsed();
1119 self.metrics.record_block_execution(&output, execution_duration);
1120 self.metrics.record_block_execution_gas_bucket(output.result.gas_used, execution_duration);
1121 debug!(target: "engine::tree::payload_validator", elapsed = ?execution_duration, "Executed block");
1122
1123 Ok((output, senders, result_rx, built_bal))
1124 }
1125
1126 fn bal_path_eligible(&self, bal: Option<&DecodedBal>) -> Result<bool, InsertBlockErrorKind> {
1135 let has_bal = bal.is_some();
1136 let parallel_execution = has_bal && !self.config.disable_bal_parallel_execution();
1137 if parallel_execution && self.config.disable_bal_parallel_state_root() {
1138 return Err(InsertBlockErrorKind::Other(
1139 "disabling parallel state root is impossible when parallel execution is enabled"
1140 .into(),
1141 ));
1142 }
1143
1144 Ok(parallel_execution)
1145 }
1146
1147 #[instrument(level = "debug", target = "engine::tree::payload_validator", skip_all)]
1158 #[expect(clippy::type_complexity)]
1159 fn execute_block_bal<Tx, Err, MakeStateProvider, T>(
1160 &self,
1161 env: ExecutionEnv<Evm>,
1162 input: &BlockOrPayload<T>,
1163 handle: &PayloadHandle<Tx, Err, N::Receipt>,
1164 make_state_provider: &MakeStateProvider,
1165 ) -> Result<
1166 (BlockExecutionOutput<N::Receipt>, Vec<Address>, ReceiptRootReceiver, Option<ExecutedBal>),
1167 InsertBlockErrorKind,
1168 >
1169 where
1170 Tx: ExecutableTxFor<Evm> + Send,
1171 Err: core::error::Error + Send + Sync + 'static,
1172 MakeStateProvider: Fn(bool) -> ProviderResult<EvmStateProviderBox> + Sync,
1173 Evm: ConfigureEngineEvm<T::ExecutionData, Primitives = N>,
1174 T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>,
1175 V: PayloadValidator<T, Block = N::Block>,
1176 {
1177 debug!(target: "engine::tree::payload_validator", "Executing block via BAL path");
1178
1179 let (receipt_tx, result_rx) = self.spawn_receipt_root_task(env.transaction_count);
1180 let input_bal = env.decoded_bal.ok_or_else(|| {
1181 InsertBlockErrorKind::Other("BAL execute path: no decoded BAL available".into())
1182 })?;
1183
1184 let make_db = |fill_on_miss| {
1185 let provider = make_state_provider(fill_on_miss)
1186 .map_err(crate::tree::payload_processor::bal::BalExecutionError::Provider)?;
1187 Ok(StateProviderDatabase::new(provider))
1188 };
1189 let execution_start = Instant::now();
1190 let ctx =
1191 self.execution_ctx_for(input).map_err(|e| InsertBlockErrorKind::Other(Box::new(e)))?;
1192 let (output, senders, built_bal, received_bal_revm) =
1193 crate::tree::payload_processor::bal::execute_block(
1194 &self.runtime,
1195 &self.evm_config,
1196 &make_db,
1197 input_bal,
1198 env.evm_env,
1199 ctx,
1200 env.transaction_count,
1201 handle.clone_transaction_receiver(),
1202 receipt_tx,
1203 )?;
1204 let execution_duration = execution_start.elapsed();
1205
1206 self.metrics.record_block_execution(&output, execution_duration);
1207 self.metrics.record_block_execution_gas_bucket(output.result.gas_used, execution_duration);
1208 debug!(
1209 target: "engine::tree::payload_validator",
1210 elapsed = ?execution_duration,
1211 "Executed block via BAL path",
1212 );
1213
1214 Ok((
1218 output,
1219 senders,
1220 result_rx,
1221 Some(ExecutedBal { alloy: built_bal, revm: received_bal_revm }),
1222 ))
1223 }
1224
1225 fn spawn_receipt_root_task(
1226 &self,
1227 receipts_len: usize,
1228 ) -> (ReceiptRootSender<N>, ReceiptRootReceiver) {
1229 let (receipt_tx, receipt_rx) = crossbeam_channel::unbounded();
1231 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1232 let task_handle = ReceiptRootTaskHandle::new(receipt_rx, result_tx);
1233 self.runtime.spawn_blocking_named("receipt-root", move || task_handle.run(receipts_len));
1234
1235 (receipt_tx, result_rx)
1236 }
1237
1238 fn execute_transactions<'a, E, Tx, InnerTx, Err, DB>(
1248 &self,
1249 mut executor: E,
1250 transaction_count: usize,
1251 transactions: impl Iterator<Item = Result<Tx, Err>>,
1252 receipt_tx: &crossbeam_channel::Sender<IndexedReceipt<N::Receipt>>,
1253 executed_tx_index: &AtomicUsize,
1254 has_bal: bool,
1255 ) -> Result<(E, Vec<Address>), BlockExecutionError>
1256 where
1257 E: BlockExecutor<Receipt = N::Receipt, Evm: alloy_evm::Evm<DB = &'a mut State<DB>>>,
1258 Tx: alloy_evm::block::ExecutableTx<E> + alloy_evm::RecoveredTx<InnerTx>,
1259 InnerTx: TxHashRef,
1260 DB: revm::Database + 'a,
1261 Err: core::error::Error + Send + Sync + 'static,
1262 {
1263 let mut senders = Vec::with_capacity(transaction_count);
1264
1265 let pre_exec_start = Instant::now();
1267 debug_span!(target: "engine::tree", "pre_execution")
1268 .in_scope(|| executor.apply_pre_execution_changes())?;
1269 self.metrics.record_pre_execution(pre_exec_start.elapsed());
1270
1271 if has_bal {
1273 executor.evm_mut().db_mut().bump_bal_index();
1274 }
1275
1276 let exec_span = debug_span!(target: "engine::tree", "execution").entered();
1278 let mut transactions = transactions.into_iter();
1279 let mut last_sent_len = 0usize;
1284 loop {
1285 let wait_start = Instant::now();
1288 let Some(tx_result) = transactions.next() else { break };
1289 self.metrics.record_transaction_wait(wait_start.elapsed());
1290
1291 let tx = tx_result.map_err(BlockValidationError::other)?;
1292 let tx_signer = *<Tx as alloy_evm::RecoveredTx<InnerTx>>::signer(&tx);
1293
1294 senders.push(tx_signer);
1295
1296 let _enter = tracing::enabled!(target: "engine::tree", Level::TRACE).then(|| {
1297 tracing::trace_span!(
1298 target: "engine::tree",
1299 "execute tx",
1300 tx_index = senders.len() - 1,
1301 )
1302 .entered()
1303 });
1304 if tracing::enabled!(target: "engine::tree", Level::TRACE) {
1305 trace!(target: "engine::tree", "Executing transaction");
1306 }
1307
1308 let tx_start = Instant::now();
1309 executor.execute_transaction(tx)?;
1310 self.metrics.record_transaction_execution(tx_start.elapsed());
1311
1312 executed_tx_index.store(senders.len(), Ordering::Relaxed);
1314
1315 let current_len = executor.receipts().len();
1316 if current_len > last_sent_len {
1317 last_sent_len = current_len;
1318 if let Some(receipt) = executor.receipts().last() {
1320 let tx_index = current_len - 1;
1321 let _ = receipt_tx.send(IndexedReceipt::new(tx_index, receipt.clone()));
1322 }
1323 }
1324 if has_bal {
1326 executor.evm_mut().db_mut().bump_bal_index();
1327 }
1328 }
1329
1330 drop(exec_span);
1331
1332 Ok((executor, senders))
1333 }
1334
1335 #[instrument(level = "debug", target = "engine::tree::payload_validator", skip_all)]
1347 fn validate_post_execution<T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>>(
1348 &mut self,
1349 block: &RecoveredBlock<N::Block>,
1350 parent_block: &SealedHeader<N::BlockHeader>,
1351 output: &BlockExecutionOutput<N::Receipt>,
1352 ctx: &mut TreeCtx<'_, N>,
1353 receipt_root_bloom: Option<ReceiptRootBloom>,
1354 built_bal: Option<BlockAccessList>,
1355 ) -> Result<(), InsertBlockErrorKind>
1356 where
1357 V: PayloadValidator<T, Block = N::Block>,
1358 {
1359 let start = Instant::now();
1360
1361 trace!(target: "engine::tree::payload_validator", block=?block.num_hash(), "Validating block consensus");
1362
1363 let _enter =
1365 debug_span!(target: "engine::tree::payload_validator", "validate_block_post_execution")
1366 .entered();
1367 let built_bal = built_bal.map(Bal::from);
1368 let block_access_list_hash =
1369 built_bal.as_ref().map(|bal| bal.compute_hash_with_buf(&mut self.bal_hash_buf));
1370
1371 let validation_result = built_bal
1372 .as_ref()
1373 .map(|bal| bal.validate_gas_limit(block.gas_limit()).map_err(ConsensusError::from))
1374 .transpose()
1375 .and_then(|_| {
1376 self.consensus.validate_block_post_execution(
1377 block,
1378 output,
1379 receipt_root_bloom,
1380 block_access_list_hash,
1381 )
1382 });
1383
1384 if let Err(err) = validation_result {
1385 self.on_invalid_block(parent_block, block, output, None, ctx.state_mut());
1387 return Err(err.into())
1388 }
1389 drop(_enter);
1390
1391 self.metrics
1393 .block_validation
1394 .post_execution_validation_duration
1395 .record(start.elapsed().as_secs_f64());
1396
1397 Ok(())
1398 }
1399
1400 #[instrument(
1405 level = "debug",
1406 target = "engine::tree::payload_validator",
1407 skip_all,
1408 fields(
1409 has_hint_stream = hint_stream.is_some(),
1410 has_hashed_update_stream = hashed_update_stream.is_some(),
1411 parallel_bal_execution
1412 )
1413 )]
1414 fn spawn_payload_processor<T: ExecutableTxIterator<Evm>>(
1415 &self,
1416 env: ExecutionEnv<Evm>,
1417 txs: T,
1418 state_provider_factory: OverlayStateProviderFactory<P, N>,
1419 hint_stream: Option<StateRootHintStream>,
1420 hashed_update_stream: Option<StateRootUpdateStream>,
1421 parallel_bal_execution: bool,
1422 ) -> Result<
1423 PayloadHandle<
1424 impl ExecutableTxFor<Evm> + use<N, P, Evm, V, T>,
1425 impl core::error::Error + Send + Sync + 'static + use<N, P, Evm, V, T>,
1426 N::Receipt,
1427 >,
1428 InsertBlockErrorKind,
1429 > {
1430 let start = Instant::now();
1431 let handle = self.payload_processor.spawn_with_state_root_streams(
1432 env,
1433 txs,
1434 state_provider_factory,
1435 hint_stream,
1436 hashed_update_stream,
1437 parallel_bal_execution,
1438 );
1439
1440 self.metrics.block_validation.spawn_payload_processor.record(start.elapsed().as_secs_f64());
1441
1442 Ok(handle)
1443 }
1444
1445 fn overlay_state_provider_factory(
1449 &self,
1450 hash: B256,
1451 state: &EngineApiTreeState<N>,
1452 ) -> ProviderResult<Option<OverlayStateProviderFactory<P, N>>> {
1453 if !state.tree_state.contains_hash(&hash) && self.provider.header(hash)?.is_none() {
1454 debug!(target: "engine::tree::payload_validator", %hash, "no canonical state found for block");
1455 return Ok(None)
1456 }
1457
1458 Ok(Some(OverlayStateProviderFactory::new(
1459 self.provider.clone(),
1460 state.tree_state.overlay_manager.overlay_builder(hash),
1461 )))
1462 }
1463
1464 fn on_invalid_block(
1466 &self,
1467 parent_header: &SealedHeader<N::BlockHeader>,
1468 block: &RecoveredBlock<N::Block>,
1469 output: &BlockExecutionOutput<N::Receipt>,
1470 trie_updates: Option<(&TrieUpdates, B256)>,
1471 state: &mut EngineApiTreeState<N>,
1472 ) {
1473 if state.invalid_headers.get(&block.hash()).is_some() {
1474 return
1476 }
1477 self.invalid_block_hook.on_invalid_block(parent_header, block, output, trie_updates);
1478 }
1479
1480 fn payload_state_root_handle_for(
1483 &self,
1484 parent_hash: B256,
1485 parent_header: &N::BlockHeader,
1486 timestamp: u64,
1487 state: &mut EngineApiTreeState<N>,
1488 ) -> Option<PayloadStateRootHandle> {
1489 let state_provider_factory = match self.overlay_state_provider_factory(parent_hash, state) {
1490 Ok(Some(state_provider_factory)) => state_provider_factory,
1491 Ok(None) => return None,
1492 Err(err) => {
1493 warn!(
1494 target: "engine::tree::payload_validator",
1495 %err,
1496 %parent_hash,
1497 "failed to prepare payload-builder state-root provider"
1498 );
1499 return None
1500 }
1501 };
1502 match self.state_root_strategy.prepare_payload_builder(PayloadStateRootJobContext::new(
1503 &self.runtime,
1504 &self.overlay_manager,
1505 parent_hash,
1506 parent_header,
1507 timestamp,
1508 state,
1509 state_provider_factory,
1510 &self.config,
1511 )) {
1512 Ok(handle) => handle,
1513 Err(err) => {
1514 warn!(
1515 target: "engine::tree::payload_validator",
1516 %err,
1517 %parent_hash,
1518 "failed to prepare payload-builder state-root job"
1519 );
1520 None
1521 }
1522 }
1523 }
1524
1525 fn spawn_deferred_hashed_state_task(
1527 &self,
1528 block: Arc<RecoveredBlock<N::Block>>,
1529 execution_outcome: Arc<BlockExecutionOutput<N::Receipt>>,
1530 hashed_state: LazyHashedPostState,
1531 trie_output: Arc<TrieUpdatesSorted>,
1532 ) -> ExecutedBlock<N> {
1533 let hashed_state = match hashed_state.try_into_inner() {
1536 Ok(state) => state,
1537 Err(handle) => handle.get().clone(),
1538 };
1539 let (deferred_hashed_state, producer) = LazyHashedPostStateSorted::pending(hashed_state);
1540 self.metrics
1541 .block_validation
1542 .trie_updates_sorted_size
1543 .record(trie_output.total_len() as f64);
1544 let block_validation_metrics = self.metrics.block_validation.clone();
1545
1546 let block_number = block.number();
1548
1549 let sort_hashed_state_task = move || {
1551 let _span = debug_span!(
1552 target: "engine::tree::payload_validator",
1553 "sort_hashed_state_task",
1554 block_number
1555 )
1556 .entered();
1557
1558 let compute_start = Instant::now();
1559 let computed = producer.compute_and_publish();
1560 block_validation_metrics
1561 .deferred_trie_compute_duration
1562 .record(compute_start.elapsed().as_secs_f64());
1563
1564 block_validation_metrics.hashed_post_state_size.record(computed.total_len() as f64);
1566 };
1567
1568 self.runtime.spawn_blocking_named(DEFERRED_TRIE_WORKER_NAME, sort_hashed_state_task);
1569
1570 ExecutedBlock::with_deferred_hashed_state(
1571 block,
1572 execution_outcome,
1573 deferred_hashed_state,
1574 trie_output,
1575 )
1576 }
1577
1578 fn calculate_timing_stats(
1579 &self,
1580 block: &RecoveredBlock<N::Block>,
1581 provider_stats: Arc<StateProviderStats>,
1582 cache_stats: Option<Arc<CacheStats>>,
1583 output: &BlockExecutionOutput<N::Receipt>,
1584 execution_duration: Duration,
1585 state_hash_duration: Duration,
1586 ) -> Box<ExecutionTimingStats> {
1587 let accounts_read = provider_stats.total_account_fetches();
1588 let storage_read = provider_stats.total_storage_fetches();
1589 let code_read = provider_stats.total_code_fetches();
1590 let code_bytes_read = provider_stats.total_code_fetched_bytes();
1591
1592 let accounts_changed = output.state.len();
1594 let accounts_deleted =
1595 output.state.state.values().filter(|acc| acc.was_destroyed()).count();
1596 let storage_slots_changed =
1597 output.state.state.values().map(|account| account.storage.len()).sum::<usize>();
1598 let storage_slots_deleted = output
1599 .state
1600 .state
1601 .values()
1602 .flat_map(|account| account.storage.values())
1603 .filter(|slot| {
1604 slot.present_value.is_zero() && !slot.previous_or_original_value.is_zero()
1605 })
1606 .count();
1607
1608 let is_new_deployment = |acc: &BundleAccount| -> bool {
1610 let has_code_now = acc.info.as_ref().is_some_and(|info| !info.is_empty_code_hash());
1611 let had_no_code_before =
1612 acc.original_info.as_ref().is_none_or(AccountInfo::is_empty_code_hash);
1613 has_code_now && had_no_code_before
1614 };
1615
1616 let bytecodes_changed =
1617 output.state.state.values().filter(|acc| is_new_deployment(acc)).count();
1618
1619 let unique_new_code_hashes: B256Set = output
1621 .state
1622 .state
1623 .values()
1624 .filter(|acc| is_new_deployment(acc))
1625 .filter_map(|acc| acc.info.as_ref().map(|info| info.code_hash))
1626 .collect();
1627 let code_bytes_written: usize = unique_new_code_hashes
1628 .iter()
1629 .filter_map(|hash| output.state.contracts.get(hash).map(|bytecode| bytecode.len()))
1630 .sum();
1631
1632 let state_read_duration = provider_stats.total_account_fetch_latency() +
1634 provider_stats.total_storage_fetch_latency() +
1635 provider_stats.total_code_fetch_latency();
1636
1637 let eip7702_delegations_set =
1640 output.state.contracts.values().filter(|bytecode| bytecode.is_eip7702()).count();
1641 let eip7702_delegations_cleared = output
1646 .state
1647 .state
1648 .values()
1649 .filter(|acc| {
1650 let original_was_eip7702 = acc
1652 .original_info
1653 .as_ref()
1654 .and_then(|info| info.code.as_ref())
1655 .map(|bytecode| bytecode.is_eip7702())
1656 .unwrap_or(false);
1657
1658 let code_now_empty = acc.info.as_ref().is_some_and(AccountInfo::is_empty_code_hash);
1660
1661 original_was_eip7702 && code_now_empty
1662 })
1663 .count();
1664
1665 let (account_cache_hits, account_cache_misses) = cache_stats
1667 .as_ref()
1668 .map(|s| (s.account_hits(), s.account_misses()))
1669 .unwrap_or_default();
1670 let (storage_cache_hits, storage_cache_misses) = cache_stats
1671 .as_ref()
1672 .map(|s| (s.storage_hits(), s.storage_misses()))
1673 .unwrap_or_default();
1674 let (code_cache_hits, code_cache_misses) =
1675 cache_stats.as_ref().map(|s| (s.code_hits(), s.code_misses())).unwrap_or_default();
1676 let (txpool_snapshot_account_hits, txpool_snapshot_account_misses) = cache_stats
1677 .as_ref()
1678 .map(|s| (s.txpool_snapshot_account_hits(), s.txpool_snapshot_account_misses()))
1679 .unwrap_or_default();
1680 let (txpool_snapshot_storage_hits, txpool_snapshot_storage_misses) = cache_stats
1681 .as_ref()
1682 .map(|s| (s.txpool_snapshot_storage_hits(), s.txpool_snapshot_storage_misses()))
1683 .unwrap_or_default();
1684 let (txpool_snapshot_code_hits, txpool_snapshot_code_misses) = cache_stats
1685 .as_ref()
1686 .map(|s| (s.txpool_snapshot_code_hits(), s.txpool_snapshot_code_misses()))
1687 .unwrap_or_default();
1688
1689 Box::new(ExecutionTimingStats {
1691 block_number: block.number(),
1692 block_hash: block.hash(),
1693 gas_used: output.result.gas_used,
1694 tx_count: block.transaction_count(),
1695 execution_duration,
1696 state_read_duration,
1697 state_hash_duration,
1698 accounts_read,
1699 storage_read,
1700 code_read,
1701 code_bytes_read,
1702 accounts_changed,
1703 accounts_deleted,
1704 storage_slots_changed,
1705 storage_slots_deleted,
1706 bytecodes_changed,
1707 code_bytes_written,
1708 eip7702_delegations_set,
1709 eip7702_delegations_cleared,
1710 account_cache_hits,
1711 account_cache_misses,
1712 storage_cache_hits,
1713 storage_cache_misses,
1714 code_cache_hits,
1715 code_cache_misses,
1716 txpool_snapshot_account_hits,
1717 txpool_snapshot_account_misses,
1718 txpool_snapshot_storage_hits,
1719 txpool_snapshot_storage_misses,
1720 txpool_snapshot_code_hits,
1721 txpool_snapshot_code_misses,
1722 })
1723 }
1724}
1725
1726pub trait EngineValidator<
1730 Types: PayloadTypes,
1731 N: NodePrimitives = <<Types as PayloadTypes>::BuiltPayload as BuiltPayload>::Primitives,
1732>: Send + Sync + 'static
1733{
1734 fn validate_payload_attributes_against_header(
1744 &self,
1745 attr: &Types::PayloadAttributes,
1746 header: &N::BlockHeader,
1747 ) -> Result<(), InvalidPayloadAttributesError>;
1748
1749 fn convert_payload_to_block(
1758 &self,
1759 payload: Types::ExecutionData,
1760 ) -> Result<SealedBlock<N::Block>, NewPayloadError>;
1761
1762 fn validate_payload(
1764 &mut self,
1765 payload: Types::ExecutionData,
1766 ctx: TreeCtx<'_, N>,
1767 ) -> ValidationOutcome<N>;
1768
1769 fn validate_block(
1771 &mut self,
1772 block: SealedBlockWithAccessList<N::Block>,
1773 ctx: TreeCtx<'_, N>,
1774 ) -> ValidationOutcome<N>;
1775
1776 fn on_inserted_executed_block(
1781 &self,
1782 block: BuiltPayloadExecutedBlock<N>,
1783 ) -> ProviderResult<ExecutedBlock<N>>;
1784
1785 fn on_canonical_head_changed(&self, _hash: B256, _state: &EngineApiTreeState<N>) {}
1789
1790 fn payload_builder_resources(
1794 &self,
1795 parent_hash: B256,
1796 parent_header: &N::BlockHeader,
1797 timestamp: u64,
1798 state: &mut EngineApiTreeState<N>,
1799 ) -> PayloadBuilderResources;
1800}
1801
1802impl<N, Types, P, Evm, V> EngineValidator<Types> for BasicEngineValidator<P, Evm, V>
1803where
1804 P: DatabaseProviderFactory<
1805 Provider: BlockReader
1806 + BlockHashReader
1807 + StageCheckpointReader
1808 + PruneCheckpointReader
1809 + ChangeSetReader
1810 + StorageChangeSetReader
1811 + StorageSettingsCache
1812 + HistoryReader
1813 + 'static,
1814 > + BlockReader<Header = N::BlockHeader>
1815 + StateProviderFactory
1816 + StateReader
1817 + ChangeSetReader
1818 + Clone
1819 + 'static,
1820 OverlayStateProviderFactory<P, N>: DatabaseProviderROFactory<
1821 Provider: TrieCursorFactory
1822 + HashedCursorFactory
1823 + HashedPostStateProvider
1824 + StateRootProvider
1825 + StateProvider
1826 + Send,
1827 > + Clone
1828 + 'static,
1829 N: NodePrimitives,
1830 V: PayloadValidator<Types, Block = N::Block> + Clone,
1831 Evm: ConfigureEngineEvm<Types::ExecutionData, Primitives = N> + 'static,
1832 Types: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = N>>,
1833{
1834 fn validate_payload_attributes_against_header(
1835 &self,
1836 attr: &Types::PayloadAttributes,
1837 header: &N::BlockHeader,
1838 ) -> Result<(), InvalidPayloadAttributesError> {
1839 self.validator.validate_payload_attributes_against_header(attr, header)
1840 }
1841
1842 fn convert_payload_to_block(
1843 &self,
1844 payload: Types::ExecutionData,
1845 ) -> Result<SealedBlock<N::Block>, NewPayloadError> {
1846 let block = self.validator.convert_payload_to_block(payload)?;
1847 Ok(block)
1848 }
1849
1850 fn validate_payload(
1851 &mut self,
1852 payload: Types::ExecutionData,
1853 ctx: TreeCtx<'_, N>,
1854 ) -> ValidationOutcome<N> {
1855 self.validate_block_with_state(BlockOrPayload::Payload(payload), ctx)
1856 }
1857
1858 fn validate_block(
1859 &mut self,
1860 block: SealedBlockWithAccessList<N::Block>,
1861 ctx: TreeCtx<'_, N>,
1862 ) -> ValidationOutcome<N> {
1863 self.validate_block_with_state(BlockOrPayload::Block(block), ctx)
1864 }
1865
1866 fn on_inserted_executed_block(
1867 &self,
1868 block: BuiltPayloadExecutedBlock<N>,
1869 ) -> ProviderResult<ExecutedBlock<N>> {
1870 self.payload_processor.on_inserted_executed_block(
1871 block.recovered_block.block_with_parent(),
1872 &block.execution_output.state,
1873 );
1874
1875 Ok(self.spawn_deferred_hashed_state_task(
1876 block.recovered_block,
1877 block.execution_output,
1878 LazyHashedPostState::ready(block.hashed_state),
1879 block.trie_updates,
1880 ))
1881 }
1882
1883 fn on_canonical_head_changed(&self, hash: B256, state: &EngineApiTreeState<N>) {
1884 let Some(txpool_prewarm) = self.txpool_prewarm.as_ref() else { return };
1885
1886 let parent = match self.sealed_header_by_hash(hash, state) {
1889 Ok(Some(header)) => header,
1890 Ok(None) => return,
1891 Err(err) => {
1892 trace!(
1893 target: "engine::tree::txpool_prewarm",
1894 %err,
1895 block_hash = ?hash,
1896 "failed to fetch canonical header for txpool prewarming"
1897 );
1898 return
1899 }
1900 };
1901 let evm_env = match self.evm_config.evm_env(parent.header()) {
1905 Ok(evm_env) => evm_env,
1906 Err(err) => {
1907 trace!(
1908 target: "engine::tree::txpool_prewarm",
1909 %err,
1910 block_hash = ?parent.hash(),
1911 "failed to derive canonical txpool prewarming environment"
1912 );
1913 return
1914 }
1915 };
1916
1917 let state_provider_factory = match self.overlay_state_provider_factory(parent.hash(), state)
1918 {
1919 Ok(Some(state_provider_factory)) => state_provider_factory,
1920 Ok(None) => return,
1921 Err(err) => {
1922 trace!(
1923 target: "engine::tree::txpool_prewarm",
1924 %err,
1925 block_hash = ?parent.hash(),
1926 "failed to derive canonical txpool prewarming provider"
1927 );
1928 return
1929 }
1930 };
1931 txpool_prewarm.start(parent.hash(), evm_env, state_provider_factory)
1932 }
1933
1934 fn payload_builder_resources(
1935 &self,
1936 parent_hash: B256,
1937 parent_header: &N::BlockHeader,
1938 timestamp: u64,
1939 state: &mut EngineApiTreeState<N>,
1940 ) -> PayloadBuilderResources {
1941 let execution_cache = self
1942 .config
1943 .share_execution_cache_with_payload_builder()
1944 .then(|| self.payload_processor.cache_for(parent_hash));
1945 let state_root_handle =
1946 self.payload_state_root_handle_for(parent_hash, parent_header, timestamp, state);
1947 let mut resources = PayloadBuilderResources::new(execution_cache, state_root_handle)
1948 .with_lease(PayloadBuilderLease::new(JitPauseGuard::new(&self.evm_config)));
1949 if let Some(txpool_prewarm) = self.txpool_prewarm.as_ref() {
1953 let txpool_lease = PayloadBuilderLease::new(txpool_prewarm.pause());
1954 resources = resources.with_lease(txpool_lease);
1955 }
1956 resources
1957 }
1958}
1959
1960impl<P, Evm, V> WaitForCaches for BasicEngineValidator<P, Evm, V>
1961where
1962 Evm: ConfigureEvm,
1963{
1964 fn wait_for_caches(&self) -> CacheWaitDurations {
1965 debug!(target: "engine::tree::payload_validator", "Waiting for execution cache and sparse trie locks");
1966
1967 let execution_cache = self.payload_processor.execution_cache();
1968 let overlay_manager = self.overlay_manager.clone();
1969 let (execution_tx, execution_rx) = std::sync::mpsc::channel();
1970 let (sparse_trie_tx, sparse_trie_rx) = std::sync::mpsc::channel();
1971
1972 self.runtime.spawn_blocking_named("wait-exec-cache", move || {
1973 let _ = execution_tx.send(execution_cache.wait_for_availability());
1974 });
1975 self.runtime.spawn_blocking_named("wait-sparse-tri", move || {
1976 let _ = sparse_trie_tx.send(overlay_manager.wait_for_sparse_trie_availability());
1977 });
1978
1979 let execution_cache =
1980 execution_rx.recv().expect("execution cache wait task failed to send result");
1981 let sparse_trie =
1982 sparse_trie_rx.recv().expect("sparse trie wait task failed to send result");
1983 debug!(
1984 target: "engine::tree::payload_validator",
1985 ?execution_cache,
1986 ?sparse_trie,
1987 "Execution cache and sparse trie locks acquired"
1988 );
1989 CacheWaitDurations { execution_cache, sparse_trie }
1990 }
1991}
1992
1993#[derive(Debug, Clone)]
1995pub enum BlockOrPayload<T: PayloadTypes> {
1996 Payload(T::ExecutionData),
1998 Block(SealedBlockWithAccessList<BlockTy<<T::BuiltPayload as BuiltPayload>::Primitives>>),
2000}
2001
2002impl<T: PayloadTypes> BlockOrPayload<T> {
2003 pub fn hash(&self) -> B256 {
2005 match self {
2006 Self::Payload(payload) => payload.block_hash(),
2007 Self::Block(block) => block.hash(),
2008 }
2009 }
2010
2011 pub fn num_hash(&self) -> NumHash {
2013 match self {
2014 Self::Payload(payload) => payload.num_hash(),
2015 Self::Block(block) => block.num_hash(),
2016 }
2017 }
2018
2019 pub fn parent_hash(&self) -> B256 {
2021 match self {
2022 Self::Payload(payload) => payload.parent_hash(),
2023 Self::Block(block) => block.parent_hash(),
2024 }
2025 }
2026
2027 pub fn block_with_parent(&self) -> BlockWithParent {
2029 match self {
2030 Self::Payload(payload) => payload.block_with_parent(),
2031 Self::Block(block) => block.block_with_parent(),
2032 }
2033 }
2034
2035 pub const fn type_name(&self) -> &'static str {
2037 match self {
2038 Self::Payload(_) => "payload",
2039 Self::Block(_) => "block",
2040 }
2041 }
2042
2043 pub const fn is_payload(&self) -> bool {
2045 matches!(self, Self::Payload(_))
2046 }
2047
2048 pub const fn is_block(&self) -> bool {
2050 matches!(self, Self::Block(_))
2051 }
2052
2053 pub fn try_decoded_access_list(&self) -> Result<Option<DecodedBal>, alloy_rlp::Error> {
2055 match self {
2056 Self::Payload(payload) => payload
2057 .block_access_list()
2058 .map(|block_access_list| DecodedBal::from_rlp_bytes(block_access_list.clone()))
2059 .transpose(),
2060 Self::Block(block) => block.data().clone().map(DecodedBal::from_raw_bal).transpose(),
2061 }
2062 }
2063
2064 pub fn has_block_access_list(&self) -> bool {
2066 match self {
2067 Self::Payload(payload) => payload.block_access_list().is_some(),
2068 Self::Block(block) => block.block_access_list_hash().is_some(),
2069 }
2070 }
2071
2072 pub fn transaction_count(&self) -> usize
2074 where
2075 T::ExecutionData: ExecutionPayload,
2076 {
2077 match self {
2078 Self::Payload(payload) => payload.transaction_count(),
2079 Self::Block(block) => block.transaction_count(),
2080 }
2081 }
2082
2083 pub fn withdrawals(&self) -> Option<&[Withdrawal]>
2085 where
2086 T::ExecutionData: ExecutionPayload,
2087 {
2088 match self {
2089 Self::Payload(payload) => payload.withdrawals().map(|w| w.as_slice()),
2090 Self::Block(block) => block.body().withdrawals().map(|w| w.as_slice()),
2091 }
2092 }
2093
2094 pub fn gas_used(&self) -> u64
2096 where
2097 T::ExecutionData: ExecutionPayload,
2098 {
2099 match self {
2100 Self::Payload(payload) => payload.gas_used(),
2101 Self::Block(block) => block.gas_used(),
2102 }
2103 }
2104
2105 pub fn gas_limit(&self) -> u64
2107 where
2108 T::ExecutionData: ExecutionPayload,
2109 {
2110 match self {
2111 Self::Payload(payload) => payload.gas_limit(),
2112 Self::Block(block) => block.gas_limit(),
2113 }
2114 }
2115}
2116
2117struct ExecutedBal {
2119 alloy: BlockAccessList,
2121 revm: Arc<RevmBal>,
2123}