Skip to main content

reth_engine_tree/tree/payload_processor/bal/
worker.rs

1use super::BalExecutionError;
2use alloy_consensus::Transaction;
3use alloy_eip7928::BlockAccessIndex;
4use alloy_evm::{
5    block::{BlockExecutionError, BlockExecutor, BlockExecutorFactory, BlockValidationError},
6    Evm,
7};
8use alloy_primitives::Address;
9use crossbeam_channel::{Receiver, Sender};
10use reth_evm::{execute::ExecutableTxFor, ConfigureEvm, Database, EvmEnvFor, ExecutionCtxFor};
11use revm::{database::State, state::bal::Bal as RevmBal};
12use std::sync::Arc;
13
14#[derive(Debug, thiserror::Error)]
15pub(super) enum BalWorkerError {
16    /// Worker state or provider setup failed.
17    #[error("BAL worker setup failed: {0}")]
18    Setup(#[source] BalExecutionError),
19    /// Transaction recovery or conversion failed before EVM execution.
20    #[error("BAL worker transaction conversion failed for transaction {tx_index}: {source}")]
21    Transaction {
22        /// Index of the transaction that failed.
23        tx_index: usize,
24        /// The underlying recovery or conversion error.
25        #[source]
26        source: Box<dyn core::error::Error + Send + Sync + 'static>,
27    },
28    /// EVM transaction execution failed.
29    #[error("BAL worker EVM execution failed for transaction {tx_index}: {source}")]
30    Execution {
31        /// Index of the transaction that failed.
32        tx_index: usize,
33        /// Gas limit of the transaction that failed.
34        tx_gas_limit: u64,
35        /// The underlying execution error.
36        #[source]
37        source: BlockExecutionError,
38    },
39}
40
41impl From<BalWorkerError> for BalExecutionError {
42    fn from(err: BalWorkerError) -> Self {
43        match err {
44            BalWorkerError::Setup(err) => err,
45            BalWorkerError::Transaction { source, .. } => {
46                Self::Execution(BlockValidationError::Other(source).into())
47            }
48            BalWorkerError::Execution { source, .. } => Self::Execution(source),
49        }
50    }
51}
52
53pub(super) struct BalWorkerOutput<R> {
54    pub(super) index: usize,
55    pub(super) signer: Address,
56    pub(super) tx_gas_limit: u64,
57    pub(super) result: R,
58}
59
60type WorkerExecutorResult<Cfg> =
61    <<Cfg as ConfigureEvm>::BlockExecutorFactory as BlockExecutorFactory>::TxExecutionResult;
62
63type WorkerResultSender<Cfg> =
64    Sender<Result<BalWorkerOutput<WorkerExecutorResult<Cfg>>, BalWorkerError>>;
65
66#[expect(clippy::too_many_arguments)]
67pub(super) fn spawn_worker<'scope, Evm, Tx, Err, DB, MakeDb>(
68    scope: &rayon::Scope<'scope>,
69    tx_rx: Receiver<(usize, Result<Tx, Err>)>,
70    abort_rx: Receiver<()>,
71    result_tx: WorkerResultSender<Evm>,
72    evm_config: &'scope Evm,
73    make_db: &'scope MakeDb,
74    received_bal_revm: Arc<RevmBal>,
75    evm_env: EvmEnvFor<Evm>,
76    ctx: ExecutionCtxFor<'scope, Evm>,
77) where
78    Evm: ConfigureEvm + 'scope,
79    Tx: ExecutableTxFor<Evm> + Send + 'scope,
80    Err: core::error::Error + Send + Sync + 'static,
81    DB: Database + Send + 'scope,
82    MakeDb: Fn(bool) -> Result<DB, BalExecutionError> + Sync + 'scope,
83{
84    scope.spawn(move |_| {
85        let worker_result = (|| -> Result<(), BalWorkerError> {
86            // Keep the cache-filling database across executor resets so a speculative failure
87            // cannot introduce an unindexed provider setup error ahead of its ordered verdict.
88            let mut database = make_db(true).map_err(BalWorkerError::Setup)?;
89            'worker: loop {
90                let mut worker_state = State::builder()
91                    .with_database(&mut database)
92                    .with_bal(Arc::clone(&received_bal_revm))
93                    .with_bundle_update()
94                    .build();
95                let evm = evm_config.evm_with_env(&mut worker_state, evm_env.clone());
96                let mut executor = evm_config.create_executor_with_state(evm, ctx.clone());
97
98                loop {
99                    let (tx_index, tx) = crossbeam_channel::select_biased! {
100                        recv(abort_rx) -> _ => break 'worker,
101                        recv(tx_rx) -> msg => match msg {
102                            Ok(ix_tx) => ix_tx,
103                            Err(_) => break 'worker,
104                        },
105                    };
106                    let tx = match tx {
107                        Ok(tx) => tx,
108                        Err(source) => {
109                            let error =
110                                BalWorkerError::Transaction { tx_index, source: Box::new(source) };
111                            if result_tx.send(Err(error)).is_err() {
112                                break 'worker;
113                            }
114                            continue;
115                        }
116                    };
117                    let signer = *tx.signer();
118                    let tx_gas_limit = tx.tx().gas_limit();
119
120                    executor
121                        .evm_mut()
122                        .db_mut()
123                        .set_bal_index(BlockAccessIndex::from_tx_index(tx_index as u64));
124                    let message = match executor.execute_transaction_without_commit(tx) {
125                        Ok(result) => {
126                            Ok(BalWorkerOutput { index: tx_index, signer, tx_gas_limit, result })
127                        }
128                        Err(source) => {
129                            Err(BalWorkerError::Execution { tx_index, tx_gas_limit, source })
130                        }
131                    };
132                    let failed = message.is_err();
133                    if result_tx.send(message).is_err() {
134                        break 'worker;
135                    }
136                    if failed {
137                        // The executor trait does not guarantee reuse after an error. Rebuild
138                        // its EVM and state before serving more work: the queue can still contain
139                        // earlier transactions whose verdict must precede this failure.
140                        break;
141                    }
142                }
143            }
144
145            Ok(())
146        })();
147
148        if let Err(err) = worker_result {
149            let _ = result_tx.send(Err(err));
150        }
151    });
152}