reth_engine_tree/tree/payload_processor/bal/
worker.rs1use 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 #[error("BAL worker setup failed: {0}")]
18 Setup(#[source] BalExecutionError),
19 #[error("BAL worker transaction conversion failed: {0}")]
21 Transaction(Box<dyn core::error::Error + Send + Sync + 'static>),
22 #[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 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}