1use alloy_consensus::{transaction::TxHashRef, EnvKzgSettings, Transaction as _};
4use alloy_eips::eip7840::BlobParams;
5use alloy_evm::env::BlockEnvironment;
6use alloy_primitives::{uint, Keccak256, U256};
7use alloy_rpc_types_mev::{EthCallBundle, EthCallBundleResponse, EthCallBundleTransactionResult};
8use jsonrpsee::core::RpcResult;
9use reth_chainspec::{ChainSpecProvider, EthChainSpec};
10use reth_evm::{ConfigureEvm, Evm};
11use reth_rpc_eth_api::{
12 helpers::{Call, EthTransactions, LoadPendingBlock},
13 EthCallBundleApiServer, FromEthApiError, FromEvmError,
14};
15use reth_rpc_eth_types::{
16 EthApiError, PendingBlockEnv, PendingBlockEnvOrigin, RpcInvalidTransactionError,
17};
18use reth_storage_api::BlockReaderIdExt;
19use reth_tasks::pool::BlockingTaskGuard;
20use reth_transaction_pool::{
21 EthBlobTransactionSidecar, EthPoolTransaction, PoolPooledTx, PoolTransaction, TransactionPool,
22};
23use revm::{
24 context::Block, context_interface::result::ResultAndState, DatabaseCommit, DatabaseRef,
25};
26use std::sync::Arc;
27
28pub struct EthBundle<Eth> {
30 inner: Arc<EthBundleInner<Eth>>,
32}
33
34impl<Eth> EthBundle<Eth> {
35 pub fn new(eth_api: Eth, blocking_task_guard: BlockingTaskGuard) -> Self {
37 Self { inner: Arc::new(EthBundleInner { eth_api, blocking_task_guard }) }
38 }
39
40 pub fn eth_api(&self) -> &Eth {
42 &self.inner.eth_api
43 }
44}
45
46impl<Eth> EthBundle<Eth>
47where
48 Eth: EthTransactions + LoadPendingBlock + Call + 'static,
49{
50 pub async fn call_bundle(
55 &self,
56 bundle: EthCallBundle,
57 ) -> Result<EthCallBundleResponse, Eth::Error> {
58 let EthCallBundle {
59 txs,
60 block_number,
61 coinbase,
62 state_block_number,
63 timeout: _,
64 timestamp,
65 gas_limit,
66 difficulty,
67 base_fee,
68 ..
69 } = bundle;
70 if txs.is_empty() {
71 return Err(EthApiError::InvalidParams(
72 EthBundleError::EmptyBundleTransactions.to_string(),
73 )
74 .into())
75 }
76 if block_number == 0 {
77 return Err(EthApiError::InvalidParams(
78 EthBundleError::BundleMissingBlockNumber.to_string(),
79 )
80 .into())
81 }
82
83 let call_gas_limit = self.inner.eth_api.call_gas_limit();
85 if let Some(gas_limit) = gas_limit &&
86 gas_limit > call_gas_limit
87 {
88 return Err(
89 EthApiError::InvalidTransaction(RpcInvalidTransactionError::GasTooHigh).into()
90 )
91 }
92
93 let transactions =
94 self.eth_api().recover_raw_transactions::<PoolPooledTx<Eth::Pool>>(txs)?;
95
96 let block_id: alloy_rpc_types_eth::BlockId = state_block_number.into();
97 let (mut evm_env, at, parent) = if block_id.is_pending() {
99 let PendingBlockEnv { evm_env, origin } = self.eth_api().pending_block_env_and_cfg()?;
100 let at = origin.state_block_id();
101 let parent = match origin {
102 PendingBlockEnvOrigin::ActualPending(block, _) => block.clone_sealed_header(),
103 PendingBlockEnvOrigin::DerivedFromLatest(header) => header,
104 };
105 (evm_env, at, parent)
106 } else {
107 let parent = self
108 .eth_api()
109 .provider()
110 .sealed_header_by_id(block_id)
111 .map_err(Eth::Error::from_eth_err)?
112 .ok_or(EthApiError::HeaderNotFound(block_id))?;
113 let evm_env = self.eth_api().evm_env_for_header(&parent)?;
114 (evm_env, parent.hash().into(), parent)
115 };
116
117 if let Some(coinbase) = coinbase {
118 evm_env.block_env.inner_mut().beneficiary = coinbase;
119 }
120
121 if let Some(timestamp) = timestamp {
123 evm_env.block_env.inner_mut().timestamp = U256::from(timestamp);
124 } else {
125 evm_env.block_env.inner_mut().timestamp += uint!(12_U256);
126 }
127
128 if let Some(difficulty) = difficulty {
129 evm_env.block_env.inner_mut().difficulty = U256::from(difficulty);
130 }
131
132 let blob_gas_used = transactions.iter().filter_map(|tx| tx.blob_gas_used()).sum::<u64>();
135 if blob_gas_used > 0 {
136 let blob_params = self
137 .eth_api()
138 .provider()
139 .chain_spec()
140 .blob_params_at_timestamp(evm_env.block_env.timestamp().saturating_to())
141 .unwrap_or_else(BlobParams::cancun);
142 if blob_gas_used > blob_params.max_blob_gas_per_block() {
143 return Err(EthApiError::InvalidParams(
144 EthBundleError::Eip4844BlobGasExceeded(blob_params.max_blob_gas_per_block())
145 .to_string(),
146 )
147 .into())
148 }
149 }
150
151 evm_env.block_env.inner_mut().gas_limit = gas_limit.unwrap_or(call_gas_limit);
153
154 if let Some(base_fee) = base_fee {
157 evm_env.block_env.inner_mut().basefee = base_fee.try_into().unwrap_or(u64::MAX);
158 } else if let Some(next_base_fee) = self
159 .eth_api()
160 .provider()
161 .chain_spec()
162 .next_block_base_fee(&parent, evm_env.block_env.timestamp().saturating_to())
163 {
164 evm_env.block_env.inner_mut().basefee = next_base_fee;
165 }
166
167 let state_block_number = evm_env.block_env.number();
168 evm_env.block_env.inner_mut().number = U256::from(block_number);
170
171 self.eth_api()
172 .spawn_with_state_at_block(at, move |eth_api, db| {
173 let coinbase = evm_env.block_env.beneficiary();
174 let basefee = evm_env.block_env.basefee();
175
176 let initial_coinbase = db
177 .basic_ref(coinbase)
178 .map_err(Eth::Error::from_eth_err)?
179 .map(|acc| acc.balance)
180 .unwrap_or_default();
181 let mut coinbase_balance_before_tx = initial_coinbase;
182 let mut coinbase_balance_after_tx = initial_coinbase;
183 let mut total_gas_used = 0u64;
184 let mut total_gas_fees = U256::ZERO;
185 let mut hasher = Keccak256::new();
186
187 let mut evm = eth_api.evm_config().evm_with_env(db, evm_env);
188
189 let mut results = Vec::with_capacity(transactions.len());
190 let mut transactions = transactions.into_iter().peekable();
191
192 while let Some(tx) = transactions.next() {
193 let signer = tx.signer();
194 let tx = {
195 let mut tx = <Eth::Pool as TransactionPool>::Transaction::from_pooled(tx);
196
197 if let EthBlobTransactionSidecar::Present(sidecar) = tx.take_blob() {
198 tx.validate_blob(&sidecar, EnvKzgSettings::Default.get()).map_err(
199 |e| {
200 Eth::Error::from_eth_err(EthApiError::InvalidParams(
201 e.to_string(),
202 ))
203 },
204 )?;
205 }
206
207 tx.into_consensus()
208 };
209
210 hasher.update(*tx.tx_hash());
211 let ResultAndState { result, state } = evm
212 .transact(eth_api.evm_config().tx_env(&tx))
213 .map_err(Eth::Error::from_evm_err)?;
214
215 let gas_price = tx
216 .effective_tip_per_gas(basefee)
217 .expect("fee is always valid; execution succeeded");
218 let gas_used = result.tx_gas_used();
219 total_gas_used += gas_used;
220
221 let gas_fees = U256::from(gas_used) * U256::from(gas_price);
222 total_gas_fees += gas_fees;
223
224 coinbase_balance_after_tx =
226 state.get(&coinbase).map(|acc| acc.info.balance).unwrap_or_default();
227 let coinbase_diff =
228 coinbase_balance_after_tx.saturating_sub(coinbase_balance_before_tx);
229 let eth_sent_to_coinbase = coinbase_diff.saturating_sub(gas_fees);
230
231 coinbase_balance_before_tx = coinbase_balance_after_tx;
233
234 let (value, revert) = if result.is_success() {
236 let value = result.into_output().unwrap_or_default();
237 (Some(value), None)
238 } else {
239 let revert = result.into_output().unwrap_or_default();
240 (None, Some(revert))
241 };
242
243 let tx_res = EthCallBundleTransactionResult {
244 coinbase_diff,
245 eth_sent_to_coinbase,
246 from_address: signer,
247 gas_fees,
248 gas_price: U256::from(gas_price),
249 gas_used,
250 to_address: tx.to(),
251 tx_hash: *tx.tx_hash(),
252 value,
253 revert,
254 };
255 results.push(tx_res);
256
257 if transactions.peek().is_some() {
260 evm.db_mut().commit(state)
263 }
264 }
265
266 let coinbase_diff = coinbase_balance_after_tx.saturating_sub(initial_coinbase);
269 let eth_sent_to_coinbase = coinbase_diff.saturating_sub(total_gas_fees);
270 let bundle_gas_price =
271 coinbase_diff.checked_div(U256::from(total_gas_used)).unwrap_or_default();
272 let res = EthCallBundleResponse {
273 bundle_gas_price,
274 bundle_hash: hasher.finalize(),
275 coinbase_diff,
276 eth_sent_to_coinbase,
277 gas_fees: total_gas_fees,
278 results,
279 state_block_number: state_block_number.to(),
280 total_gas_used,
281 };
282
283 Ok(res)
284 })
285 .await
286 }
287}
288
289#[async_trait::async_trait]
290impl<Eth> EthCallBundleApiServer for EthBundle<Eth>
291where
292 Eth: EthTransactions + LoadPendingBlock + Call + 'static,
293{
294 async fn call_bundle(&self, request: EthCallBundle) -> RpcResult<EthCallBundleResponse> {
295 Self::call_bundle(self, request).await.map_err(Into::into)
296 }
297}
298
299#[derive(Debug)]
301struct EthBundleInner<Eth> {
302 eth_api: Eth,
304 #[expect(dead_code)]
306 blocking_task_guard: BlockingTaskGuard,
307}
308
309impl<Eth> std::fmt::Debug for EthBundle<Eth> {
310 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
311 f.debug_struct("EthBundle").finish_non_exhaustive()
312 }
313}
314
315impl<Eth> Clone for EthBundle<Eth> {
316 fn clone(&self) -> Self {
317 Self { inner: Arc::clone(&self.inner) }
318 }
319}
320
321#[derive(Debug, thiserror::Error)]
323pub enum EthBundleError {
324 #[error("bundle missing txs")]
326 EmptyBundleTransactions,
327 #[error("bundle missing blockNumber")]
329 BundleMissingBlockNumber,
330 #[error("blob gas usage exceeds the limit of {0} gas per block.")]
332 Eip4844BlobGasExceeded(u64),
333}
334
335#[cfg(test)]
336mod tests {
337 use super::*;
338 use crate::EthApiBuilder;
339 use alloy_consensus::{transaction::SignerRecoverable, TxEip1559};
340 use alloy_eips::{eip1559::INITIAL_BASE_FEE, eip2718::Encodable2718, BlockNumberOrTag};
341 use alloy_genesis::{Genesis, GenesisAccount};
342 use alloy_primitives::{Address, TxKind};
343 use reth_chain_state::ExecutedBlock;
344 use reth_chainspec::ChainSpecBuilder;
345 use reth_db_common::init::init_genesis;
346 use reth_ethereum_primitives::Block as EthBlock;
347 use reth_evm_ethereum::EthEvmConfig;
348 use reth_network_api::noop::NoopNetwork;
349 use reth_primitives_traits::RecoveredBlock;
350 use reth_provider::{
351 providers::BlockchainProvider, test_utils::create_test_provider_factory_with_chain_spec,
352 };
353 use reth_storage_overlay::OverlayManager;
354 use reth_testing_utils::generators::{self, generate_key, sign_tx_with_key_pair};
355 use reth_transaction_pool::test_utils::testing_pool;
356
357 #[tokio::test(flavor = "multi_thread")]
360 async fn call_bundle_uses_next_block_base_fee() {
361 let mut rng = generators::rng();
362 let key = generate_key(&mut rng);
363 let next_base_fee = INITIAL_BASE_FEE - INITIAL_BASE_FEE / 8;
365 let tx = sign_tx_with_key_pair(
366 key,
367 reth_ethereum_primitives::Transaction::Eip1559(TxEip1559 {
368 chain_id: 1,
369 nonce: 0,
370 gas_limit: 21_000,
371 max_fee_per_gas: (next_base_fee + 1) as u128,
372 max_priority_fee_per_gas: (next_base_fee + 1) as u128,
373 to: TxKind::Call(Address::random()),
374 ..Default::default()
375 }),
376 );
377 let sender = tx.recover_signer().unwrap();
378
379 let genesis = Genesis::default()
380 .with_gas_limit(30_000_000)
381 .with_base_fee(Some(INITIAL_BASE_FEE as u128))
382 .extend_accounts([(
383 sender,
384 GenesisAccount::default().with_balance(U256::from(10u128.pow(18))),
385 )]);
386 let chain_spec =
387 Arc::new(ChainSpecBuilder::mainnet().cancun_activated().genesis(genesis).build());
388 let overlay_manager = OverlayManager::default();
389 let factory = create_test_provider_factory_with_chain_spec(chain_spec)
390 .with_overlay_manager(overlay_manager.clone());
391 init_genesis(&factory).unwrap();
392 let provider = BlockchainProvider::new(factory).unwrap();
393
394 let eth_api = EthApiBuilder::new(
395 provider.clone(),
396 testing_pool(),
397 NoopNetwork::default(),
398 EthEvmConfig::new(provider.chain_spec()),
399 )
400 .build();
401 let api = EthBundle::new(eth_api, BlockingTaskGuard::new(4));
402
403 let bundle = EthCallBundle {
404 txs: vec![tx.encoded_2718().into()],
405 block_number: 1,
406 state_block_number: BlockNumberOrTag::Number(0),
407 ..Default::default()
408 };
409 for state_block_number in
410 [BlockNumberOrTag::Number(0), BlockNumberOrTag::Latest, BlockNumberOrTag::Pending]
411 {
412 for base_fee in [None, Some((next_base_fee - 1) as u128)] {
413 let response = api
414 .call_bundle(EthCallBundle { state_block_number, base_fee, ..bundle.clone() })
415 .await
416 .unwrap();
417 let expected_tip =
418 (next_base_fee + 1) as u128 - base_fee.unwrap_or(next_base_fee as u128);
419 assert_bundle_fees(&response, expected_tip);
420 }
421 }
422
423 let chain_spec = provider.chain_spec();
425 let mut header = chain_spec.genesis_header().clone();
426 header.parent_hash = chain_spec.genesis_hash();
427 header.number = 1;
428 header.timestamp = 12;
429 header.base_fee_per_gas = Some(next_base_fee);
430 let pending_block = ExecutedBlock {
431 recovered_block: Arc::new(RecoveredBlock::new_unhashed(
432 EthBlock { header, body: Default::default() },
433 vec![],
434 )),
435 ..Default::default()
436 };
437 overlay_manager.insert_block(pending_block.clone());
438 provider.canonical_in_memory_state().set_pending_block(pending_block);
439 let pending_next_base_fee = next_base_fee - next_base_fee / 8;
440 for base_fee in [None, Some((pending_next_base_fee - 1) as u128)] {
441 let response = api
442 .call_bundle(EthCallBundle {
443 block_number: 2,
444 state_block_number: BlockNumberOrTag::Pending,
445 base_fee,
446 ..bundle.clone()
447 })
448 .await
449 .unwrap();
450 let expected_tip =
451 (next_base_fee + 1) as u128 - base_fee.unwrap_or(pending_next_base_fee as u128);
452 assert_bundle_fees(&response, expected_tip);
453 }
454 }
455
456 fn assert_bundle_fees(response: &EthCallBundleResponse, expected_tip: u128) {
457 assert_eq!(response.results.len(), 1);
458 assert_eq!(response.total_gas_used, 21_000);
459 let gas_price = U256::from(expected_tip);
460 let gas_fees = gas_price * U256::from(21_000);
461 assert_eq!(response.results[0].gas_price, gas_price);
462 assert_eq!(response.results[0].gas_fees, gas_fees);
463 assert_eq!(response.gas_fees, gas_fees);
464 assert_eq!(response.bundle_gas_price, gas_price);
465 assert_eq!(response.coinbase_diff, gas_fees);
466 }
467}