1use super::{bal_prewarm_pool::BalPrewarmPool, StateRootHintStream, StateRootUpdateStream};
15use crate::tree::{
16 precompile_cache::{CachedPrecompile, PrecompileCacheMap},
17 CachedStateCacheMetrics, CachedStateMetrics, CachedStateProvider, ExecutionEnv,
18 PayloadExecutionCache, SavedCache,
19};
20use alloy_consensus::transaction::TxHashRef;
21use alloy_eip7928::bal::DecodedBal;
22use alloy_eips::eip4895::Withdrawal;
23use alloy_primitives::keccak256;
24use metrics::{Counter, Gauge, Histogram};
25use rayon::prelude::*;
26use reth_evm::{execute::ExecutableTxFor, ConfigureEvm, Evm, EvmFor, RecoveredTx, SpecFor};
27use reth_metrics::Metrics;
28use reth_primitives_traits::{FastInstant as Instant, NodePrimitives};
29use reth_provider::{
30 BlockExecutionOutput, BlockNumReader, ChangeSetReader, DatabaseProviderFactory,
31 DatabaseProviderROFactory, EvmStateProvider, EvmStateProviderBox, HistoryReader,
32 PruneCheckpointReader, StageCheckpointReader, StateProvider, StorageChangeSetReader,
33 StorageSettingsCache,
34};
35use reth_revm::database::StateProviderDatabase;
36use reth_storage_overlay::OverlayStateProviderFactory;
37use reth_tasks::{pool::WorkerPool, Runtime};
38use reth_trie_common::MultiProofTargetsV2;
39use std::sync::{
40 atomic::{AtomicBool, AtomicUsize, Ordering},
41 mpsc::{self, channel, Receiver, Sender},
42 Arc,
43};
44use tokio::sync::oneshot;
45use tracing::{debug, debug_span, instrument, trace, trace_span, warn, Span};
46
47#[derive(Debug)]
52pub enum PrewarmMode<Tx> {
53 Transactions {
55 pending: Receiver<(usize, Tx)>,
57 hints: Option<StateRootHintStream>,
59 },
60 BlockAccessList {
62 bal: Arc<DecodedBal>,
64 updates: Option<StateRootUpdateStream>,
66 },
67 Skipped,
70}
71
72#[derive(Debug)]
77pub struct PrewarmCacheTask<N, P, Evm>
78where
79 N: NodePrimitives,
80 Evm: ConfigureEvm<Primitives = N>,
81{
82 executor: Runtime,
84 execution_cache: PayloadExecutionCache,
86 ctx: PrewarmContext<N, P, Evm>,
88 actions_rx: Receiver<PrewarmTaskEvent<N::Receipt>>,
90 parent_span: Span,
92}
93
94impl<N, P, Evm> PrewarmCacheTask<N, P, Evm>
95where
96 N: NodePrimitives,
97 P: DatabaseProviderFactory + Clone + 'static,
98 P::Provider: BlockNumReader
99 + PruneCheckpointReader
100 + StageCheckpointReader
101 + ChangeSetReader
102 + StorageChangeSetReader
103 + StorageSettingsCache
104 + HistoryReader
105 + 'static,
106 Evm: ConfigureEvm<Primitives = N> + 'static,
107{
108 pub fn new(
110 executor: Runtime,
111 execution_cache: PayloadExecutionCache,
112 ctx: PrewarmContext<N, P, Evm>,
113 ) -> (Self, Sender<PrewarmTaskEvent<N::Receipt>>) {
114 let (actions_tx, actions_rx) = channel();
115
116 trace!(
117 target: "engine::tree::payload_processor::prewarm",
118 prewarming_threads = executor.prewarming_pool().current_num_threads(),
119 transaction_count = ctx.env.transaction_count,
120 "Initialized prewarm task"
121 );
122
123 (
124 Self { executor, execution_cache, ctx, actions_rx, parent_span: Span::current() },
125 actions_tx,
126 )
127 }
128
129 fn spawn_txs_prewarm<Tx>(
136 &self,
137 pending: mpsc::Receiver<(usize, Tx)>,
138 actions_tx: Sender<PrewarmTaskEvent<N::Receipt>>,
139 state_root_hint_stream: Option<StateRootHintStream>,
140 ) where
141 Tx: ExecutableTxFor<Evm> + Send + 'static,
142 {
143 let executor = self.executor.clone();
144 let ctx = self.ctx.clone();
145 let span = Span::current();
146
147 self.executor.spawn_blocking_named("prewarm-txs", move || {
148 let _enter = debug_span!(
149 target: "engine::tree::payload_processor::prewarm",
150 parent: &span,
151 "prewarm_txs"
152 )
153 .entered();
154
155 let ctx = &ctx;
156 let pool = executor.prewarming_pool();
157
158 let mut tx_count = 0usize;
159 let state_root_hint_stream = state_root_hint_stream.as_ref();
160 pool.in_place_scope(|s| {
161 s.spawn(|_| {
162 pool.init::<PrewarmEvmState<Evm>>(|_| ctx.evm_for_ctx());
163 });
164
165 while let Ok((index, tx)) = pending.recv() {
166 if ctx.should_stop() {
167 trace!(
168 target: "engine::tree::payload_processor::prewarm",
169 "Termination requested, stopping transaction distribution"
170 );
171 break;
172 }
173
174 if index < ctx.executed_tx_index.load(Ordering::Relaxed) {
176 continue;
177 }
178
179 tx_count += 1;
180 let parent_span = Span::current();
181 s.spawn(move |_| {
182 let _enter = trace_span!(
183 target: "engine::tree::payload_processor::prewarm",
184 parent: parent_span,
185 "prewarm_tx",
186 i = index,
187 )
188 .entered();
189 Self::transact_worker(ctx, index, tx, state_root_hint_stream);
190 });
191 }
192
193 if let Some(state_root_hint_stream) = state_root_hint_stream &&
195 let Some(withdrawals) = &ctx.env.withdrawals &&
196 !withdrawals.is_empty()
197 {
198 let targets = multiproof_targets_from_withdrawals(withdrawals);
199 state_root_hint_stream.on_access_hint(targets.into());
200 }
201 });
202
203 pool.clear();
205
206 let _ = actions_tx
207 .send(PrewarmTaskEvent::FinishedTxExecution { executed_transactions: tx_count });
208 });
209 }
210
211 fn transact_worker<Tx>(
216 ctx: &PrewarmContext<N, P, Evm>,
217 index: usize,
218 tx: Tx,
219 state_root_hint_stream: Option<&StateRootHintStream>,
220 ) where
221 Tx: ExecutableTxFor<Evm>,
222 {
223 WorkerPool::with_worker_mut(|worker| {
224 let Some(evm) =
225 worker.get_or_init::<PrewarmEvmState<Evm>>(|| ctx.evm_for_ctx()).as_mut()
226 else {
227 return;
228 };
229
230 if ctx.should_stop() {
231 return;
232 }
233
234 if index < ctx.executed_tx_index.load(Ordering::Relaxed) {
236 return;
237 }
238
239 let start = Instant::now();
240
241 let (tx_env, tx) = tx.into_parts();
242 let res = match evm.transact(tx_env) {
243 Ok(res) => res,
244 Err(err) => {
245 trace!(
246 target: "engine::tree::payload_processor::prewarm",
247 %err,
248 tx_hash=%tx.tx().tx_hash(),
249 sender=%tx.signer(),
250 "Error when executing prewarm transaction",
251 );
252 ctx.metrics.transaction_errors.increment(1);
253 return;
254 }
255 };
256 ctx.metrics.execution_duration.record(start.elapsed());
257
258 if ctx.should_stop() {
259 return;
260 }
261
262 if index > 0 {
263 let (targets, storage_targets) = MultiProofTargetsV2::from_state(res.state);
264 ctx.metrics.prefetch_storage_targets.record(storage_targets as f64);
265 if let Some(state_root_hint_stream) = state_root_hint_stream {
266 state_root_hint_stream.on_access_hint(targets.into());
267 }
268 }
269
270 ctx.metrics.total_runtime.record(start.elapsed());
271 });
272 }
273
274 #[instrument(level = "debug", target = "engine::tree::payload_processor::prewarm", skip_all)]
290 fn save_cache(
291 self,
292 execution_outcome: Arc<BlockExecutionOutput<N::Receipt>>,
293 valid_block_rx: mpsc::Receiver<()>,
294 ) {
295 let start = Instant::now();
296
297 let Self {
298 execution_cache,
299 ctx: PrewarmContext { env, metrics, cache_state_metrics, saved_cache, .. },
300 ..
301 } = self;
302 let hash = env.hash;
303
304 if let Some(saved_cache) = saved_cache {
305 debug!(target: "engine::caching", parent_hash=?hash, "Updating execution cache");
306 let (previous, rejected) = execution_cache.update_with_guard(|cached| {
307 let new_cache = SavedCache::new(hash, saved_cache.into_cache());
308
309 if new_cache.cache().insert_state(&execution_outcome.state).is_err() {
311 debug!(target: "engine::caching", "cleared execution cache on update error");
312 return (cached.take(), Some(new_cache));
313 }
314
315 new_cache.update_metrics(cache_state_metrics.as_ref());
316
317 if valid_block_rx.recv().is_err() {
321 debug!(target: "engine::caching", "cleared execution cache on invalid block");
322 return (cached.take(), Some(new_cache));
323 }
324
325 let reused =
326 cached.as_ref().is_some_and(|previous| previous.shares_cache_with(&new_cache));
327 let previous = cached.replace(new_cache);
328 if reused {
329 drop(previous);
332 (None, None)
333 } else {
334 (previous, None)
336 }
337 });
338
339 drop((previous, rejected));
344
345 let elapsed = start.elapsed();
346 debug!(target: "engine::caching", parent_hash=?hash, elapsed=?elapsed, "Updated execution cache");
347
348 metrics.cache_saving_duration.set(elapsed.as_secs_f64());
349 }
350 }
351
352 #[instrument(level = "debug", target = "engine::tree::payload_processor::prewarm", skip_all)]
360 fn run_bal_prewarm(
361 &self,
362 decoded_bal: Arc<DecodedBal>,
363 actions_tx: Sender<PrewarmTaskEvent<N::Receipt>>,
364 hashed_update_stream: Option<StateRootUpdateStream>,
365 ) {
366 let bal = decoded_bal.as_bal();
367 if bal.is_empty() {
368 if let Some(hashed_update_stream) = hashed_update_stream {
369 hashed_update_stream.finish();
370 }
371 let _ =
372 actions_tx.send(PrewarmTaskEvent::FinishedTxExecution { executed_transactions: 0 });
373 return;
374 }
375
376 trace!(
377 target: "engine::tree::payload_processor::prewarm",
378 accounts = bal.len(),
379 "Starting BAL prewarm"
380 );
381
382 let ctx = self.ctx.clone();
383 let executor = self.executor.clone();
384 let parent_span = Span::current();
385 let stream_parent_span = parent_span;
386 let prefetch_bal = Arc::clone(&decoded_bal);
387 let stream_bal = Arc::clone(&decoded_bal);
388 let (stream_tx, stream_rx) = oneshot::channel();
389
390 if let Some(hashed_update_stream) = hashed_update_stream {
391 let ctx = ctx.clone();
392 executor.bal_streaming_pool().spawn(move || {
393 let branch_span = debug_span!(
394 target: "engine::tree::payload_processor::prewarm",
395 parent: &stream_parent_span,
396 "bal_hashed_state_stream",
397 bal_accounts = stream_bal.as_bal().len(),
398 );
399 let parent_span = branch_span.clone();
400 let _span = branch_span.entered();
401
402 stream_bal.as_bal().par_iter().for_each(|account_changes| {
403 WorkerPool::with_worker_mut(|worker| {
404 let provider = worker.get_or_init::<Option<EvmStateProviderBox>>(|| None);
405 ctx.send_bal_hashed_state(
406 &parent_span,
407 provider,
408 account_changes,
409 &hashed_update_stream,
410 );
411 });
412 });
413
414 hashed_update_stream.finish();
415 let _ = stream_tx.send(());
416 });
417 } else {
418 let _ = stream_tx.send(());
419 }
420
421 if let Some(saved_cache) = ctx.saved_cache &&
422 !ctx.disable_bal_batch_io &&
423 let Some(pool) = ctx.bal_prewarm_pool.as_ref()
424 {
425 let caches = saved_cache.cache().clone();
437 let state_provider_factory = ctx.provider.clone();
438 let build = Arc::new(move || {
439 state_provider_factory.database_provider_ro().map(|provider| {
440 Box::new(provider.into_evm_state_provider()) as EvmStateProviderBox
441 })
442 });
443
444 pool.begin_block(build, caches, ctx.env.txpool_snapshot.clone());
445 let dispatch_start = Instant::now();
446 for account in prefetch_bal.as_bal() {
447 pool.warm_account(account.address, account.storage_slots().map(Into::into));
448 }
449 ctx.metrics.bal_slot_iteration_duration.record(dispatch_start.elapsed());
450 pool.end_block();
451 }
452
453 stream_rx
454 .blocking_recv()
455 .expect("BAL hashed-state streaming task dropped without signaling completion");
456
457 executor.bal_streaming_pool().clear();
459
460 let _ = actions_tx.send(PrewarmTaskEvent::FinishedTxExecution { executed_transactions: 0 });
461 }
462
463 #[instrument(
468 parent = &self.parent_span,
469 level = "debug",
470 target = "engine::tree::payload_processor::prewarm",
471 name = "prewarm and caching",
472 skip_all
473 )]
474 pub fn run<Tx>(self, mode: PrewarmMode<Tx>, actions_tx: Sender<PrewarmTaskEvent<N::Receipt>>)
475 where
476 Tx: ExecutableTxFor<Evm> + Send + 'static,
477 {
478 match mode {
482 PrewarmMode::Transactions { pending, hints } => {
483 self.spawn_txs_prewarm(pending, actions_tx, hints);
484 }
485 PrewarmMode::BlockAccessList { bal, updates } => {
486 self.run_bal_prewarm(bal, actions_tx, updates);
487 }
488 PrewarmMode::Skipped => {
489 let _ = actions_tx
490 .send(PrewarmTaskEvent::FinishedTxExecution { executed_transactions: 0 });
491 }
492 }
493
494 let mut final_execution_outcome = None;
495 let mut finished_execution = false;
496 while let Ok(event) = self.actions_rx.recv() {
497 match event {
498 PrewarmTaskEvent::TerminateTransactionExecution => {
499 debug!(target: "engine::tree::prewarm", "Terminating prewarm execution");
501 self.ctx.stop();
502 }
503 PrewarmTaskEvent::Terminate { execution_outcome, valid_block_rx } => {
504 trace!(target: "engine::tree::payload_processor::prewarm", "Received termination signal");
505 self.ctx.stop();
508 final_execution_outcome =
509 Some(execution_outcome.map(|outcome| (outcome, valid_block_rx)));
510
511 if finished_execution {
512 break
514 }
515 }
516 PrewarmTaskEvent::FinishedTxExecution { executed_transactions } => {
517 trace!(target: "engine::tree::payload_processor::prewarm", "Finished prewarm execution signal");
518 self.ctx.metrics.transactions.set(executed_transactions as f64);
519 self.ctx.metrics.transactions_histogram.record(executed_transactions as f64);
520
521 finished_execution = true;
522
523 if final_execution_outcome.is_some() {
524 break
526 }
527 }
528 }
529 }
530
531 debug!(target: "engine::tree::payload_processor::prewarm", "Completed prewarm execution");
532
533 if let Some(Some((execution_outcome, valid_block_rx))) = final_execution_outcome {
535 self.save_cache(execution_outcome, valid_block_rx);
536 }
537 }
538}
539
540#[derive(Debug, Clone)]
542pub struct PrewarmContext<N, P, Evm>
543where
544 N: NodePrimitives,
545 Evm: ConfigureEvm<Primitives = N>,
546{
547 pub env: ExecutionEnv<Evm>,
549 pub evm_config: Evm,
551 pub saved_cache: Option<SavedCache>,
553 pub provider: OverlayStateProviderFactory<P, N>,
555 pub(crate) bal_prewarm_pool: Option<Arc<BalPrewarmPool>>,
558 pub metrics: PrewarmMetrics,
560 pub cache_metrics: Option<CachedStateMetrics>,
563 pub cache_state_metrics: Option<CachedStateCacheMetrics>,
565 pub terminate_execution: Arc<AtomicBool>,
567 pub executed_tx_index: Arc<AtomicUsize>,
571 pub precompile_cache_disabled: bool,
573 pub precompile_cache_map: PrecompileCacheMap<SpecFor<Evm>>,
575 pub disable_bal_parallel_state_root: bool,
578 pub disable_bal_batch_io: bool,
580}
581
582type PrewarmEvmState<Evm> = Option<EvmFor<Evm, StateProviderDatabase<EvmStateProviderBox>>>;
585
586impl<N, P, Evm> PrewarmContext<N, P, Evm>
587where
588 N: NodePrimitives,
589 P: DatabaseProviderFactory,
590 P::Provider: BlockNumReader
591 + PruneCheckpointReader
592 + StageCheckpointReader
593 + ChangeSetReader
594 + StorageChangeSetReader
595 + StorageSettingsCache
596 + HistoryReader
597 + 'static,
598 Evm: ConfigureEvm<Primitives = N> + 'static,
599{
600 #[instrument(level = "debug", target = "engine::tree::payload_processor::prewarm", skip_all)]
602 fn evm_for_ctx(&self) -> PrewarmEvmState<Evm> {
603 let mut state_provider = match self.provider.database_provider_ro() {
604 Ok(provider) => Box::new(provider.into_evm_state_provider()) as EvmStateProviderBox,
605 Err(err) => {
606 trace!(
607 target: "engine::tree::payload_processor::prewarm",
608 %err,
609 "Failed to build state provider in prewarm thread"
610 );
611 return None
612 }
613 };
614
615 if let Some(saved_cache) = &self.saved_cache {
617 let caches = saved_cache.cache().clone();
618 state_provider = Box::new(
619 CachedStateProvider::new_prewarm(state_provider, caches)
620 .with_txpool_snapshot(self.env.txpool_snapshot.clone()),
621 );
622 }
623
624 let state_provider = StateProviderDatabase::new(state_provider);
625
626 let mut evm_env = self.env.evm_env.clone();
627
628 evm_env.cfg_env.disable_nonce_check = true;
631
632 evm_env.cfg_env.disable_balance_check = true;
635
636 let spec_id = *evm_env.spec_id();
638 let mut evm = self.evm_config.evm_with_env(state_provider, evm_env);
639
640 if !self.precompile_cache_disabled {
641 evm.precompiles_mut().map_cacheable_precompiles(|address, precompile| {
643 CachedPrecompile::wrap(
644 precompile,
645 self.precompile_cache_map.cache_for_address(*address),
646 spec_id,
647 None, )
649 });
650 }
651
652 Some(evm)
653 }
654
655 #[inline]
657 pub fn should_stop(&self) -> bool {
658 self.terminate_execution.load(Ordering::Relaxed)
659 }
660
661 #[inline]
663 pub fn stop(&self) {
664 self.terminate_execution.store(true, Ordering::Relaxed);
665 }
666
667 fn send_bal_hashed_state(
677 &self,
678 parent_span: &Span,
679 provider: &mut Option<EvmStateProviderBox>,
680 account_changes: &alloy_eip7928::AccountChanges,
681 hashed_update_stream: &StateRootUpdateStream,
682 ) {
683 if self.disable_bal_parallel_state_root {
684 return;
685 }
686 let address = account_changes.address;
687 let mut hashed_address = None;
688 let account_info = account_changes.account_info();
689
690 if !account_info.changes_state_root(account_changes) {
691 return;
692 }
693
694 if account_changes.has_storage_changes() {
698 let hashed_address = *hashed_address.get_or_insert_with(|| keccak256(address));
699 let storage_map = reth_trie::HashedStorage::from_iter(
700 account_changes
701 .storage_post_states()
702 .map(|(slot, value)| (keccak256(slot.to_be_bytes::<32>()), value)),
703 );
704
705 let mut hashed_state = reth_trie::HashedPostState::default();
706 hashed_state.storages.insert(hashed_address, storage_map);
707 hashed_update_stream.on_hashed_state_update(hashed_state);
708 }
709
710 let existing_account = if account_info.is_complete() {
711 None
712 } else {
713 if provider.is_none() {
714 let _span = debug_span!(
715 target: "engine::tree::payload_processor::prewarm",
716 parent: parent_span,
717 "bal_hashed_state_provider_init",
718 has_saved_cache = !self.disable_bal_batch_io && self.saved_cache.is_some(),
719 )
720 .entered();
721
722 let inner = match self.provider.database_provider_ro() {
723 Ok(p) => p.into_evm_state_provider(),
724 Err(err) => {
725 warn!(
726 target: "engine::tree::payload_processor::prewarm",
727 ?err,
728 "Failed to build provider for BAL account reads"
729 );
730 return;
731 }
732 };
733 let boxed: EvmStateProviderBox =
734 match (self.disable_bal_batch_io, &self.saved_cache) {
735 (false, Some(saved)) => {
736 let caches = saved.cache().clone();
737 Box::new(
738 CachedStateProvider::new_prewarm(inner, caches)
739 .with_txpool_snapshot(self.env.txpool_snapshot.clone()),
740 )
741 }
742 _ => Box::new(inner),
743 };
744 *provider = Some(boxed);
745 }
746 let account_reader = provider.as_ref().expect("provider just initialized");
747 account_reader.basic_account(&address).ok().flatten()
748 };
749
750 let mut account = existing_account.unwrap_or_default();
751 account.apply_bal_info(account_info);
752 let hashed_address = hashed_address.unwrap_or_else(|| keccak256(address));
753
754 let account = (!account.is_empty()).then_some(account);
766
767 let mut hashed_state = reth_trie::HashedPostState::default();
768 hashed_state.accounts.insert(hashed_address, account);
769 hashed_update_stream.on_hashed_state_update(hashed_state);
770 }
771}
772
773fn multiproof_targets_from_withdrawals(withdrawals: &[Withdrawal]) -> MultiProofTargetsV2 {
778 MultiProofTargetsV2 {
779 account_targets: withdrawals.iter().map(|w| keccak256(w.address).into()).collect(),
780 ..Default::default()
781 }
782}
783
784#[cfg(test)]
785mod tests {
786 use super::*;
787 use alloy_consensus::transaction::Recovered;
788 use alloy_eip7928::{AccountChanges, BalanceChange, BlockAccessIndex};
789 use alloy_eips::eip7702::constants::EIP7702_CLEARED_DELEGATION;
790 use alloy_primitives::{address, B256, U256};
791 use reth_chainspec::ChainSpec;
792 use reth_ethereum_primitives::{EthPrimitives, TransactionSigned};
793 use reth_evm::{execute::WithTxEnv, TxEnvFor};
794 use reth_evm_ethereum::EthEvmConfig;
795 use reth_primitives_traits::Account;
796 use reth_provider::test_utils::MockEthProvider;
797 use reth_storage_overlay::OverlayManager;
798
799 #[test]
800 fn terminate_event_stops_transaction_execution() {
801 let terminate_execution = Arc::new(AtomicBool::new(false));
802 let ctx = PrewarmContext {
803 env: ExecutionEnv::test_default(),
804 evm_config: EthEvmConfig::new(Arc::new(ChainSpec::default())),
805 saved_cache: None,
806 provider: OverlayStateProviderFactory::new(
807 MockEthProvider::default(),
808 OverlayManager::default().overlay_builder(B256::ZERO),
809 ),
810 bal_prewarm_pool: None,
811 metrics: PrewarmMetrics::default(),
812 cache_metrics: None,
813 cache_state_metrics: None,
814 terminate_execution: Arc::clone(&terminate_execution),
815 executed_tx_index: Arc::new(AtomicUsize::new(0)),
816 precompile_cache_disabled: false,
817 precompile_cache_map: PrecompileCacheMap::default(),
818 disable_bal_parallel_state_root: false,
819 disable_bal_batch_io: false,
820 };
821 let (task, actions_tx) =
822 PrewarmCacheTask::new(Runtime::test(), PayloadExecutionCache::default(), ctx);
823 actions_tx
824 .send(PrewarmTaskEvent::Terminate {
825 execution_outcome: None,
826 valid_block_rx: mpsc::channel().1,
827 })
828 .unwrap();
829
830 task.run::<WithTxEnv<TxEnvFor<EthEvmConfig>, Recovered<TransactionSigned>>>(
831 PrewarmMode::Skipped,
832 actions_tx,
833 );
834
835 assert!(terminate_execution.load(Ordering::Relaxed));
836 }
837
838 fn test_prewarm_context(
839 saved_cache: SavedCache,
840 saving_duration: Gauge,
841 ) -> PrewarmContext<EthPrimitives, MockEthProvider, EthEvmConfig> {
842 PrewarmContext {
843 env: ExecutionEnv { hash: B256::repeat_byte(2), ..ExecutionEnv::test_default() },
844 evm_config: EthEvmConfig::new(Arc::new(ChainSpec::default())),
845 saved_cache: Some(saved_cache),
846 provider: OverlayStateProviderFactory::new(
847 MockEthProvider::default(),
848 OverlayManager::default().overlay_builder(B256::ZERO),
849 ),
850 bal_prewarm_pool: None,
851 metrics: PrewarmMetrics {
852 cache_saving_duration: saving_duration,
853 ..Default::default()
854 },
855 cache_metrics: None,
856 cache_state_metrics: None,
857 terminate_execution: Arc::new(AtomicBool::new(false)),
858 executed_tx_index: Arc::new(AtomicUsize::new(0)),
859 precompile_cache_disabled: false,
860 precompile_cache_map: PrecompileCacheMap::default(),
861 disable_bal_parallel_state_root: false,
862 disable_bal_batch_io: false,
863 }
864 }
865
866 fn save_test_cache(
867 runtime: &Runtime,
868 execution_cache: &PayloadExecutionCache,
869 saved_cache: SavedCache,
870 state: reth_revm::db::BundleState,
871 valid: bool,
872 saving_duration: Gauge,
873 ) {
874 let ctx = test_prewarm_context(saved_cache, saving_duration);
875 let (task, _) = PrewarmCacheTask::new(runtime.clone(), execution_cache.clone(), ctx);
876 let (valid_tx, valid_rx) = mpsc::channel();
877 if valid {
878 valid_tx.send(()).unwrap();
879 }
880 drop(valid_tx);
881 task.save_cache(
882 Arc::new(BlockExecutionOutput { state, result: Default::default() }),
883 valid_rx,
884 );
885 }
886
887 struct CacheSaveObserver {
890 cache: PayloadExecutionCache,
891 observed: Arc<AtomicBool>,
892 }
893
894 impl metrics::GaugeFn for CacheSaveObserver {
895 fn increment(&self, _: f64) {}
896 fn decrement(&self, _: f64) {}
897 fn set(&self, _: f64) {
898 assert!(
899 self.cache.get_cache_for(B256::repeat_byte(2)).is_some(),
900 "saved cache must be available before the saving task returns"
901 );
902 self.observed.store(true, Ordering::Relaxed);
903 }
904 }
905
906 #[test]
907 fn save_cache_releases_warm_cache_before_duration_metric() {
908 let runtime = Runtime::test();
909 let execution_cache = PayloadExecutionCache::default();
910 let saved = SavedCache::new(B256::repeat_byte(1), crate::tree::ExecutionCache::new(1_000));
911 let address = address!("0000000000000000000000000000000000000001");
912 saved.cache().insert_storage(address, B256::ZERO, Some(U256::from(7)));
913 execution_cache.update_with_guard(|slot| *slot = Some(saved.clone()));
914 let (release_tx, release_rx) = mpsc::channel::<()>();
916 runtime.spawn_blocking_named("drop", move || {
917 let _ = release_rx.recv();
918 });
919 let observed = Arc::new(AtomicBool::new(false));
920 let observer = Arc::new(CacheSaveObserver {
921 cache: execution_cache.clone(),
922 observed: observed.clone(),
923 });
924 save_test_cache(
925 &runtime,
926 &execution_cache,
927 saved,
928 Default::default(),
929 true,
930 Gauge::from_arc(observer),
931 );
932 assert!(observed.load(Ordering::Relaxed), "save duration was not recorded");
933 let saved = execution_cache.get_cache_for(B256::repeat_byte(2)).unwrap();
934 assert_eq!(
935 saved.cache().get_or_try_insert_storage_with(address, B256::ZERO, || Err(())),
936 Ok(reth_execution_cache::CachedStatus::Cached(U256::from(7))),
937 "handoff must preserve the warmed contents",
938 );
939 drop(release_tx);
940 }
941
942 #[test]
943 fn save_cache_blocks_reuse_while_execution_cache_is_cloned() {
944 let runtime = Runtime::test();
945 let execution_cache = PayloadExecutionCache::default();
946 let saved = SavedCache::new(B256::repeat_byte(1), crate::tree::ExecutionCache::new(1_000));
947 execution_cache.update_with_guard(|slot| *slot = Some(saved.clone()));
948 let other_cache_clone = saved.cache().clone();
949
950 save_test_cache(&runtime, &execution_cache, saved, Default::default(), true, Gauge::noop());
951
952 assert!(execution_cache.get_cache_for(B256::repeat_byte(2)).is_none());
953 drop(other_cache_clone);
954 assert!(execution_cache.get_cache_for(B256::repeat_byte(2)).is_some());
955 }
956
957 #[derive(Clone, Copy, Debug)]
958 enum CacheSlot {
959 Empty,
960 Shared,
961 Distinct,
962 }
963
964 struct CacheDropProbe {
966 started: Sender<std::thread::ThreadId>,
967 inspected: Receiver<()>,
968 result: Sender<bool>,
969 }
970
971 impl AsRef<[u8]> for CacheDropProbe {
972 fn as_ref(&self) -> &[u8] {
973 &EIP7702_CLEARED_DELEGATION
974 }
975 }
976
977 impl Drop for CacheDropProbe {
978 fn drop(&mut self) {
979 let _ = self.started.send(std::thread::current().id());
980 let unlocked = self.inspected.recv_timeout(std::time::Duration::from_secs(5)).is_ok();
982 let _ = self.result.send(unlocked);
983 }
984 }
985
986 fn observe_cache_drop(
987 saved: &SavedCache,
988 execution_cache: &PayloadExecutionCache,
989 expect_saved_cache: bool,
990 ) -> (Receiver<bool>, std::thread::JoinHandle<std::thread::ThreadId>) {
991 let (started_tx, started_rx) = mpsc::channel();
992 let (inspected_tx, inspected_rx) = mpsc::channel();
993 let (result_tx, result_rx) = mpsc::channel();
994 let cache = execution_cache.clone();
995 let reader = std::thread::spawn(move || {
996 let drop_thread = started_rx.recv().unwrap();
997 cache.wait_for_availability();
999 assert_eq!(cache.get_cache_for(B256::repeat_byte(2)).is_some(), expect_saved_cache);
1000 cache.update_with_guard(|slot| assert_eq!(slot.is_some(), expect_saved_cache));
1001 let _ = inspected_tx.send(());
1002 drop_thread
1003 });
1004 let bytes = alloy_primitives::bytes::Bytes::from_owner(CacheDropProbe {
1005 started: started_tx,
1006 inspected: inspected_rx,
1007 result: result_tx,
1008 });
1009 let code = reth_revm::bytecode::Bytecode::new_eip7702_raw(bytes.into()).unwrap();
1010 saved
1011 .cache()
1012 .insert_code(B256::repeat_byte(3), Some(reth_primitives_traits::Bytecode(code)));
1013 (result_rx, reader)
1014 }
1015
1016 #[test]
1017 fn save_cache_freeing_waits_for_validator_cache_references() {
1018 use crate::tree::payload_processor::{CacheTaskHandle, PayloadHandle};
1019
1020 let runtime = Runtime::test();
1021 let execution_cache = PayloadExecutionCache::default();
1022 let saved = SavedCache::new(B256::repeat_byte(1), crate::tree::ExecutionCache::new(1_000));
1023 execution_cache.update_with_guard(|cached| *cached = Some(saved.clone()));
1024 let (dropped, reader) = observe_cache_drop(&saved, &execution_cache, false);
1025 let ctx = test_prewarm_context(saved, Gauge::noop());
1026 let (task, actions_tx) =
1027 PrewarmCacheTask::new(runtime.clone(), execution_cache.clone(), ctx);
1028 let mut payload = PayloadHandle {
1029 prewarm_handle: CacheTaskHandle {
1030 saved_cache: task.ctx.saved_cache.clone(),
1031 to_prewarm_task: Some(actions_tx.clone()),
1032 executed_tx_index: task.ctx.executed_tx_index.clone(),
1033 cache_metrics: None,
1034 },
1035 transactions: crossbeam_channel::never::<(usize, Result<(), ()>)>(),
1036 _span: Span::none(),
1037 };
1038 let validator_cache = payload.caches().unwrap();
1040 let valid_block_tx = payload.terminate_caching(Some(Arc::new(BlockExecutionOutput {
1041 state: Default::default(),
1042 result: Default::default(),
1043 })));
1044 drop(valid_block_tx);
1045 let prewarm = runtime.spawn_blocking_named("prewarm", move || {
1046 task.run::<WithTxEnv<TxEnvFor<EthEvmConfig>, Recovered<TransactionSigned>>>(
1047 PrewarmMode::Skipped,
1048 actions_tx,
1049 );
1050 std::thread::current().id()
1051 });
1052 let prewarm_thread = *prewarm.get();
1053
1054 execution_cache.update_with_guard(|cached| assert!(cached.is_none()));
1055 assert_eq!(payload.prewarm_handle.saved_cache.as_ref().unwrap().usage_count(), 2);
1056 assert!(matches!(dropped.try_recv(), Err(mpsc::TryRecvError::Empty)));
1057 drop(payload);
1058 assert!(matches!(dropped.try_recv(), Err(mpsc::TryRecvError::Empty)));
1059
1060 drop(validator_cache);
1061 assert!(dropped.try_recv().expect("last ExecutionCache drop must free the cache"));
1062 let drop_thread = reader.join().unwrap();
1063 assert_eq!(drop_thread, std::thread::current().id());
1064 assert_ne!(drop_thread, prewarm_thread);
1065 }
1066
1067 fn assert_save_cache_drops_removed_caches(slot: CacheSlot, valid: bool, insert_error: bool) {
1068 use reth_revm::db::{AccountStatus, BundleAccount, BundleState};
1069
1070 let runtime = Runtime::test();
1071 let execution_cache = PayloadExecutionCache::default();
1072 let cache_to_save =
1073 SavedCache::new(B256::repeat_byte(1), crate::tree::ExecutionCache::new(1_000));
1074 let distinct_previous = matches!(slot, CacheSlot::Distinct).then(|| {
1075 SavedCache::new(B256::repeat_byte(2), crate::tree::ExecutionCache::new(1_000))
1077 });
1078 execution_cache.update_with_guard(|cached| {
1079 *cached = match slot {
1080 CacheSlot::Empty => None,
1081 CacheSlot::Shared => Some(cache_to_save.clone()),
1082 CacheSlot::Distinct => distinct_previous.clone(),
1083 };
1084 });
1085 let expect_saved_cache = valid && !insert_error;
1086 let mut drops = Vec::new();
1087 if let Some(previous) = &distinct_previous {
1088 drops.push(observe_cache_drop(previous, &execution_cache, expect_saved_cache));
1089 }
1090 if !expect_saved_cache {
1091 drops.push(observe_cache_drop(&cache_to_save, &execution_cache, expect_saved_cache));
1092 }
1093 drop(distinct_previous);
1095
1096 let (release_tx, release_rx) = mpsc::channel::<()>();
1098 runtime.spawn_blocking_named("drop", move || {
1099 let _ = release_rx.recv();
1100 });
1101
1102 let mut state = BundleState::default();
1103 if insert_error {
1104 state.state.insert(
1106 address!("0000000000000000000000000000000000000001"),
1107 BundleAccount::new(None, None, Default::default(), AccountStatus::Changed),
1108 );
1109 }
1110 save_test_cache(&runtime, &execution_cache, cache_to_save, state, valid, Gauge::noop());
1111
1112 for (result, reader) in drops {
1113 assert!(
1114 result.try_recv().expect("removed cache must be destroyed before save returns"),
1115 "cache mutex must be unlocked during destruction"
1116 );
1117 assert_eq!(reader.join().unwrap(), std::thread::current().id());
1118 }
1119 execution_cache.update_with_guard(|slot| {
1120 if expect_saved_cache {
1121 let saved = slot.as_ref().expect("valid cache saved");
1122 assert_eq!(saved.executed_block_hash(), B256::repeat_byte(2));
1123 } else {
1124 assert!(slot.is_none(), "polluted cache must be removed");
1125 }
1126 });
1127 assert_eq!(
1128 execution_cache.get_cache_for(B256::repeat_byte(2)).is_some(),
1129 expect_saved_cache
1130 );
1131 drop(release_tx);
1132 }
1133
1134 #[test]
1135 fn save_cache_drops_replaced_allocation_after_unlock() {
1136 assert_save_cache_drops_removed_caches(CacheSlot::Distinct, true, false);
1137 }
1138
1139 #[test]
1140 fn save_cache_drops_invalid_allocations_after_unlock() {
1141 assert_save_cache_drops_removed_caches(CacheSlot::Distinct, false, false);
1142 }
1143
1144 #[test]
1145 fn save_cache_drops_allocations_after_unlock_on_insert_error() {
1146 assert_save_cache_drops_removed_caches(CacheSlot::Distinct, true, true);
1147 }
1148
1149 #[test]
1150 fn save_cache_drops_shared_allocation_after_unlock_on_invalid_block() {
1151 assert_save_cache_drops_removed_caches(CacheSlot::Shared, false, false);
1152 }
1153
1154 #[test]
1155 fn save_cache_drops_shared_allocation_after_unlock_on_insert_error() {
1156 assert_save_cache_drops_removed_caches(CacheSlot::Shared, true, true);
1157 }
1158
1159 #[test]
1160 fn save_cache_handles_empty_slot() {
1161 for (valid, insert_error) in [(true, false), (false, false), (true, true)] {
1162 assert_save_cache_drops_removed_caches(CacheSlot::Empty, valid, insert_error);
1163 }
1164 }
1165
1166 #[test]
1167 fn bal_read_only_account_does_not_change_state_root() {
1168 let changes = AccountChanges::new(address!("0000000000000000000000000000000000000001"))
1169 .with_storage_read(U256::from(1));
1170
1171 assert!(!changes.account_info().changes_state_root(&changes));
1172 }
1173
1174 #[test]
1175 #[allow(clippy::needless_update)]
1176 fn bal_account_uses_existing_fields_only_when_missing() {
1177 let changes = AccountChanges::new(address!("0000000000000000000000000000000000000001"))
1178 .with_balance_change(BalanceChange::new(BlockAccessIndex::new(1), U256::from(10)));
1179 let info = changes.account_info();
1180
1181 assert!(!info.is_complete());
1182 let mut account = Account {
1183 balance: U256::from(1),
1184 nonce: 3,
1185 bytecode_hash: Some(B256::repeat_byte(0xaa)),
1186 ..Default::default()
1187 };
1188 account.apply_bal_info(info);
1189
1190 assert_eq!(account.balance, U256::from(10));
1191 assert_eq!(account.nonce, 3);
1192 assert_eq!(account.bytecode_hash, Some(B256::repeat_byte(0xaa)));
1193 }
1194}
1195
1196#[derive(Debug)]
1201pub enum PrewarmTaskEvent<R> {
1202 TerminateTransactionExecution,
1208 Terminate {
1216 execution_outcome: Option<Arc<BlockExecutionOutput<R>>>,
1220 valid_block_rx: mpsc::Receiver<()>,
1225 },
1226 FinishedTxExecution {
1229 executed_transactions: usize,
1231 },
1232}
1233
1234#[derive(Metrics, Clone)]
1236#[metrics(scope = "sync.prewarm")]
1237pub struct PrewarmMetrics {
1238 pub(crate) transactions: Gauge,
1240 pub(crate) transactions_histogram: Histogram,
1242 pub(crate) total_runtime: Histogram,
1244 pub(crate) execution_duration: Histogram,
1246 pub(crate) prefetch_storage_targets: Histogram,
1248 pub(crate) cache_saving_duration: Gauge,
1251 pub(crate) transaction_errors: Counter,
1253 pub(crate) bal_slot_iteration_duration: Histogram,
1255}