Skip to main content

reth_engine_tree/tree/payload_processor/
mod.rs

1//! Entrypoint for payload processing.
2
3use 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
49/// Blocks with fewer transactions than this skip prewarming, since the fixed overhead of spawning
50/// prewarm workers exceeds the execution time saved.
51pub const SMALL_BLOCK_TX_THRESHOLD: usize = 5;
52
53/// Type alias for [`PayloadHandle`] returned by payload processor spawn methods.
54type 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/// Entrypoint for executing the payload.
81#[derive(Debug)]
82pub struct PayloadProcessor<Evm>
83where
84    Evm: ConfigureEvm,
85{
86    /// The executor used by to spawn tasks.
87    executor: Runtime,
88    /// The most recent cache used for execution.
89    execution_cache: PayloadExecutionCache,
90    /// Metrics for the execution cache.
91    cache_metrics: Option<CachedStateMetrics>,
92    /// Metrics for shared execution cache state.
93    cache_state_metrics: Option<CachedStateCacheMetrics>,
94    /// Cross-block cache size in bytes.
95    cross_block_cache_size: usize,
96    /// Whether transactions should not be executed on prewarming task.
97    disable_transaction_prewarming: bool,
98    /// Whether state cache should be disable
99    disable_state_cache: bool,
100    /// Determines how to configure the evm for execution.
101    evm_config: Evm,
102    /// Whether precompile cache should be disabled.
103    precompile_cache_disabled: bool,
104    /// Precompile cache map.
105    precompile_cache_map: PrecompileCacheMap<SpecFor<Evm>>,
106    /// Whether to disable BAL-driven parallel state root computation.
107    /// Only valid when BAL parallel execution is also disabled.
108    disable_bal_parallel_state_root: bool,
109    /// Whether BAL state prefetching during prewarm is disabled.
110    disable_bal_batch_io: bool,
111    /// Dedicated blocking pool for warming the BAL read-set, created lazily on the first BAL block
112    /// (see [`Self::bal_prewarm_pool`]). Its threads exit when the processor is dropped.
113    bal_prewarm_pool: OnceLock<Arc<bal_prewarm_pool::BalPrewarmPool>>,
114}
115
116impl<Evm> PayloadProcessor<Evm>
117where
118    Evm: ConfigureEvm,
119{
120    /// Creates a new payload processor.
121    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    /// Returns the dedicated BAL read-set prewarm pool, spawning its blocking worker threads on
147    /// first use (only the BAL parallel execution path calls this).
148    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    /// Returns the shared execution cache handle used for engine backpressure.
157    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    /// Spawns transaction conversion and cache prewarming, optionally wiring prewarm output into
167    /// an externally-owned state-root task.
168    #[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    /// Whether the prewarm task will consume converted transactions, i.e. whether
209    /// [`Self::spawn_caching_with`] ends up in [`PrewarmMode::Transactions`].
210    ///
211    /// This is the only place the decision is made: the tx iterator uses it to skip creating the
212    /// prewarm channel and cloning every transaction into it, and the resulting `Option` receiver
213    /// then selects the mode, so the two cannot disagree.
214    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    /// Transaction count threshold below which sequential conversion is used.
225    ///
226    /// For blocks with fewer than this many transactions, the rayon parallel iterator overhead
227    /// (work-stealing setup, channel-based reorder) exceeds the cost of sequential conversion.
228    /// Inspired by Nethermind's `RecoverSignature` which uses sequential `foreach` for small
229    /// blocks.
230    const SMALL_BLOCK_TX_THRESHOLD: usize = 30;
231
232    /// Number of leading transactions to convert sequentially before entering the rayon
233    /// parallel path.
234    ///
235    /// Rayon's work-stealing does not guarantee that index 0 is processed first, so the
236    /// ordered consumer can block for up to ~1ms waiting for the first slot. By converting
237    /// a small head sequentially and sending it immediately, execution can start without
238    /// waiting for rayon scheduling.
239    const PARALLEL_PREFETCH_COUNT: usize = 4;
240
241    /// Size of the first parallel tx batch in non-BAL path.
242    const FIRST_PARALLEL_TX_WINDOW_SIZE: usize = 64;
243
244    /// Spawns a task advancing transaction env iterator and streaming updates through a channel.
245    ///
246    /// For blocks with fewer than [`Self::SMALL_BLOCK_TX_THRESHOLD`] transactions, uses
247    /// sequential iteration to avoid rayon overhead. For larger blocks, uses rayon parallel
248    /// iteration to convert transactions in parallel while streaming results to execution.
249    ///
250    /// When `parallel_bal_execution` is disabled, preserves the original transaction order.
251    /// Otherwise, streams results as they become available.
252    ///
253    /// The prewarm channel is only created when `prewarm_transactions` is set, see
254    /// [`Self::prewarms_transactions`]; otherwise no transaction is cloned into it.
255    #[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            // Empty block — nothing to do.
270        } else if transaction_count < Self::SMALL_BLOCK_TX_THRESHOLD {
271            // Sequential path for small blocks — avoids rayon work-stealing setup and
272            // channel-based reorder overhead when it costs more than sequential conversion.
273            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                // Restore the parent even when the worker span is filtered out.
280                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            // Parallel path — recover signatures in parallel on rayon, stream results
292            // to prewarming and execution.
293            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                    // With BALs, we don't care about the order of transactions in execution and
300                    // prewarming, so we don't have to use `for_each_ordered_in`.
301                    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                            // Feed the ordered commit loop before recovering distant transactions.
310                            .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                                // Recovery failures are adjudicated in transaction order by
321                                // the BAL commit loop. Earlier slots may still be unconverted.
322                                if disconnected {
323                                    Err(())
324                                } else {
325                                    Ok(())
326                                }
327                            });
328                    });
329                } else {
330                    // To avoid a ~1ms stall waiting for rayon to schedule index 0, the first
331                    // few transactions are recovered sequentially and sent immediately before
332                    // entering the parallel iterator for the remainder.
333                    let prefetch = Self::PARALLEL_PREFETCH_COUNT.min(transaction_count);
334                    let mut iter = transactions.into_iter();
335
336                    // Convert the first few transactions sequentially so execution can
337                    // start immediately without waiting for rayon work-stealing.
338                    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                    // Without BALs, we need to preserve the initial order of transactions.
352                    // Process exponentially increasing windows to make sure that first transactions are prioritized.
353                    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    /// Spawn prewarming optionally wired to the sparse trie task for target updates.
400    ///
401    /// `parallel_bal_execution` is true when the BAL execute path will execute this block. In
402    /// that case prewarm runs in BAL mode: it streams BAL-derived sparse-trie updates and,
403    /// unless `disable_bal_batch_io` is set, prefetches BAL-declared state into the shared cache.
404    #[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        // Each mode carries the capability its producers use; the rest is dropped here, so
428        // unused capabilities do not keep the state-root task's update channel open.
429        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        // configure prewarming
443        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    /// Returns the cache for the given parent hash.
478    ///
479    /// If the given hash is different then what is recently cached, then this will create a new
480    /// instance.
481    #[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    /// Updates the execution cache with the post-execution state from an inserted block.
498    ///
499    /// This is used when blocks are inserted directly (e.g., locally built blocks by sequencers)
500    /// to ensure the cache remains warm for subsequent block execution.
501    ///
502    /// The cache enables subsequent blocks to reuse account, storage, and bytecode data without
503    /// hitting the database, maintaining performance consistency.
504    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            // Take existing cache (if any) or create fresh caches
531            let caches = match cached.take() {
532                Some(existing) => existing.cache().clone(),
533                None => ExecutionCache::new(self.cross_block_cache_size),
534            };
535
536            // Insert the block's bundle state into cache
537            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            // Replace with the updated cache
546            *cached = Some(new_cache);
547            debug!(target: "engine::caching", ?block_with_parent, "Updated execution cache for inserted block");
548        });
549    }
550}
551
552/// Converts transactions sequentially and sends them to the execute channel, and to the prewarm
553/// channel if there is one. Returns false on conversion failure or disconnection.
554fn 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/// Handle to all the spawned tasks.
584///
585/// Generic over `R` (receipt type) to allow sharing `Arc<ExecutionOutcome<R>>` with the
586/// caching task without cloning the expensive `BundleState`.
587#[derive(Debug)]
588pub struct PayloadHandle<Tx, Err, R> {
589    prewarm_handle: CacheTaskHandle<R>,
590    /// Stream of block transactions and their indices in the block.
591    transactions: IndexedTxReceiver<Tx, Err>,
592    /// Span for tracing
593    _span: Span,
594}
595
596impl<Tx, Err, R: Send + Sync + 'static> PayloadHandle<Tx, Err, R> {
597    /// Returns a clone of the caches used by prewarming
598    pub fn caches(&self) -> Option<ExecutionCache> {
599        self.prewarm_handle.saved_cache.as_ref().map(|cache| cache.cache().clone())
600    }
601
602    /// Returns engine cache metrics if a cache exists for prewarming.
603    pub fn cache_metrics(&self) -> Option<CachedStateMetrics> {
604        self.prewarm_handle.cache_metrics.clone()
605    }
606
607    /// Returns a reference to the shared executed transaction index counter.
608    ///
609    /// The main execution loop should store `index + 1` after executing each transaction so that
610    /// prewarm workers can skip transactions that have already been processed.
611    pub const fn executed_tx_index(&self) -> &Arc<AtomicUsize> {
612        &self.prewarm_handle.executed_tx_index
613    }
614
615    /// Terminates the pre-warming transaction processing.
616    ///
617    /// Note: This does not terminate the task yet.
618    pub fn stop_prewarming_execution(&self) {
619        self.prewarm_handle.stop_prewarming_execution()
620    }
621
622    /// Terminates the entire caching task.
623    ///
624    /// If the [`BlockExecutionOutput`] is provided it will update the shared cache using its
625    /// bundle state. Using `Arc<ExecutionOutcome>` allows sharing with the main execution
626    /// path without cloning the expensive `BundleState`.
627    ///
628    /// Returns a sender for the channel that should be notified on block validation success.
629    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    /// Returns iterator yielding transactions from the stream.
637    pub fn iter_transactions(&mut self) -> impl Iterator<Item = Result<Tx, Err>> + '_ {
638        self.transactions.iter().map(|(_, tx)| tx)
639    }
640
641    /// Returns a clone of the indexed transaction receiver.
642    pub fn clone_transaction_receiver(&self) -> IndexedTxReceiver<Tx, Err> {
643        self.transactions.clone()
644    }
645}
646
647/// Access to the spawned [`PrewarmCacheTask`].
648///
649/// Generic over `R` (receipt type) to allow sharing `Arc<ExecutionOutcome<R>>` with the
650/// prewarm task without cloning the expensive `BundleState`.
651#[derive(Debug)]
652pub struct CacheTaskHandle<R> {
653    /// The shared cache the task operates with.
654    saved_cache: Option<SavedCache>,
655    /// Channel to the spawned prewarm task if any
656    to_prewarm_task: Option<std::sync::mpsc::Sender<PrewarmTaskEvent<R>>>,
657    /// Shared counter tracking the next transaction index to be executed by the main execution
658    /// loop. Prewarm workers skip transactions below this index.
659    executed_tx_index: Arc<AtomicUsize>,
660    /// Metrics for the execution cache.
661    cache_metrics: Option<CachedStateMetrics>,
662}
663
664impl<R: Send + Sync + 'static> CacheTaskHandle<R> {
665    /// Terminates the pre-warming transaction processing.
666    ///
667    /// Note: This does not terminate the task yet.
668    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    /// Terminates the entire pre-warming task.
675    ///
676    /// If the [`BlockExecutionOutput`] is provided it will update the shared cache using its
677    /// bundle state. Using `Arc<ExecutionOutcome>` avoids cloning the expensive `BundleState`.
678    #[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        // Ensure we always terminate on drop - send None without needing Send + Sync bounds
698        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            // The converter owns the last sender, so disconnection confirms recovery exited.
888            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            // Guard dropped at end of scope
934        }
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        // When the parent hash doesn't match (fork block), the cache is cleared,
948        // hash updated on the original, and clone returned for reuse
949        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        // The stored cache now has the fork block's parent hash.
956        // Canonical chain looking for original hash sees a mismatch → clears and reuses.
957        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        // Cache should be empty initially
998        assert!(payload_processor.execution_cache.get_cache_for(block_hash).is_none());
999
1000        // Update cache with inserted block
1001        payload_processor.on_inserted_executed_block(block_with_parent, &bundle_state);
1002
1003        // Cache should now exist for the block hash
1004        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        // Setup: populate cache with block 1
1019        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        // Try to insert block 3 with wrong parent (should skip and keep block 1's cache)
1025        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        // Cache should still be for block 1 (unchanged)
1036        let cached = payload_processor.execution_cache.get_cache_for(block1_hash);
1037        assert!(cached.is_some(), "Original cache should be preserved");
1038
1039        // Cache for block 3 should not exist
1040        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        // Checking out the cache bumps its `ExecutionCache` refcount, marking the slot as in-use.
1059        // The returned SavedCache shares the same underlying ExecutionCache Arc as the slot,
1060        // so any writes through the slot are observable here
1061        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        // Make parent match the cached slot so we bypass the parent-mismatch guard and exercise
1081        // the in-use guard specifically.
1082        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        // The closure runs only on a cache miss, so NotCached(None) means polluted_address was
1090        // absent and Cached(Some(_)) means it was written by on_inserted_executed_block.
1091        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    /// Tests the full prewarm lifecycle for a fork block:
1104    ///
1105    /// 1. Cache is at canonical block 4.
1106    /// 2. Fork block (parent = block 2) checks out the cache via `get_cache_for`, simulating what
1107    ///    `PrewarmCacheTask` does when it receives a `SavedCache`.
1108    /// 3. Prewarm populates the shared cache with fork-specific state.
1109    /// 4. While the prewarm clone is alive, the cache is unavailable (`usage_count` > 1).
1110    /// 5. Prewarm drops without calling `save_cache` (fork block was invalid).
1111    /// 6. Canonical block 5 (parent = block 4) must get a cache with correct hash and no stale fork
1112    ///    data.
1113    #[test]
1114    fn fork_prewarm_dropped_without_save_does_not_corrupt_cache() {
1115        let execution_cache = PayloadExecutionCache::default();
1116
1117        // Canonical chain at block 4.
1118        let block4_hash = B256::from([4u8; 32]);
1119        execution_cache.update_with_guard(|slot| *slot = Some(make_saved_cache(block4_hash)));
1120
1121        // Fork block arrives with parent = block 2. Prewarm task checks out the cache.
1122        // This simulates PrewarmCacheTask receiving a SavedCache clone from get_cache_for.
1123        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        // Prewarm populates cache with fork-specific state (ancestor data for block 2).
1130        // Since ExecutionCache uses Arc<Inner>, this data is shared with the stored original.
1131        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        // While prewarm holds the clone, the cache handle count > 1 so the cache is in use.
1136        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        // Fork block fails — prewarm task drops without calling save_cache/update_with_guard.
1143        drop(prewarm_cache);
1144
1145        // Canonical block 5 arrives (parent = block 4).
1146        // Stored hash = fork_parent (our fix), so get_cache_for sees a mismatch,
1147        // clears the stale fork data, and returns a cache with hash = block4_hash.
1148        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        // Worker threads use the global subscriber. Isolate it from other tests and run each
1163        // filter level in a fresh process so filtered child spans also exercise context fallback.
1164        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        // Exercise both sides of the sequential cutoff and the first parallel window, as well
1191        // as randomly sized blocks. BAL execution bypasses the sequential prefetch/windowing.
1192        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        // Retain the spans so registry IDs cannot be recycled between conversions.
1214        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}