1use super::precompile_cache::PrecompileCacheMap;
4use crate::tree::{
5 payload_processor::prewarm::{PrewarmCacheTask, PrewarmContext, PrewarmMode, PrewarmTaskEvent},
6 CachedStateCacheMetrics, CachedStateMetrics, CachedStateMetricsSource, ExecutionCache,
7 ExecutionEnv, PayloadExecutionCache, SavedCache, TreeConfig,
8};
9use alloy_eips::eip1898::BlockWithParent;
10use alloy_primitives::B256;
11use crossbeam_channel::{Receiver as CrossbeamReceiver, Sender as CrossbeamSender};
12use prewarm::PrewarmMetrics;
13use rayon::prelude::*;
14use reth_evm::{
15 block::ExecutableTxParts,
16 execute::{ExecutableTxFor, WithTxEnv},
17 ConfigureEvm, ConvertTx, ExecutableTxIterator, ExecutableTxTuple, SpecFor, TxEnvFor,
18};
19use reth_primitives_traits::{FastInstant as Instant, NodePrimitives};
20use reth_provider::{
21 BlockExecutionOutput, BlockNumReader, ChangeSetReader, DatabaseProviderFactory, HistoryReader,
22 PruneCheckpointReader, StageCheckpointReader, StorageChangeSetReader, StorageSettingsCache,
23};
24use reth_revm::db::BundleState;
25use reth_storage_overlay::OverlayStateProviderFactory;
26use reth_tasks::Runtime;
27pub use reth_trie_parallel::{
28 error::StateRootTaskError,
29 state_root_task::{
30 evm_state_to_hashed_post_state, PayloadStateRootHandle, StateAccessHint,
31 StateRootComputeOutcome, StateRootHandle, StateRootHintStream, StateRootMessage,
32 StateRootSink, StateRootTaskCancelGuard, StateRootUpdateHook, StateRootUpdateStream,
33 },
34};
35use std::{
36 ops::Not,
37 sync::{
38 atomic::{AtomicBool, AtomicUsize},
39 mpsc, Arc, OnceLock,
40 },
41};
42use tracing::{debug, debug_span, instrument, trace, trace_span, warn, Span};
43
44pub mod bal;
45pub mod bal_prewarm_pool;
46pub mod prewarm;
47pub mod receipt_root_task;
48
49pub const SMALL_BLOCK_TX_THRESHOLD: usize = 5;
52
53type IteratorTx<Evm, I> = RecoveredTx<TxEnvFor<Evm>, <I as ExecutableTxIterator<Evm>>::Recovered>;
55
56type IteratorPayloadHandle<Evm, I> = PayloadHandle<
57 IteratorTx<Evm, I>,
58 <I as ExecutableTxTuple>::Error,
59 <<Evm as ConfigureEvm>::Primitives as NodePrimitives>::Receipt,
60>;
61
62type IteratorPrewarmTxReceiver<Evm, I> =
63 PrewarmTxReceiver<TxEnvFor<Evm>, <I as ExecutableTxIterator<Evm>>::Recovered>;
64
65type IteratorExecuteTxReceiver<Evm, I> = ExecuteTxReceiver<
66 TxEnvFor<Evm>,
67 <I as ExecutableTxIterator<Evm>>::Recovered,
68 <I as ExecutableTxTuple>::Error,
69>;
70
71type RecoveredTx<TxEnv, Recovered> = WithTxEnv<TxEnv, Recovered>;
72type IndexedTxResult<Tx, Err> = (usize, Result<Tx, Err>);
73type IndexedTxReceiver<Tx, Err> = CrossbeamReceiver<IndexedTxResult<Tx, Err>>;
74type IndexedTxSender<Tx, Err> = CrossbeamSender<IndexedTxResult<Tx, Err>>;
75type PrewarmTxReceiver<TxEnv, Recovered> = mpsc::Receiver<(usize, RecoveredTx<TxEnv, Recovered>)>;
76type ExecuteTxReceiver<TxEnv, Recovered, Err> =
77 IndexedTxReceiver<RecoveredTx<TxEnv, Recovered>, Err>;
78type ExecuteTxSender<TxEnv, Recovered, Err> = IndexedTxSender<RecoveredTx<TxEnv, Recovered>, Err>;
79
80#[derive(Debug)]
82pub struct PayloadProcessor<Evm>
83where
84 Evm: ConfigureEvm,
85{
86 executor: Runtime,
88 execution_cache: PayloadExecutionCache,
90 cache_metrics: Option<CachedStateMetrics>,
92 cache_state_metrics: Option<CachedStateCacheMetrics>,
94 cross_block_cache_size: usize,
96 disable_transaction_prewarming: bool,
98 disable_state_cache: bool,
100 evm_config: Evm,
102 precompile_cache_disabled: bool,
104 precompile_cache_map: PrecompileCacheMap<SpecFor<Evm>>,
106 disable_bal_parallel_state_root: bool,
109 disable_bal_batch_io: bool,
111 bal_prewarm_pool: OnceLock<Arc<bal_prewarm_pool::BalPrewarmPool>>,
114}
115
116impl<Evm> PayloadProcessor<Evm>
117where
118 Evm: ConfigureEvm,
119{
120 pub fn new(
122 executor: Runtime,
123 evm_config: Evm,
124 config: &TreeConfig,
125 precompile_cache_map: PrecompileCacheMap<SpecFor<Evm>>,
126 ) -> Self {
127 Self {
128 executor,
129 execution_cache: Default::default(),
130 cross_block_cache_size: config.cross_block_cache_size(),
131 disable_transaction_prewarming: config.disable_prewarming(),
132 evm_config,
133 disable_state_cache: config.disable_state_cache(),
134 precompile_cache_disabled: config.precompile_cache_disabled(),
135 precompile_cache_map,
136 cache_metrics: (!config.disable_cache_metrics())
137 .then(|| CachedStateMetrics::zeroed(CachedStateMetricsSource::Engine)),
138 cache_state_metrics: (!config.disable_cache_metrics())
139 .then(CachedStateCacheMetrics::default),
140 disable_bal_parallel_state_root: config.disable_bal_parallel_state_root(),
141 disable_bal_batch_io: config.disable_bal_batch_io(),
142 bal_prewarm_pool: OnceLock::new(),
143 }
144 }
145
146 fn bal_prewarm_pool(&self) -> Arc<bal_prewarm_pool::BalPrewarmPool> {
149 self.bal_prewarm_pool
150 .get_or_init(|| {
151 bal_prewarm_pool::BalPrewarmPool::new(bal_prewarm_pool::DEFAULT_BAL_PREWARM_THREADS)
152 })
153 .clone()
154 }
155
156 pub(crate) fn execution_cache(&self) -> PayloadExecutionCache {
158 self.execution_cache.clone()
159 }
160}
161
162impl<Evm> PayloadProcessor<Evm>
163where
164 Evm: ConfigureEvm + 'static,
165{
166 #[instrument(level = "debug", target = "engine::tree::payload_processor", skip_all)]
169 pub fn spawn_with_state_root_streams<P, I: ExecutableTxIterator<Evm>>(
170 &self,
171 env: ExecutionEnv<Evm>,
172 transactions: I,
173 state_provider_factory: OverlayStateProviderFactory<P, Evm::Primitives>,
174 hint_stream: Option<StateRootHintStream>,
175 hashed_update_stream: Option<StateRootUpdateStream>,
176 parallel_bal_execution: bool,
177 ) -> IteratorPayloadHandle<Evm, I>
178 where
179 P: DatabaseProviderFactory + Clone + 'static,
180 P::Provider: BlockNumReader
181 + PruneCheckpointReader
182 + StageCheckpointReader
183 + ChangeSetReader
184 + StorageChangeSetReader
185 + StorageSettingsCache
186 + HistoryReader
187 + 'static,
188 {
189 let prewarm_transactions =
190 self.prewarms_transactions(env.transaction_count, parallel_bal_execution);
191 let (prewarm_rx, execution_rx) = self.spawn_tx_iterator(
192 transactions,
193 env.transaction_count,
194 parallel_bal_execution,
195 prewarm_transactions,
196 );
197 let prewarm_handle = self.spawn_caching_with(
198 env,
199 prewarm_rx,
200 state_provider_factory,
201 hint_stream,
202 hashed_update_stream,
203 parallel_bal_execution,
204 );
205 PayloadHandle { prewarm_handle, transactions: execution_rx, _span: Span::current() }
206 }
207
208 const fn prewarms_transactions(
215 &self,
216 transaction_count: usize,
217 parallel_bal_execution: bool,
218 ) -> bool {
219 !parallel_bal_execution &&
220 !self.disable_transaction_prewarming &&
221 transaction_count >= SMALL_BLOCK_TX_THRESHOLD
222 }
223
224 const SMALL_BLOCK_TX_THRESHOLD: usize = 30;
231
232 const PARALLEL_PREFETCH_COUNT: usize = 4;
240
241 const FIRST_PARALLEL_TX_WINDOW_SIZE: usize = 64;
243
244 #[instrument(level = "debug", target = "engine::tree::payload_processor", skip_all)]
256 fn spawn_tx_iterator<I: ExecutableTxIterator<Evm>>(
257 &self,
258 transactions: I,
259 transaction_count: usize,
260 parallel_bal_execution: bool,
261 prewarm_transactions: bool,
262 ) -> (Option<IteratorPrewarmTxReceiver<Evm, I>>, IteratorExecuteTxReceiver<Evm, I>) {
263 let (prewarm_tx, prewarm_rx) =
264 prewarm_transactions.then(|| mpsc::sync_channel(transaction_count)).unzip();
265 let (execute_tx, execute_rx) = crossbeam_channel::bounded(transaction_count);
266 let parent_span = Span::current();
267
268 if transaction_count == 0 {
269 } else if transaction_count < Self::SMALL_BLOCK_TX_THRESHOLD {
271 debug!(
274 target: "engine::tree::payload_processor",
275 transaction_count,
276 "using sequential sig recovery for small block"
277 );
278 self.executor.spawn_blocking_named("tx-iterator", move || {
279 let _parent = parent_span.enter();
281 let _span = debug_span!(target: "engine::tree::payload_processor", "tx_iterator", transaction_count).entered();
282 let (transactions, convert) = transactions.into_parts();
283 convert_serial(
284 transactions.into_iter(),
285 &convert,
286 prewarm_tx.as_ref(),
287 &execute_tx,
288 );
289 });
290 } else {
291 let executor = self.executor.clone();
294 self.executor.spawn_blocking_named("tx-iterator", move || {
295 let _parent = parent_span.enter();
296 let _span = debug_span!(target: "engine::tree::payload_processor", "tx_iterator", transaction_count).entered();
297 let (transactions, convert) = transactions.into_parts();
298 if parallel_bal_execution {
299 let parent_span = Span::current();
302 executor.cpu_pool().install(|| {
303 let _parent = parent_span.enter();
304 let span = debug_span!(target: "engine::tree::payload_processor", "convert_transactions").or_current();
305 let _entered = span.enter();
306 let _ = transactions
307 .into_par_iter()
308 .enumerate()
309 .by_exponential_blocks()
311 .try_for_each(|(idx, tx)| {
312 let _parent = span.enter();
313 let _span = trace_span!(target: "engine::tree::payload_processor", "convert_transaction", idx).entered();
314 let tx = convert.convert(tx).map(WithTxEnv::new);
315 if let (Some(prewarm_tx), Ok(tx)) = (&prewarm_tx, &tx) {
316 let _ = prewarm_tx.send((idx, tx.clone()));
317 }
318 let disconnected = execute_tx.send((idx, tx)).is_err();
319 trace!(target: "engine::tree::payload_processor", idx, "yielded transaction");
320 if disconnected {
323 Err(())
324 } else {
325 Ok(())
326 }
327 });
328 });
329 } else {
330 let prefetch = Self::PARALLEL_PREFETCH_COUNT.min(transaction_count);
334 let mut iter = transactions.into_iter();
335
336 if !convert_serial(
339 iter.by_ref().take(prefetch),
340 &convert,
341 prewarm_tx.as_ref(),
342 &execute_tx,
343 ) {
344 return
345 }
346
347 let mut iter = iter.enumerate();
348
349 let mut batch_size = Self::FIRST_PARALLEL_TX_WINDOW_SIZE;
350
351 let parent_span = Span::current();
354 executor.cpu_pool().install(move || {
355 let _parent = parent_span.enter();
356 let span = debug_span!(target: "engine::tree::payload_processor", "convert_transactions").or_current();
357 let _entered = span.enter();
358 loop {
359 let chunk = iter
360 .by_ref()
361 .take(batch_size)
362 .collect::<Vec<_>>();
363 if chunk.is_empty() {
364 break;
365 }
366
367 batch_size = batch_size.saturating_mul(2);
368
369 let chunk = chunk
370 .into_par_iter()
371 .map(|(i, tx)| {
372 let idx = i + prefetch;
373 let _parent = span.enter();
374 let _span = trace_span!(target: "engine::tree::payload_processor", "convert_transaction", idx).entered();
375 let tx = convert.convert(tx).map(WithTxEnv::new);
376 (idx, tx)
377 })
378 .collect::<Vec<_>>();
379
380 for (idx, tx) in chunk {
381 let failed = tx.is_err();
382 if let (Some(prewarm_tx), Ok(tx)) = (&prewarm_tx, &tx) {
383 let _ = prewarm_tx.send((idx, tx.clone()));
384 }
385 if execute_tx.send((idx, tx)).is_err() || failed {
386 return
387 }
388 trace!(target: "engine::tree::payload_processor", idx, "yielded transaction");
389 }
390 }
391 });
392 }
393 });
394 }
395
396 (prewarm_rx, execute_rx)
397 }
398
399 #[instrument(level = "debug", target = "engine::tree::payload_processor", skip_all)]
405 fn spawn_caching_with<P>(
406 &self,
407 env: ExecutionEnv<Evm>,
408 transactions: Option<
409 mpsc::Receiver<(usize, impl ExecutableTxFor<Evm> + Clone + Send + 'static)>,
410 >,
411 state_provider_factory: OverlayStateProviderFactory<P, Evm::Primitives>,
412 hint_stream: Option<StateRootHintStream>,
413 hashed_update_stream: Option<StateRootUpdateStream>,
414 parallel_bal_execution: bool,
415 ) -> CacheTaskHandle<<Evm::Primitives as NodePrimitives>::Receipt>
416 where
417 P: DatabaseProviderFactory + Clone + 'static,
418 P::Provider: BlockNumReader
419 + PruneCheckpointReader
420 + StageCheckpointReader
421 + ChangeSetReader
422 + StorageChangeSetReader
423 + StorageSettingsCache
424 + HistoryReader
425 + 'static,
426 {
427 let mode = if parallel_bal_execution {
430 PrewarmMode::BlockAccessList {
431 bal: env.decoded_bal.clone().expect("BAL dispatch implies decoded BAL"),
432 updates: hashed_update_stream,
433 }
434 } else if let Some(pending) = transactions {
435 PrewarmMode::Transactions { pending, hints: hint_stream }
436 } else {
437 PrewarmMode::Skipped
438 };
439 let saved_cache = self.disable_state_cache.not().then(|| self.cache_for(env.parent_hash));
440
441 let executed_tx_index = Arc::new(AtomicUsize::new(0));
442 let prewarm_ctx = PrewarmContext {
444 env,
445 evm_config: self.evm_config.clone(),
446 saved_cache: saved_cache.clone(),
447 provider: state_provider_factory,
448 bal_prewarm_pool: parallel_bal_execution.then(|| self.bal_prewarm_pool()),
449 metrics: PrewarmMetrics::default(),
450 cache_metrics: self.cache_metrics.clone(),
451 cache_state_metrics: self.cache_state_metrics.clone(),
452 terminate_execution: Arc::new(AtomicBool::new(false)),
453 executed_tx_index: Arc::clone(&executed_tx_index),
454 precompile_cache_disabled: self.precompile_cache_disabled,
455 precompile_cache_map: self.precompile_cache_map.clone(),
456 disable_bal_parallel_state_root: self.disable_bal_parallel_state_root,
457 disable_bal_batch_io: self.disable_bal_batch_io,
458 };
459
460 let (prewarm_task, to_prewarm_task) =
461 PrewarmCacheTask::new(self.executor.clone(), self.execution_cache.clone(), prewarm_ctx);
462 {
463 let to_prewarm_task = to_prewarm_task.clone();
464 self.executor.spawn_blocking_named("prewarm", move || {
465 prewarm_task.run(mode, to_prewarm_task);
466 });
467 }
468
469 CacheTaskHandle {
470 saved_cache,
471 to_prewarm_task: Some(to_prewarm_task),
472 executed_tx_index,
473 cache_metrics: self.cache_metrics.clone(),
474 }
475 }
476
477 #[instrument(level = "debug", target = "engine::caching", skip(self))]
482 pub fn cache_for(&self, parent_hash: B256) -> SavedCache {
483 if let Some(cache) = self.execution_cache.get_cache_for(parent_hash) {
484 debug!("reusing execution cache");
485 cache
486 } else {
487 debug!("creating new execution cache on cache miss");
488 let start = Instant::now();
489 let cache = ExecutionCache::new(self.cross_block_cache_size);
490 if let Some(metrics) = &self.cache_metrics {
491 metrics.record_cache_creation(start.elapsed());
492 }
493 SavedCache::new(parent_hash, cache)
494 }
495 }
496
497 pub fn on_inserted_executed_block(
505 &self,
506 block_with_parent: BlockWithParent,
507 bundle_state: &BundleState,
508 ) {
509 let cache_state_metrics = self.cache_state_metrics.clone();
510 self.execution_cache.update_with_guard(|cached| {
511 if cached.as_ref().is_some_and(|c| c.executed_block_hash() != block_with_parent.parent) {
512 debug!(
513 target: "engine::caching",
514 parent_hash = %block_with_parent.parent,
515 "Cannot find cache for parent hash, skip updating cache with new state for inserted executed block",
516 );
517 return
518 }
519
520 if let Some(cache) = cached.as_ref().filter(|cache| !cache.is_available()) {
521 debug!(
522 target: "engine::caching",
523 parent_hash = %block_with_parent.parent,
524 usage_count = cache.usage_count(),
525 "Execution cache is in use, skip updating cache with new state for inserted executed block",
526 );
527 return
528 }
529
530 let caches = match cached.take() {
532 Some(existing) => existing.cache().clone(),
533 None => ExecutionCache::new(self.cross_block_cache_size),
534 };
535
536 let new_cache = SavedCache::new(block_with_parent.block.hash, caches);
538 if new_cache.cache().insert_state(bundle_state).is_err() {
539 *cached = None;
540 debug!(target: "engine::caching", "cleared execution cache on update error");
541 return
542 }
543 new_cache.update_metrics(cache_state_metrics.as_ref());
544
545 *cached = Some(new_cache);
547 debug!(target: "engine::caching", ?block_with_parent, "Updated execution cache for inserted block");
548 });
549 }
550}
551
552fn convert_serial<RawTx, Tx, TxEnv, InnerTx, Recovered, Err, C>(
555 iter: impl Iterator<Item = RawTx>,
556 convert: &C,
557 prewarm_tx: Option<&mpsc::SyncSender<(usize, WithTxEnv<TxEnv, Recovered>)>>,
558 execute_tx: &ExecuteTxSender<TxEnv, Recovered, Err>,
559) -> bool
560where
561 Tx: ExecutableTxParts<TxEnv, InnerTx, Recovered = Recovered>,
562 TxEnv: Clone,
563 C: ConvertTx<RawTx, Tx = Tx, Error = Err>,
564{
565 for (idx, raw_tx) in iter.enumerate() {
566 let _span =
567 trace_span!(target: "engine::tree::payload_processor", "convert_transaction", idx)
568 .entered();
569 let tx = convert.convert(raw_tx);
570 let failed = tx.is_err();
571 let tx = tx.map(WithTxEnv::new);
572 if let (Some(prewarm_tx), Ok(tx)) = (prewarm_tx, &tx) {
573 let _ = prewarm_tx.send((idx, tx.clone()));
574 }
575 if execute_tx.send((idx, tx)).is_err() || failed {
576 return false
577 }
578 trace!(target: "engine::tree::payload_processor", idx, "yielded transaction");
579 }
580 true
581}
582
583#[derive(Debug)]
588pub struct PayloadHandle<Tx, Err, R> {
589 prewarm_handle: CacheTaskHandle<R>,
590 transactions: IndexedTxReceiver<Tx, Err>,
592 _span: Span,
594}
595
596impl<Tx, Err, R: Send + Sync + 'static> PayloadHandle<Tx, Err, R> {
597 pub fn caches(&self) -> Option<ExecutionCache> {
599 self.prewarm_handle.saved_cache.as_ref().map(|cache| cache.cache().clone())
600 }
601
602 pub fn cache_metrics(&self) -> Option<CachedStateMetrics> {
604 self.prewarm_handle.cache_metrics.clone()
605 }
606
607 pub const fn executed_tx_index(&self) -> &Arc<AtomicUsize> {
612 &self.prewarm_handle.executed_tx_index
613 }
614
615 pub fn stop_prewarming_execution(&self) {
619 self.prewarm_handle.stop_prewarming_execution()
620 }
621
622 pub fn terminate_caching(
630 &mut self,
631 execution_outcome: Option<Arc<BlockExecutionOutput<R>>>,
632 ) -> Option<mpsc::Sender<()>> {
633 self.prewarm_handle.terminate_caching(execution_outcome)
634 }
635
636 pub fn iter_transactions(&mut self) -> impl Iterator<Item = Result<Tx, Err>> + '_ {
638 self.transactions.iter().map(|(_, tx)| tx)
639 }
640
641 pub fn clone_transaction_receiver(&self) -> IndexedTxReceiver<Tx, Err> {
643 self.transactions.clone()
644 }
645}
646
647#[derive(Debug)]
652pub struct CacheTaskHandle<R> {
653 saved_cache: Option<SavedCache>,
655 to_prewarm_task: Option<std::sync::mpsc::Sender<PrewarmTaskEvent<R>>>,
657 executed_tx_index: Arc<AtomicUsize>,
660 cache_metrics: Option<CachedStateMetrics>,
662}
663
664impl<R: Send + Sync + 'static> CacheTaskHandle<R> {
665 pub fn stop_prewarming_execution(&self) {
669 self.to_prewarm_task
670 .as_ref()
671 .map(|tx| tx.send(PrewarmTaskEvent::TerminateTransactionExecution).ok());
672 }
673
674 #[must_use = "sender must be used and notified on block validation success"]
679 pub fn terminate_caching(
680 &mut self,
681 execution_outcome: Option<Arc<BlockExecutionOutput<R>>>,
682 ) -> Option<mpsc::Sender<()>> {
683 if let Some(tx) = self.to_prewarm_task.take() {
684 let (valid_block_tx, valid_block_rx) = mpsc::channel();
685 let event = PrewarmTaskEvent::Terminate { execution_outcome, valid_block_rx };
686 let _ = tx.send(event);
687
688 Some(valid_block_tx)
689 } else {
690 None
691 }
692 }
693}
694
695impl<R> Drop for CacheTaskHandle<R> {
696 fn drop(&mut self) {
697 if let Some(tx) = self.to_prewarm_task.take() {
699 let _ = tx.send(PrewarmTaskEvent::Terminate {
700 execution_outcome: None,
701 valid_block_rx: mpsc::channel().1,
702 });
703 }
704 }
705}
706
707#[cfg(test)]
708mod tests {
709 use super::*;
710 use crate::tree::{
711 payload_processor::PayloadProcessor, precompile_cache::PrecompileCacheMap, ExecutionCache,
712 PayloadExecutionCache, SavedCache, TreeConfig,
713 };
714 use alloy_consensus::constants::KECCAK_EMPTY;
715 use alloy_eips::eip1898::{BlockNumHash, BlockWithParent};
716 use alloy_primitives::{Address, B256, U256};
717 use proptest::{
718 prelude::any,
719 test_runner::{Config, TestRunner},
720 };
721 use reth_chainspec::ChainSpec;
722 use reth_evm_ethereum::EthEvmConfig;
723 use reth_execution_cache::CachedStatus;
724 use reth_revm::db::BundleState;
725 use revm::state::AccountInfo;
726 use std::{
727 collections::HashSet,
728 process::Command,
729 sync::{atomic::Ordering, Arc, Mutex},
730 };
731 use tracing_subscriber::{
732 filter::LevelFilter, layer::SubscriberExt, registry::LookupSpan, Registry,
733 };
734
735 type TestTx = reth_evm::execute::WithTxEnv<
736 reth_evm::TxEnvFor<EthEvmConfig>,
737 reth_primitives_traits::Recovered<reth_ethereum_primitives::TransactionSigned>,
738 >;
739
740 fn converted_tx() -> TestTx {
741 TestTx {
742 tx_env: Default::default(),
743 tx: Arc::new(reth_primitives_traits::Recovered::new_unchecked(
744 reth_ethereum_primitives::TransactionSigned::Legacy(
745 alloy_consensus::Signed::new_unchecked(
746 alloy_consensus::TxLegacy::default(),
747 alloy_primitives::Signature::test_signature(),
748 B256::ZERO,
749 ),
750 ),
751 Address::ZERO,
752 )),
753 }
754 }
755
756 fn test_processor() -> PayloadProcessor<EthEvmConfig> {
757 PayloadProcessor::new(
758 reth_tasks::Runtime::test(),
759 EthEvmConfig::new(Arc::new(ChainSpec::default())),
760 &TreeConfig::default(),
761 PrecompileCacheMap::default(),
762 )
763 }
764
765 #[test]
766 fn transaction_conversion_preserves_results() {
767 for (count, bal) in [(10, false), (200, false), (200, true)] {
768 let processor = test_processor();
769 let (prewarm, receiver) = processor.spawn_tx_iterator(
770 ((0..count).collect::<Vec<_>>(), |_| Ok::<_, std::io::Error>(converted_tx())),
771 count,
772 bal,
773 true,
774 );
775 let mut indices = Vec::new();
776 for _ in 0..count {
777 let (idx, tx) = receiver.recv_timeout(std::time::Duration::from_secs(10)).unwrap();
778 assert!(tx.is_ok());
779 indices.push(idx);
780 }
781 if bal {
782 indices.sort_unstable();
783 }
784 assert_eq!(indices, (0..count).collect::<Vec<_>>());
785 assert_eq!(prewarm.unwrap().iter().count(), count);
786 }
787 }
788
789 #[test]
790 fn transaction_conversion_stops_on_error() {
791 for (count, bal, fail_at) in
792 [(10, false, 0), (10, true, 0), (200, false, 0), (200, false, 4), (200, false, 20)]
793 {
794 let processor = test_processor();
795 let (_, receiver) = processor.spawn_tx_iterator(
796 ((0..count).collect::<Vec<_>>(), move |idx| {
797 if idx >= fail_at {
798 Err(std::io::Error::other("invalid transaction"))
799 } else {
800 Ok(converted_tx())
801 }
802 }),
803 count,
804 bal,
805 false,
806 );
807 let mut results = Vec::new();
808 loop {
809 match receiver.recv_timeout(std::time::Duration::from_secs(10)) {
810 Ok(tx) => results.push(tx),
811 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
812 Err(err) => panic!("conversion did not terminate: {err}"),
813 }
814 }
815 assert!(results.iter().any(|(_, tx)| tx.is_err()));
816 assert!(results.len() < count);
817 if !bal {
818 assert_eq!(results.len(), fail_at + 1);
819 assert!(results.last().unwrap().1.is_err());
820 }
821 }
822 }
823
824 #[test]
825 fn bal_transaction_conversion_preserves_all_error_slots() {
826 let count = 200;
827 let processor = test_processor();
828 let (_, receiver) = processor.spawn_tx_iterator(
829 ((0..count).collect::<Vec<_>>(), |idx| {
830 if idx % 3 == 0 {
831 Err(std::io::Error::other("invalid transaction"))
832 } else {
833 Ok(converted_tx())
834 }
835 }),
836 count,
837 true,
838 false,
839 );
840 let mut indices = Vec::new();
841 for _ in 0..count {
842 let (idx, tx) = receiver.recv_timeout(std::time::Duration::from_secs(10)).unwrap();
843 assert_eq!(tx.is_err(), idx % 3 == 0);
844 indices.push(idx);
845 }
846 indices.sort_unstable();
847 assert_eq!(indices, (0..count).collect::<Vec<_>>());
848 }
849
850 #[test]
851 fn dropping_payload_handle_stops_transaction_conversion() {
852 for (count, bal, pause_at) in
853 [(10, false, 0), (1000, false, 0), (1000, false, 4), (1000, true, 0)]
854 {
855 let processor = test_processor();
856 let calls = Arc::new(super::AtomicUsize::new(0));
857 let converted = calls.clone();
858 let (started_tx, started_rx) = crossbeam_channel::unbounded();
859 let (release_tx, release_rx) = crossbeam_channel::bounded::<()>(0);
860 let (_, receiver) = processor.spawn_tx_iterator(
861 ((0..count).collect::<Vec<_>>(), move |idx| {
862 converted.fetch_add(1, Ordering::Relaxed);
863 if idx >= pause_at {
864 started_tx.send(()).unwrap();
865 let _ = release_rx.recv_timeout(std::time::Duration::from_secs(10));
866 }
867 Ok::<_, std::io::Error>(converted_tx())
868 }),
869 count,
870 bal,
871 false,
872 );
873 let handle = super::PayloadHandle {
874 prewarm_handle: super::CacheTaskHandle::<()> {
875 saved_cache: None,
876 to_prewarm_task: None,
877 executed_tx_index: Default::default(),
878 cache_metrics: None,
879 },
880 transactions: receiver,
881 _span: tracing::Span::none(),
882 };
883 started_rx.recv_timeout(std::time::Duration::from_secs(10)).unwrap();
884 drop(handle);
885 drop(release_tx);
886
887 loop {
889 match started_rx.recv_timeout(std::time::Duration::from_secs(10)) {
890 Ok(()) => {}
891 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
892 Err(err) => panic!("conversion did not terminate: {err}"),
893 }
894 }
895 assert!(calls.load(Ordering::Relaxed) < count);
896 }
897 }
898
899 fn make_saved_cache(hash: B256) -> SavedCache {
900 let execution_cache = ExecutionCache::new(1_000);
901 SavedCache::new(hash, execution_cache)
902 }
903
904 #[test]
905 fn execution_cache_allows_single_checkout() {
906 let execution_cache = PayloadExecutionCache::default();
907 let hash = B256::from([1u8; 32]);
908
909 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(hash)));
910
911 let first = execution_cache.get_cache_for(hash);
912 assert!(first.is_some(), "expected initial checkout to succeed");
913
914 let second = execution_cache.get_cache_for(hash);
915 assert!(second.is_none(), "second checkout should be blocked while guard is active");
916
917 drop(first);
918
919 let third = execution_cache.get_cache_for(hash);
920 assert!(third.is_some(), "third checkout should succeed after guard is dropped");
921 }
922
923 #[test]
924 fn execution_cache_checkout_releases_on_drop() {
925 let execution_cache = PayloadExecutionCache::default();
926 let hash = B256::from([2u8; 32]);
927
928 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(hash)));
929
930 {
931 let guard = execution_cache.get_cache_for(hash);
932 assert!(guard.is_some(), "expected checkout to succeed");
933 }
935
936 let retry = execution_cache.get_cache_for(hash);
937 assert!(retry.is_some(), "checkout should succeed after guard drop");
938 }
939
940 #[test]
941 fn execution_cache_mismatch_parent_clears_and_returns() {
942 let execution_cache = PayloadExecutionCache::default();
943 let hash = B256::from([3u8; 32]);
944
945 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(hash)));
946
947 let different_hash = B256::from([4u8; 32]);
950 let cache = execution_cache.get_cache_for(different_hash);
951 assert!(cache.is_some(), "cache should be returned for reuse after clearing");
952
953 drop(cache);
954
955 let original = execution_cache.get_cache_for(hash);
958 assert!(original.is_some(), "canonical chain gets cache back via mismatch+clear");
959 }
960
961 #[test]
962 fn execution_cache_update_after_release_succeeds() {
963 let execution_cache = PayloadExecutionCache::default();
964 let initial = B256::from([5u8; 32]);
965
966 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(initial)));
967
968 let guard =
969 execution_cache.get_cache_for(initial).expect("expected initial checkout to succeed");
970
971 drop(guard);
972
973 let updated = B256::from([6u8; 32]);
974 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(updated)));
975
976 let new_checkout = execution_cache.get_cache_for(updated);
977 assert!(new_checkout.is_some(), "new checkout should succeed after release and update");
978 }
979
980 #[test]
981 fn on_inserted_executed_block_populates_cache() {
982 let payload_processor = PayloadProcessor::new(
983 reth_tasks::Runtime::test(),
984 EthEvmConfig::new(Arc::new(ChainSpec::default())),
985 &TreeConfig::default(),
986 PrecompileCacheMap::default(),
987 );
988
989 let parent_hash = B256::from([1u8; 32]);
990 let block_hash = B256::from([10u8; 32]);
991 let block_with_parent = BlockWithParent {
992 block: BlockNumHash { hash: block_hash, number: 1 },
993 parent: parent_hash,
994 };
995 let bundle_state = BundleState::default();
996
997 assert!(payload_processor.execution_cache.get_cache_for(block_hash).is_none());
999
1000 payload_processor.on_inserted_executed_block(block_with_parent, &bundle_state);
1002
1003 let cached = payload_processor.execution_cache.get_cache_for(block_hash);
1005 assert!(cached.is_some());
1006 assert_eq!(cached.unwrap().executed_block_hash(), block_hash);
1007 }
1008
1009 #[test]
1010 fn on_inserted_executed_block_skips_on_parent_mismatch() {
1011 let payload_processor = PayloadProcessor::new(
1012 reth_tasks::Runtime::test(),
1013 EthEvmConfig::new(Arc::new(ChainSpec::default())),
1014 &TreeConfig::default(),
1015 PrecompileCacheMap::default(),
1016 );
1017
1018 let block1_hash = B256::from([1u8; 32]);
1020 payload_processor
1021 .execution_cache
1022 .update_with_guard(|slot| *slot = Some(make_saved_cache(block1_hash)));
1023
1024 let wrong_parent = B256::from([99u8; 32]);
1026 let block3_hash = B256::from([3u8; 32]);
1027 let block_with_parent = BlockWithParent {
1028 block: BlockNumHash { hash: block3_hash, number: 3 },
1029 parent: wrong_parent,
1030 };
1031 let bundle_state = BundleState::default();
1032
1033 payload_processor.on_inserted_executed_block(block_with_parent, &bundle_state);
1034
1035 let cached = payload_processor.execution_cache.get_cache_for(block1_hash);
1037 assert!(cached.is_some(), "Original cache should be preserved");
1038
1039 let cached3 = payload_processor.execution_cache.get_cache_for(block3_hash);
1041 assert!(cached3.is_none(), "New block cache should not be created on mismatch");
1042 }
1043
1044 #[test]
1045 fn on_inserted_executed_block_does_not_mutate_checked_out_parent_cache() {
1046 let payload_processor = PayloadProcessor::new(
1047 reth_tasks::Runtime::test(),
1048 EthEvmConfig::new(Arc::new(ChainSpec::default())),
1049 &TreeConfig::default(),
1050 PrecompileCacheMap::default(),
1051 );
1052
1053 let parent_hash = B256::from([1u8; 32]);
1054 payload_processor
1055 .execution_cache
1056 .update_with_guard(|slot| *slot = Some(make_saved_cache(parent_hash)));
1057
1058 let checked_out = payload_processor
1062 .execution_cache
1063 .get_cache_for(parent_hash)
1064 .expect("expected parent cache checkout to succeed");
1065
1066 let polluted_address = Address::random();
1067 let bundle_state = BundleState::builder(2..=2)
1068 .state_present_account_info(
1069 polluted_address,
1070 AccountInfo {
1071 balance: U256::from(1337),
1072 nonce: 7,
1073 code_hash: KECCAK_EMPTY,
1074 code: None,
1075 ..Default::default()
1076 },
1077 )
1078 .build();
1079
1080 let block_with_parent = BlockWithParent {
1083 block: BlockNumHash { hash: B256::from([2u8; 32]), number: 2 },
1084 parent: parent_hash,
1085 };
1086
1087 payload_processor.on_inserted_executed_block(block_with_parent, &bundle_state);
1088
1089 let account = checked_out
1092 .cache()
1093 .get_or_try_insert_account_with(polluted_address, || Ok::<_, ()>(None))
1094 .expect("cache read should succeed");
1095
1096 assert_eq!(
1097 account,
1098 CachedStatus::NotCached(None),
1099 "checked-out parent cache should not observe state from inserted local block"
1100 );
1101 }
1102
1103 #[test]
1114 fn fork_prewarm_dropped_without_save_does_not_corrupt_cache() {
1115 let execution_cache = PayloadExecutionCache::default();
1116
1117 let block4_hash = B256::from([4u8; 32]);
1119 execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(block4_hash)));
1120
1121 let fork_parent = B256::from([2u8; 32]);
1124 let prewarm_cache = execution_cache.get_cache_for(fork_parent);
1125 assert!(prewarm_cache.is_some(), "prewarm should obtain cache for fork block");
1126 let prewarm_cache = prewarm_cache.unwrap();
1127 assert_eq!(prewarm_cache.executed_block_hash(), fork_parent);
1128
1129 let fork_addr = Address::from([0xBB; 20]);
1132 let fork_key = B256::from([0xCC; 32]);
1133 prewarm_cache.cache().insert_storage(fork_addr, fork_key, Some(U256::from(999)));
1134
1135 let during_prewarm = execution_cache.get_cache_for(block4_hash);
1137 assert!(
1138 during_prewarm.is_none(),
1139 "cache must be unavailable while prewarm holds a reference"
1140 );
1141
1142 drop(prewarm_cache);
1144
1145 let block5_cache = execution_cache.get_cache_for(block4_hash);
1149 assert!(
1150 block5_cache.is_some(),
1151 "canonical chain must get cache after fork prewarm is dropped"
1152 );
1153 assert_eq!(
1154 block5_cache.as_ref().unwrap().executed_block_hash(),
1155 block4_hash,
1156 "cache must carry the canonical parent hash, not the fork parent"
1157 );
1158 }
1159
1160 #[test]
1161 fn transaction_conversion_has_child_spans() {
1162 const LEVEL_ENV: &str = "RETH_TEST_CONVERSION_SPAN_LEVEL";
1165 let Ok(level) = std::env::var(LEVEL_ENV) else {
1166 for level in ["trace", "debug", "info"] {
1167 let output = Command::new(std::env::current_exe().unwrap())
1168 .args([
1169 "--exact",
1170 "tree::payload_processor::tests::transaction_conversion_has_child_spans",
1171 "--nocapture",
1172 ])
1173 .env(LEVEL_ENV, level)
1174 .output()
1175 .unwrap();
1176 assert!(
1177 output.status.success(),
1178 "{level}: {}\n{}",
1179 String::from_utf8_lossy(&output.stdout),
1180 String::from_utf8_lossy(&output.stderr),
1181 );
1182 }
1183 return;
1184 };
1185 let level = level.parse::<LevelFilter>().unwrap();
1186 tracing::subscriber::set_global_default(tracing_subscriber::registry().with(level))
1187 .unwrap();
1188 let processor = test_processor();
1189
1190 for count in [0, 1, 29, 30, 68, 69] {
1193 for bal in [false, true] {
1194 assert_conversion_spans(&processor, count, bal, level);
1195 }
1196 }
1197 TestRunner::new(Config::with_cases(32))
1198 .run(&(0..200usize, any::<bool>()), |(count, bal)| {
1199 assert_conversion_spans(&processor, count, bal, level);
1200 Ok(())
1201 })
1202 .unwrap();
1203 }
1204
1205 fn assert_conversion_spans(
1206 processor: &PayloadProcessor<EthEvmConfig>,
1207 count: usize,
1208 bal: bool,
1209 level: LevelFilter,
1210 ) {
1211 let request = tracing::info_span!(parent: None, "request");
1212 let _entered = request.enter();
1213 let observed = Arc::new(Mutex::new(Vec::new()));
1215 let converted = observed.clone();
1216 let (_, receiver) = processor.spawn_tx_iterator(
1217 ((0..count).collect::<Vec<_>>(), move |idx| {
1218 converted.lock().unwrap().push((idx, Span::current()));
1219 Ok::<_, std::io::Error>(converted_tx())
1220 }),
1221 count,
1222 bal,
1223 false,
1224 );
1225 let mut indices = Vec::new();
1226 loop {
1227 match receiver.recv_timeout(std::time::Duration::from_secs(10)) {
1228 Ok((idx, tx)) => {
1229 assert!(tx.is_ok());
1230 indices.push(idx);
1231 }
1232 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
1233 Err(err) => panic!("conversion did not terminate: {err}"),
1234 }
1235 }
1236 if bal {
1237 indices.sort_unstable();
1238 }
1239 assert_eq!(indices, (0..count).collect::<Vec<_>>());
1240 let observed = observed.lock().unwrap();
1241 assert_eq!(observed.len(), count);
1242 let mut transaction_ids = HashSet::new();
1243 for (idx, span) in observed.iter() {
1244 let scope = span
1245 .with_subscriber(|(id, dispatch)| {
1246 dispatch
1247 .downcast_ref::<Registry>()
1248 .unwrap()
1249 .span(id)
1250 .unwrap()
1251 .scope()
1252 .from_root()
1253 .map(|span| (span.id(), span.name()))
1254 .collect::<Vec<_>>()
1255 })
1256 .expect("conversion lost its request context");
1257 assert_eq!(Some(scope[0].0.clone()), request.id());
1258 let mut expected = vec!["request"];
1259 if level >= LevelFilter::DEBUG {
1260 expected.push("spawn_tx_iterator");
1261 expected.push("tx_iterator");
1262 if count >= 30 && (bal || *idx >= 4) {
1263 expected.push("convert_transactions");
1264 }
1265 }
1266 if level == LevelFilter::TRACE {
1267 expected.push("convert_transaction");
1268 assert!(transaction_ids.insert(span.id().unwrap()), "reused a transaction span");
1269 }
1270 assert_eq!(scope.iter().map(|(_, name)| *name).collect::<Vec<_>>(), expected);
1271 }
1272 }
1273}