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: {0}")]
21    Transaction(Box<dyn core::error::Error + Send + Sync + 'static>),
22    /// EVM transaction execution failed.
23    #[error("BAL worker EVM execution failed: {0}")]
24    Execution(BlockExecutionError),
25}
26
27impl From<BalWorkerError> for BalExecutionError {
28    fn from(err: BalWorkerError) -> Self {
29        match err {
30            BalWorkerError::Setup(err) => err,
31            BalWorkerError::Transaction(err) => {
32                Self::Execution(BlockValidationError::Other(err).into())
33            }
34            BalWorkerError::Execution(err) => Self::Execution(err),
35        }
36    }
37}
38
39pub(super) struct BalWorkerOutput<R> {
40    pub(super) index: usize,
41    pub(super) signer: Address,
42    pub(super) tx_gas_limit: u64,
43    pub(super) result: R,
44}
45
46type WorkerExecutorResult<Cfg> =
47    <<Cfg as ConfigureEvm>::BlockExecutorFactory as BlockExecutorFactory>::TxExecutionResult;
48
49type WorkerResultSender<Cfg> =
50    Sender<Result<BalWorkerOutput<WorkerExecutorResult<Cfg>>, BalWorkerError>>;
51
52#[expect(clippy::too_many_arguments)]
53pub(super) fn spawn_worker<'scope, Evm, Tx, Err, DB, MakeDb>(
54    scope: &rayon::Scope<'scope>,
55    tx_rx: Receiver<(usize, Result<Tx, Err>)>,
56    abort_rx: Receiver<()>,
57    result_tx: WorkerResultSender<Evm>,
58    evm_config: &'scope Evm,
59    make_db: &'scope MakeDb,
60    received_bal_revm: Arc<RevmBal>,
61    evm_env: EvmEnvFor<Evm>,
62    ctx: ExecutionCtxFor<'scope, Evm>,
63) where
64    Evm: ConfigureEvm + 'scope,
65    Tx: ExecutableTxFor<Evm> + Send + 'scope,
66    Err: core::error::Error + Send + Sync + 'static,
67    DB: Database + Send + 'scope,
68    MakeDb: Fn(bool) -> Result<DB, BalExecutionError> + Sync + 'scope,
69{
70    scope.spawn(move |_| {
71        let worker_result = (|| -> Result<(), BalWorkerError> {
72            // Create a database with fill_on_miss=true ensuring misses
73            // are inserted for the other workers.
74            let database = make_db(true).map_err(BalWorkerError::Setup)?;
75            let mut worker_state = State::builder()
76                .with_database(database)
77                .with_bal(received_bal_revm)
78                .with_bundle_update()
79                .build();
80            let evm = evm_config.evm_with_env(&mut worker_state, evm_env);
81            let mut executor = evm_config.create_executor_with_state(evm, ctx.clone());
82
83            loop {
84                let (index, tx) = crossbeam_channel::select_biased! {
85                    recv(abort_rx) -> _ => break,
86                    recv(tx_rx) -> msg => match msg {
87                        Ok(ix_tx) => ix_tx,
88                        Err(_) => break,
89                    },
90                };
91                let tx = tx.map_err(|e| BalWorkerError::Transaction(Box::new(e)))?;
92                let signer = *tx.signer();
93                let tx_gas_limit = tx.tx().gas_limit();
94
95                executor
96                    .evm_mut()
97                    .db_mut()
98                    .set_bal_index(BlockAccessIndex::from_tx_index(index as u64));
99                let result = executor
100                    .execute_transaction_without_commit(tx)
101                    .map_err(BalWorkerError::Execution)?;
102
103                if result_tx
104                    .send(Ok(BalWorkerOutput { index, signer, tx_gas_limit, result }))
105                    .is_err()
106                {
107                    break;
108                }
109            }
110
111            Ok(())
112        })();
113
114        if let Err(err) = worker_result {
115            let _ = result_tx.send(Err(err));
116        }
117    });
118}