Skip to main content

reth_rpc_eth_api/helpers/
state.rs

1//! Loads a pending block from database. Helper trait for `eth_` block, transaction, call and trace
2//! RPC methods.
3
4use super::{
5    pending_block::PendingStateSource, EthApiSpec, LoadBlock, LoadPendingBlock, SpawnBlocking,
6};
7use crate::{EthApiTypes, FromEthApiError, RpcNodeCore, RpcNodeCoreExt};
8use alloy_consensus::{constants::KECCAK_EMPTY, BlockHeader};
9use alloy_eips::BlockId;
10use alloy_primitives::{keccak256, Address, Bytes, B256, U256};
11use alloy_rpc_types_eth::{Account, AccountInfo, EIP1186AccountProofResponse};
12use alloy_serde::JsonStorageKey;
13use futures::Future;
14use reth_errors::RethError;
15use reth_evm::{ConfigureEvm, EvmEnvFor};
16use reth_primitives_traits::{BlockTy, RecoveredBlock, SealedHeaderFor};
17use reth_rpc_convert::{RpcConvert, RpcTxReq};
18use reth_rpc_eth_types::{
19    error::{FromEvmError, IntoEthApiError},
20    EthApiError, PendingBlockEnv, RpcInvalidTransactionError, SignError,
21};
22use reth_rpc_server_types::constants::DEFAULT_MAX_STORAGE_VALUES_SLOTS;
23use reth_storage_api::{
24    BlockIdReader, BlockReaderIdExt, StateProvider, StateProviderBox, StateProviderFactory,
25};
26use reth_transaction_pool::TransactionPool;
27use reth_trie_common::{MultiProofTargetsV2, ProofV2Target};
28use std::{collections::HashMap, sync::Arc};
29
30/// Helper methods for `eth_` methods relating to state (accounts).
31pub trait EthState: LoadState + SpawnBlocking {
32    /// Returns the maximum number of blocks into the past for generating state proofs.
33    fn max_proof_window(&self) -> u64 {
34        self.eth_api_settings().eth_proof_window
35    }
36
37    /// Validates that the given block is within the configured proof window.
38    ///
39    /// Returns an error if the distance between the chain tip and the requested block exceeds
40    /// [`Self::max_proof_window`].
41    fn ensure_within_proof_window(&self, block_id: BlockId) -> Result<(), Self::Error>
42    where
43        Self: EthApiSpec,
44    {
45        let chain_info = self.chain_info().map_err(Self::Error::from_eth_err)?;
46        let block_number = self
47            .provider()
48            .block_number_for_id(block_id)
49            .map_err(Self::Error::from_eth_err)?
50            .ok_or(EthApiError::HeaderNotFound(block_id))?;
51        if chain_info.best_number.saturating_sub(block_number) > self.max_proof_window() {
52            return Err(EthApiError::ExceedsMaxProofWindow.into())
53        }
54        Ok(())
55    }
56
57    /// Returns the number of transactions sent from an address at the given block identifier.
58    ///
59    /// If this is [`BlockNumberOrTag::Pending`](alloy_eips::BlockNumberOrTag) then this will
60    /// look up the highest transaction in pool and return the next nonce (highest + 1).
61    fn transaction_count(
62        &self,
63        address: Address,
64        block_id: Option<BlockId>,
65    ) -> impl Future<Output = Result<U256, Self::Error>> + Send {
66        LoadState::transaction_count(self, address, block_id)
67    }
68
69    /// Returns code of given account, at given blocknumber.
70    fn get_code(
71        &self,
72        address: Address,
73        block_id: Option<BlockId>,
74    ) -> impl Future<Output = Result<Bytes, Self::Error>> + Send {
75        LoadState::get_code(self, address, block_id)
76    }
77
78    /// Returns balance of given account, at given blocknumber.
79    fn balance(
80        &self,
81        address: Address,
82        block_id: Option<BlockId>,
83    ) -> impl Future<Output = Result<U256, Self::Error>> + Send {
84        self.spawn_blocking_io_with_state(block_id.unwrap_or_default(), move |_, state| {
85            Ok(state
86                .account_balance(&address)
87                .map_err(Self::Error::from_eth_err)?
88                .unwrap_or_default())
89        })
90    }
91
92    /// Returns values stored of given account, at given blocknumber.
93    fn storage_at(
94        &self,
95        address: Address,
96        index: JsonStorageKey,
97        block_id: Option<BlockId>,
98    ) -> impl Future<Output = Result<B256, Self::Error>> + Send {
99        self.spawn_blocking_io_with_state(block_id.unwrap_or_default(), move |_, state| {
100            Ok(B256::new(
101                state
102                    .storage(address, index.as_b256())
103                    .map_err(Self::Error::from_eth_err)?
104                    .unwrap_or_default()
105                    .to_be_bytes(),
106            ))
107        })
108    }
109
110    /// Returns values from multiple storage positions across multiple addresses.
111    ///
112    /// Enforces a cap on total slot count (sum of all slot arrays) and returns an error if
113    /// exceeded.
114    fn storage_values(
115        &self,
116        requests: HashMap<Address, Vec<JsonStorageKey>>,
117        block_id: Option<BlockId>,
118    ) -> impl Future<Output = Result<HashMap<Address, Vec<B256>>, Self::Error>> + Send {
119        async move {
120            if requests.is_empty() {
121                return Err(Self::Error::from_eth_err(EthApiError::InvalidParams(
122                    "empty request".to_string(),
123                )));
124            }
125            let total_slots: usize = requests.values().map(|slots| slots.len()).sum();
126            if total_slots > DEFAULT_MAX_STORAGE_VALUES_SLOTS {
127                return Err(Self::Error::from_eth_err(EthApiError::InvalidParams(
128                    format!(
129                        "total slot count {total_slots} exceeds limit {DEFAULT_MAX_STORAGE_VALUES_SLOTS}",
130                    ),
131                )));
132            }
133
134            self.spawn_blocking_io_with_state(block_id.unwrap_or_default(), move |_, state| {
135                let mut result = HashMap::with_capacity(requests.len());
136                for (address, slots) in requests {
137                    let mut values = Vec::with_capacity(slots.len());
138                    for slot in &slots {
139                        let value = state
140                            .storage(address, slot.as_b256())
141                            .map_err(Self::Error::from_eth_err)?
142                            .unwrap_or_default();
143                        values.push(B256::new(value.to_be_bytes()));
144                    }
145                    result.insert(address, values);
146                }
147
148                Ok(result)
149            })
150            .await
151        }
152    }
153
154    /// Returns values stored of given account, with Merkle-proof, at given blocknumber.
155    fn get_proof(
156        &self,
157        address: Address,
158        keys: Vec<JsonStorageKey>,
159        block_id: Option<BlockId>,
160    ) -> Result<
161        impl Future<Output = Result<EIP1186AccountProofResponse, Self::Error>> + Send,
162        Self::Error,
163    >
164    where
165        Self: EthApiSpec,
166    {
167        Ok(async move {
168            let permit = self
169                .acquire_owned_tracing()
170                .await
171                .map_err(RethError::other)
172                .map_err(EthApiError::Internal)?;
173
174            let block_id = block_id.unwrap_or_default();
175            self.ensure_within_proof_window(block_id)?;
176
177            self.spawn_blocking_io_with_state(block_id, move |_, state| {
178                let _permit = permit;
179                let storage_keys = keys.iter().map(|key| key.as_b256()).collect::<Vec<_>>();
180                let proof = state
181                    .proof(Default::default(), address, &storage_keys)
182                    .map_err(Self::Error::from_eth_err)?;
183                Ok(proof.into_eip1186_response(keys))
184            })
185            .await
186        })
187    }
188
189    /// Returns account and storage proofs for multiple targets at the given block number.
190    fn get_multi_proof(
191        &self,
192        targets: Vec<(Address, Vec<B256>)>,
193        block_id: Option<BlockId>,
194    ) -> Result<
195        impl Future<Output = Result<Vec<EIP1186AccountProofResponse>, Self::Error>> + Send,
196        Self::Error,
197    >
198    where
199        Self: EthApiSpec,
200    {
201        Ok(async move {
202            let permit = self
203                .acquire_owned_tracing()
204                .await
205                .map_err(RethError::other)
206                .map_err(EthApiError::Internal)?;
207
208            let block_id = block_id.unwrap_or_default();
209            self.ensure_within_proof_window(block_id)?;
210
211            self.spawn_blocking_io_with_state(block_id, move |_, state| {
212                let _permit = permit;
213                let mut proof_targets = MultiProofTargetsV2::default();
214                proof_targets.account_targets.reserve(targets.len());
215                proof_targets.storage_targets.reserve(targets.len());
216                for (address, slots) in &targets {
217                    let hashed_address = keccak256(address);
218                    proof_targets.account_targets.push(ProofV2Target::new(hashed_address));
219                    proof_targets
220                        .storage_targets
221                        .entry(hashed_address)
222                        .or_default()
223                        .extend(slots.iter().map(|slot| ProofV2Target::new(keccak256(slot))));
224                }
225
226                let multiproof = state
227                    .multiproof_v2(Default::default(), proof_targets)
228                    .map_err(Self::Error::from_eth_err)?;
229
230                targets
231                    .into_iter()
232                    .map(|(address, slots)| {
233                        let proof = multiproof
234                            .account_proof(address, &slots)
235                            .map_err(RethError::other)
236                            .map_err(Self::Error::from_eth_err)?;
237                        let storage_keys =
238                            slots.into_iter().map(JsonStorageKey::from).collect::<Vec<_>>();
239                        Ok(proof.into_eip1186_response(storage_keys))
240                    })
241                    .collect::<Result<Vec<_>, Self::Error>>()
242            })
243            .await
244        })
245    }
246
247    /// Returns the account at the given address for the provided block identifier.
248    fn get_account(
249        &self,
250        address: Address,
251        block_id: BlockId,
252    ) -> impl Future<Output = Result<Option<Account>, Self::Error>> + Send
253    where
254        Self: EthApiSpec,
255    {
256        async move {
257            self.ensure_within_proof_window(block_id)?;
258
259            self.spawn_blocking_io_with_state(block_id, move |_, state| {
260                let account = state.basic_account(&address).map_err(Self::Error::from_eth_err)?;
261                let Some(account) = account else { return Ok(None) };
262
263                // Provide a default `HashedStorage` value in order to
264                // get the storage root hash of the current state.
265                let storage_root = state
266                    .storage_root(address, Default::default())
267                    .map_err(Self::Error::from_eth_err)?;
268
269                Ok(Some(account.into_trie_account(storage_root)))
270            })
271            .await
272        }
273    }
274
275    /// Retrieves the account's balance, nonce, and code for a given address.
276    #[allow(clippy::needless_update)]
277    fn get_account_info(
278        &self,
279        address: Address,
280        block_id: BlockId,
281    ) -> impl Future<Output = Result<AccountInfo, Self::Error>> + Send {
282        self.spawn_blocking_io_with_state(block_id, move |_, state| {
283            let account = state
284                .basic_account(&address)
285                .map_err(Self::Error::from_eth_err)?
286                .unwrap_or_default();
287
288            let balance = account.balance;
289            let nonce = account.nonce;
290            let code = if account.get_bytecode_hash() == KECCAK_EMPTY {
291                Default::default()
292            } else {
293                state
294                    .account_code(&address)
295                    .map_err(Self::Error::from_eth_err)?
296                    .unwrap_or_default()
297                    .original_bytes()
298            };
299
300            Ok(AccountInfo {
301                balance,
302                nonce,
303                code,
304                #[cfg(feature = "account-ext")]
305                extension: account.extension,
306                ..Default::default()
307            })
308        })
309    }
310}
311
312/// Loads state from database.
313///
314/// Behaviour shared by several `eth_` RPC methods, not exclusive to `eth_` state RPC methods.
315pub trait LoadState:
316    LoadPendingBlock
317    + EthApiTypes<
318        Error: FromEvmError<Self::Evm> + FromEthApiError,
319        RpcConvert: RpcConvert<Network = Self::NetworkTypes>,
320    > + RpcNodeCoreExt
321{
322    /// Returns the state at the given block number
323    fn state_at_hash(&self, block_hash: B256) -> Result<StateProviderBox, Self::Error> {
324        self.provider().history_by_block_hash(block_hash).map_err(Self::Error::from_eth_err)
325    }
326
327    /// Returns the state at the given [`BlockId`], preferring locally built state for `pending`.
328    ///
329    /// If local pending state is unavailable or fails to build, falls back to the provider's
330    /// pending state. Other block IDs resolve only canonical state. See
331    /// <https://github.com/paradigmxyz/reth/issues/4515>.
332    ///
333    /// Pending block construction may spawn blocking work, so await this outside a blocking task.
334    /// The fallback provider lookup is synchronous on the calling task; RPC handlers should use
335    /// [`Self::spawn_blocking_io_with_state`] to keep that read on the blocking pool.
336    fn state_at_block_id(
337        &self,
338        at: BlockId,
339    ) -> impl Future<Output = Result<StateProviderBox, Self::Error>> + Send
340    where
341        Self: SpawnBlocking,
342    {
343        async move {
344            if at.is_pending() &&
345                let Ok(Some(state)) = self.local_pending_state().await
346            {
347                return Ok(state)
348            }
349
350            self.provider().state_by_block_id(at).map_err(Self::Error::from_eth_err)
351        }
352    }
353
354    /// Returns the _latest_ state
355    fn latest_state(&self) -> Result<StateProviderBox, Self::Error> {
356        self.provider().latest().map_err(Self::Error::from_eth_err)
357    }
358
359    /// Returns the state at the given [`BlockId`] enum or the latest.
360    ///
361    /// Convenience function to interprets `None` as `BlockId::Number(BlockNumberOrTag::Latest)`
362    fn state_at_block_id_or_latest(
363        &self,
364        block_id: Option<BlockId>,
365    ) -> impl Future<Output = Result<StateProviderBox, Self::Error>> + Send
366    where
367        Self: SpawnBlocking,
368    {
369        async move {
370            if let Some(block_id) = block_id {
371                self.state_at_block_id(block_id).await
372            } else {
373                Ok(self.latest_state()?)
374            }
375        }
376    }
377
378    /// Executes `f` with the state at the given [`BlockId`] on a blocking IO task.
379    ///
380    /// Resolves the chain's pending-state source before occupying a blocking thread, since
381    /// building a pending block may need a blocking thread of its own.
382    fn spawn_blocking_io_with_state<F, R>(
383        &self,
384        at: BlockId,
385        f: F,
386    ) -> impl Future<Output = Result<R, Self::Error>> + Send
387    where
388        Self: SpawnBlocking,
389        F: FnOnce(Self, StateProviderBox) -> Result<R, Self::Error> + Send + 'static,
390        R: Send + 'static,
391    {
392        async move {
393            let pending = if at.is_pending() {
394                self.local_pending_block_or_state().await.ok().flatten()
395            } else {
396                None
397            };
398
399            self.spawn_blocking_io(move |this| {
400                let state = match pending {
401                    Some(PendingStateSource::Block(pending)) => this
402                        .provider()
403                        .state_with_block_appended(
404                            pending.block().parent_hash(),
405                            pending.executed_block,
406                        )
407                        .or_else(|_| this.provider().state_by_block_id(BlockId::pending()))
408                        .map_err(Self::Error::from_eth_err)?,
409                    Some(PendingStateSource::State(state)) => state,
410                    None if at.is_latest() => this.latest_state()?,
411                    None => {
412                        this.provider().state_by_block_id(at).map_err(Self::Error::from_eth_err)?
413                    }
414                };
415                f(this, state)
416            })
417            .await
418        }
419    }
420
421    /// Returns the EVM environment for the given sealed header.
422    fn evm_env_for_header(
423        &self,
424        header: &SealedHeaderFor<Self::Primitives>,
425    ) -> Result<EvmEnvFor<Self::Evm>, Self::Error> {
426        self.evm_config()
427            .evm_env(header)
428            .map_err(RethError::other)
429            .map_err(Self::Error::from_eth_err)
430    }
431
432    /// Returns the EVM environment for the requested [`BlockId`]
433    ///
434    /// If the [`BlockId`] this will return the [`BlockId`] of the block the env was configured
435    /// for.
436    /// If the [`BlockId`] is pending, this will return the "Pending" tag, otherwise this returns
437    /// the hash of the exact block.
438    fn evm_env_at(
439        &self,
440        at: BlockId,
441    ) -> impl Future<Output = Result<(EvmEnvFor<Self::Evm>, BlockId), Self::Error>> + Send
442    where
443        Self: SpawnBlocking,
444    {
445        async move {
446            if at.is_pending() {
447                let PendingBlockEnv { evm_env, origin } = self.pending_block_env_and_cfg()?;
448                Ok((evm_env, origin.state_block_id()))
449            } else {
450                // we can assume that the blockid will be predominantly `Latest` (e.g. for
451                // `eth_call`) and if requested by number or hash we can quickly fetch just the
452                // header
453                let header = RpcNodeCore::provider(self)
454                    .sealed_header_by_id(at)
455                    .map_err(Self::Error::from_eth_err)?
456                    .ok_or_else(|| EthApiError::HeaderNotFound(at))?;
457                let evm_env = self.evm_env_for_header(&header)?;
458
459                Ok((evm_env, header.hash().into()))
460            }
461        }
462    }
463
464    /// Returns the recovered block, EVM environment, and state block id for the requested
465    /// [`BlockId`].
466    ///
467    /// For pending blocks, this preserves the state id returned by [`Self::evm_env_at`], which can
468    /// be the pending tag for an actual pending block or the latest block hash when the pending env
469    /// is derived from latest.
470    #[expect(clippy::type_complexity)]
471    fn evm_env_and_recovered_block_at(
472        &self,
473        at: BlockId,
474    ) -> impl Future<
475        Output = Result<
476            (Arc<RecoveredBlock<BlockTy<Self::Primitives>>>, EvmEnvFor<Self::Evm>, BlockId),
477            Self::Error,
478        >,
479    > + Send
480    where
481        Self: SpawnBlocking + LoadBlock,
482    {
483        async move {
484            if at.is_pending() {
485                let (evm_env, block_id) = self.evm_env_at(at).await?;
486                let block = self
487                    .recovered_block(block_id)
488                    .await?
489                    .ok_or_else(|| EthApiError::HeaderNotFound(at))?;
490
491                Ok((block, evm_env, block_id))
492            } else {
493                let block = self
494                    .recovered_block(at)
495                    .await?
496                    .ok_or_else(|| EthApiError::HeaderNotFound(at))?;
497                let evm_env = self.evm_env_for_header(block.sealed_block().sealed_header())?;
498                let block_id = block.hash().into();
499
500                Ok((block, evm_env, block_id))
501            }
502        }
503    }
504
505    /// Returns the next available nonce without gaps for the given address
506    /// Next available nonce is either the on chain nonce of the account or the highest consecutive
507    /// nonce in the pool + 1
508    ///
509    /// The provided request must have a from address set.
510    fn next_available_nonce_for(
511        &self,
512        request: &RpcTxReq<Self::NetworkTypes>,
513    ) -> impl Future<Output = Result<u64, Self::Error>> + Send
514    where
515        Self: SpawnBlocking,
516    {
517        let address = request.as_ref().from;
518        self.spawn_blocking_io(move |this| {
519            let address = match address {
520                Some(address) => address,
521                None => return Err(SignError::NoAccount.into_eth_err()),
522            };
523
524            // first fetch the on chain nonce of the account
525            let mut next_nonce = this
526                .latest_state()?
527                .account_nonce(&address)
528                .map_err(Self::Error::from_eth_err)?
529                .unwrap_or_default();
530
531            // Retrieve the highest consecutive transaction for the sender from the transaction pool
532            if let Some(highest_tx) =
533                this.pool().get_highest_consecutive_transaction_by_sender(address, next_nonce)
534            {
535                // Return the nonce of the highest consecutive transaction + 1
536                next_nonce = highest_tx.nonce().checked_add(1).ok_or_else(|| {
537                    Self::Error::from(EthApiError::InvalidTransaction(
538                        RpcInvalidTransactionError::NonceMaxValue,
539                    ))
540                })?;
541            }
542
543            Ok(next_nonce)
544        })
545    }
546
547    /// Returns the number of transactions sent from an address at the given block identifier.
548    ///
549    /// If this is [`BlockNumberOrTag::Pending`](alloy_eips::BlockNumberOrTag) then this will
550    /// look up the highest transaction in pool and return the next nonce (highest + 1).
551    fn transaction_count(
552        &self,
553        address: Address,
554        block_id: Option<BlockId>,
555    ) -> impl Future<Output = Result<U256, Self::Error>> + Send
556    where
557        Self: SpawnBlocking,
558    {
559        let at = block_id.unwrap_or_default();
560        self.spawn_blocking_io_with_state(at, move |this, state| {
561            // first fetch the on chain nonce of the account
562            let on_chain_account_nonce = state
563                .account_nonce(&address)
564                .map_err(Self::Error::from_eth_err)?
565                .unwrap_or_default();
566
567            if at.is_pending() {
568                // for pending tag we need to find the highest nonce of txn in the pending state.
569                if let Some(highest_pool_tx) = this
570                    .pool()
571                    .get_highest_consecutive_transaction_by_sender(address, on_chain_account_nonce)
572                {
573                    {
574                        // and the corresponding txcount is nonce + 1 of the highest tx in the pool
575                        // (on chain nonce is increased after tx)
576                        let next_tx_nonce =
577                            highest_pool_tx.nonce().checked_add(1).ok_or_else(|| {
578                                Self::Error::from(EthApiError::InvalidTransaction(
579                                    RpcInvalidTransactionError::NonceMaxValue,
580                                ))
581                            })?;
582
583                        // guard against drifts in the pool
584                        let next_tx_nonce = on_chain_account_nonce.max(next_tx_nonce);
585
586                        let tx_count = on_chain_account_nonce.max(next_tx_nonce);
587                        return Ok(U256::from(tx_count));
588                    }
589                }
590            }
591            Ok(U256::from(on_chain_account_nonce))
592        })
593    }
594
595    /// Returns code of given account, at the given identifier.
596    fn get_code(
597        &self,
598        address: Address,
599        block_id: Option<BlockId>,
600    ) -> impl Future<Output = Result<Bytes, Self::Error>> + Send
601    where
602        Self: SpawnBlocking,
603    {
604        self.spawn_blocking_io_with_state(block_id.unwrap_or_default(), move |_, state| {
605            Ok(state
606                .account_code(&address)
607                .map_err(Self::Error::from_eth_err)?
608                .unwrap_or_default()
609                .original_bytes())
610        })
611    }
612}