Skip to main content

reth_rpc/eth/
bundle.rs

1//! `Eth` bundle implementation and helpers.
2
3use 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
28/// `Eth` bundle implementation.
29pub struct EthBundle<Eth> {
30    /// All nested fields bundled together.
31    inner: Arc<EthBundleInner<Eth>>,
32}
33
34impl<Eth> EthBundle<Eth> {
35    /// Create a new `EthBundle` instance.
36    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    /// Access the underlying `Eth` API.
41    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    /// Simulates a bundle of transactions at the top of a given block number with the state of
51    /// another (or the same) block. This can be used to simulate future blocks with the current
52    /// state, or it can be used to simulate a past block. The sender is responsible for signing the
53    /// transactions and using the correct nonce and ensuring validity
54    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        // Validate gas limit against the configured call gas limit before any DB calls
84        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        // Note: the block number is considered the `parent` block: <https://github.com/flashbots/mev-geth/blob/fddf97beec5877483f879a77b7dea2e58a58d653/internal/ethapi/api.go#L2104>
98        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        // need to adjust the timestamp for the next block
122        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        // Validate that the bundle does not contain more than MAX_BLOB_NUMBER_PER_BLOCK blob
133        // transactions.
134        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        // Apply gas limit: default to call gas limit unless user requests a smaller limit
152        evm_env.block_env.inner_mut().gas_limit = gas_limit.unwrap_or(call_gas_limit);
153
154        // The bundle is simulated on top of the state block, so default to the next block's base
155        // fee. <https://github.com/flashbots/mev-geth/blob/fddf97beec5877483f879a77b7dea2e58a58d653/internal/ethapi/api.go#L2130>
156        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        // use the block number of the request
169        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 is always present in the result state
225                    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                    // update the coinbase balance
232                    coinbase_balance_before_tx = coinbase_balance_after_tx;
233
234                    // set the return data for the response
235                    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                    // need to apply the state changes of this call before executing the
258                    // next call
259                    if transactions.peek().is_some() {
260                        // need to apply the state changes of this call before executing
261                        // the next call
262                        evm.db_mut().commit(state)
263                    }
264                }
265
266                // populate the response
267
268                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/// Container type for `EthBundle` internals
300#[derive(Debug)]
301struct EthBundleInner<Eth> {
302    /// Access to commonly used code of the `eth` namespace
303    eth_api: Eth,
304    // restrict the number of concurrent tracing calls.
305    #[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/// [`EthBundle`] specific errors.
322#[derive(Debug, thiserror::Error)]
323pub enum EthBundleError {
324    /// Thrown if the bundle does not contain any transactions.
325    #[error("bundle missing txs")]
326    EmptyBundleTransactions,
327    /// Thrown if the bundle does not contain a block number, or block number is 0.
328    #[error("bundle missing blockNumber")]
329    BundleMissingBlockNumber,
330    /// Thrown when the blob gas usage of the blob transactions in a bundle exceed the maximum.
331    #[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    /// The bundle is simulated in the block after the state block, so a transaction that only
358    /// covers the next block's base fee, but not the state block's, must be accepted.
359    #[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        // genesis is empty, so the next base fee is 7/8 of the initial base fee
364        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        // An actual pending block is not canonical and has its own base fee.
424        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}