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 for transaction {tx_index}: {source}")]
21 Transaction {
22 tx_index: usize,
24 #[source]
26 source: Box<dyn core::error::Error + Send + Sync + 'static>,
27 },
28 #[error("BAL worker EVM execution failed for transaction {tx_index}: {source}")]
30 Execution {
31 tx_index: usize,
33 tx_gas_limit: u64,
35 #[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 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 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}