Skip to main content

reth_rpc/
reth.rs

1use std::{future::Future, sync::Arc};
2
3use alloy_consensus::BlockHeader;
4use alloy_eips::BlockId;
5use alloy_primitives::{map::AddressMap, U256, U64};
6use async_trait::async_trait;
7use futures::{Stream, StreamExt};
8use jsonrpsee::{core::RpcResult, PendingSubscriptionSink, SubscriptionMessage, SubscriptionSink};
9use reth_chain_state::{
10    CanonStateNotification, CanonStateSubscriptions, ForkChoiceSubscriptions,
11    PersistedBlockSubscriptions,
12};
13use reth_errors::{RethError, RethResult};
14use reth_evm::{execute::Executor, ConfigureEvm};
15use reth_execution_types::{Chain, ExecutionOutcome};
16use reth_primitives_traits::{NodePrimitives, SealedHeader};
17use reth_rpc_api::{RethApiServer, RethJitAction};
18use reth_rpc_eth_types::{EthApiError, EthResult};
19use reth_storage_api::{
20    BlockReader, BlockReaderIdExt, ChangeSetReader, StateProvider, StateProviderFactory,
21    TransactionVariant,
22};
23use reth_tasks::{pool::BlockingTaskGuard, CancelOnDrop, Runtime};
24use serde::Serialize;
25use tokio::sync::oneshot;
26
27/// `reth` API implementation.
28///
29/// This type provides the functionality for handling `reth` prototype RPC requests.
30pub struct RethApi<Provider, EvmConfig> {
31    inner: Arc<RethApiInner<Provider, EvmConfig>>,
32}
33
34// === impl RethApi ===
35
36impl<Provider, EvmConfig> RethApi<Provider, EvmConfig> {
37    /// The provider that can interact with the chain.
38    pub fn provider(&self) -> &Provider {
39        &self.inner.provider
40    }
41
42    /// The evm config.
43    pub fn evm_config(&self) -> &EvmConfig {
44        &self.inner.evm_config
45    }
46
47    /// Create a new instance of the [`RethApi`]
48    pub fn new(
49        provider: Provider,
50        evm_config: EvmConfig,
51        blocking_task_guard: BlockingTaskGuard,
52        task_spawner: Runtime,
53    ) -> Self {
54        let inner =
55            Arc::new(RethApiInner { provider, evm_config, blocking_task_guard, task_spawner });
56        Self { inner }
57    }
58}
59
60impl<Provider, EvmConfig> RethApi<Provider, EvmConfig>
61where
62    Provider: BlockReaderIdExt + ChangeSetReader + StateProviderFactory + 'static,
63    EvmConfig: Send + Sync + 'static,
64{
65    /// Executes the future on a new blocking task.
66    async fn on_blocking_task<C, F, R>(&self, c: C) -> EthResult<R>
67    where
68        C: FnOnce(Self) -> F,
69        F: Future<Output = EthResult<R>> + Send + 'static,
70        R: Send + 'static,
71    {
72        let (tx, rx) = oneshot::channel();
73        let this = self.clone();
74        let f = c(this);
75        self.inner.task_spawner.spawn_blocking_task(async move {
76            let res = f.await;
77            let _ = tx.send(res);
78        });
79        rx.await.map_err(|_| EthApiError::InternalEthError)?
80    }
81
82    /// Returns a map of addresses to changed account balanced for a particular block.
83    pub async fn balance_changes_in_block(&self, block_id: BlockId) -> EthResult<AddressMap<U256>> {
84        self.on_blocking_task(async move |this| this.try_balance_changes_in_block(block_id)).await
85    }
86
87    fn try_balance_changes_in_block(&self, block_id: BlockId) -> EthResult<AddressMap<U256>> {
88        let Some(block_number) = self.provider().block_number_for_id(block_id)? else {
89            return Err(EthApiError::HeaderNotFound(block_id))
90        };
91
92        let state = self.provider().state_by_block_id(block_id)?;
93        let accounts_before = self.provider().account_block_changeset(block_number)?;
94        let hash_map = accounts_before.iter().try_fold(
95            AddressMap::default(),
96            |mut hash_map, account_before| -> RethResult<_> {
97                let current_balance = state.account_balance(&account_before.address)?;
98                let prev_balance = account_before.info.as_ref().map(|info| info.balance);
99                if current_balance != prev_balance {
100                    hash_map.insert(account_before.address, current_balance.unwrap_or_default());
101                }
102                Ok(hash_map)
103            },
104        )?;
105        Ok(hash_map)
106    }
107}
108
109impl<N, Provider, EvmConfig> RethApi<Provider, EvmConfig>
110where
111    N: NodePrimitives,
112    Provider: BlockReaderIdExt
113        + ChangeSetReader
114        + StateProviderFactory
115        + BlockReader<Block = N::Block>
116        + CanonStateSubscriptions<Primitives = N>
117        + 'static,
118    EvmConfig: ConfigureEvm<Primitives = N> + 'static,
119{
120    /// Re-executes one or more consecutive blocks and returns the execution outcome.
121    pub async fn block_execution_outcome(
122        &self,
123        block_id: BlockId,
124        count: Option<U64>,
125    ) -> EthResult<Option<ExecutionOutcome<N::Receipt>>> {
126        const MAX_BLOCK_COUNT: u64 = 128;
127
128        let block_count = count.map(|c| c.to::<u64>()).unwrap_or(1);
129        if block_count == 0 || block_count > MAX_BLOCK_COUNT {
130            return Err(EthApiError::InvalidParams(format!(
131                "block count must be between 1 and {MAX_BLOCK_COUNT}, got {block_count}"
132            )))
133        }
134
135        let permit = self
136            .inner
137            .blocking_task_guard
138            .clone()
139            .acquire_owned()
140            .await
141            .map_err(|_| EthApiError::InternalEthError)?;
142        let guard = CancelOnDrop::default();
143        let cancel = guard.clone();
144        let outcome = self
145            .on_blocking_task(async move |this| {
146                let _permit = permit;
147                this.try_block_execution_outcome(block_id, block_count, &cancel)
148            })
149            .await;
150        drop(guard);
151        outcome
152    }
153
154    fn try_block_execution_outcome(
155        &self,
156        block_id: BlockId,
157        block_count: u64,
158        cancel: &CancelOnDrop,
159    ) -> EthResult<Option<ExecutionOutcome<N::Receipt>>> {
160        let Some(start_block) = self.provider().block_number_for_id(block_id)? else {
161            return Ok(None)
162        };
163
164        if start_block == 0 {
165            return Ok(Some(ExecutionOutcome::default()))
166        }
167
168        let state_provider = self.provider().history_by_block_number(start_block - 1)?;
169        let db = reth_revm::database::StateProviderDatabase::new(
170            (&state_provider).into_evm_state_provider(),
171        );
172
173        let mut blocks = Vec::with_capacity(block_count as usize);
174        for block_number in start_block..start_block + block_count {
175            let Some(block) = self
176                .provider()
177                .recovered_block(block_number.into(), TransactionVariant::WithHash)?
178            else {
179                if block_number == start_block {
180                    return Ok(None)
181                }
182                break;
183            };
184            blocks.push(block);
185        }
186
187        // stop between blocks once the request is dropped
188        let blocks = blocks.iter().take_while(|_| !cancel.is_cancelled());
189        let outcome = self.evm_config().executor(db).execute_batch(blocks).map_err(
190            |e: reth_evm::execute::BlockExecutionError| {
191                EthApiError::Internal(reth_errors::RethError::Other(e.into()))
192            },
193        )?;
194        if cancel.is_cancelled() {
195            return Err(EthApiError::InternalEthError)
196        }
197
198        Ok(Some(outcome))
199    }
200}
201
202#[async_trait]
203impl<Provider, EvmConfig> RethApiServer for RethApi<Provider, EvmConfig>
204where
205    Provider: BlockReaderIdExt
206        + ChangeSetReader
207        + StateProviderFactory
208        + BlockReader<
209            Block = <<Provider as CanonStateSubscriptions>::Primitives as NodePrimitives>::Block,
210        > + CanonStateSubscriptions
211        + ForkChoiceSubscriptions<
212            Header = <<Provider as CanonStateSubscriptions>::Primitives as NodePrimitives>::BlockHeader,
213        >
214        + PersistedBlockSubscriptions
215        + 'static,
216    EvmConfig: ConfigureEvm<Primitives = <Provider as CanonStateSubscriptions>::Primitives> + 'static,
217{
218    /// Handler for `reth_getBalanceChangesInBlock`
219    async fn reth_get_balance_changes_in_block(
220        &self,
221        block_id: BlockId,
222    ) -> RpcResult<AddressMap<U256>> {
223        Ok(Self::balance_changes_in_block(self, block_id).await?)
224    }
225
226    /// Handler for `reth_getBlockExecutionOutcome`
227    async fn reth_get_block_execution_outcome(
228        &self,
229        block_id: BlockId,
230        count: Option<U64>,
231    ) -> RpcResult<Option<serde_json::Value>> {
232        let outcome = Self::block_execution_outcome(self, block_id, count).await?;
233        match outcome {
234            Some(outcome) => {
235                let value = serde_json::to_value(&outcome).map_err(|e| {
236                    EthApiError::Internal(reth_errors::RethError::msg(e.to_string()))
237                })?;
238                Ok(Some(value))
239            }
240            None => Ok(None),
241        }
242    }
243
244    /// Handler for `reth_jit`
245    async fn reth_jit(&self, action: RethJitAction) -> RpcResult<()> {
246        let Some(jit_backend) = self.evm_config().jit_backend() else {
247            return Ok(());
248        };
249
250        match action {
251            RethJitAction::Enable => jit_backend
252                .set_enabled(true)
253                .map_err(|err| EthApiError::Internal(RethError::msg(err)))?,
254            RethJitAction::Disable => jit_backend
255                .set_enabled(false)
256                .map_err(|err| EthApiError::Internal(RethError::msg(err)))?,
257            RethJitAction::Pause => jit_backend.pause(),
258            RethJitAction::Unpause => jit_backend.resume(),
259            RethJitAction::Clear => jit_backend.clear(),
260        }
261
262        Ok(())
263    }
264
265    /// Handler for `reth_subscribeChainNotifications`
266    async fn reth_subscribe_chain_notifications(
267        &self,
268        pending: PendingSubscriptionSink,
269    ) -> jsonrpsee::core::SubscriptionResult {
270        let sink = pending.accept().await?;
271        let stream = self.provider().canonical_state_stream();
272        self.inner.task_spawner.spawn_task(pipe_from_stream(sink, stream));
273
274        Ok(())
275    }
276
277    /// Handler for `reth_subscribePersistedBlock`
278    async fn reth_subscribe_persisted_block(
279        &self,
280        pending: PendingSubscriptionSink,
281    ) -> jsonrpsee::core::SubscriptionResult {
282        let sink = pending.accept().await?;
283        let stream = self.provider().persisted_block_stream();
284        self.inner.task_spawner.spawn_task(pipe_from_stream(sink, stream));
285
286        Ok(())
287    }
288
289    /// Handler for `reth_subscribeFinalizedChainNotifications`
290    async fn reth_subscribe_finalized_chain_notifications(
291        &self,
292        pending: PendingSubscriptionSink,
293    ) -> jsonrpsee::core::SubscriptionResult {
294        let sink = pending.accept().await?;
295        let canon_stream = self.provider().canonical_state_stream();
296        let finalized_stream = self.provider().finalized_block_stream();
297        self.inner.task_spawner.spawn_task(finalized_chain_notifications(
298            sink,
299            canon_stream,
300            finalized_stream,
301        ));
302
303        Ok(())
304    }
305}
306
307/// Pipes all stream items to the subscription sink.
308async fn pipe_from_stream<S, T>(sink: SubscriptionSink, mut stream: S)
309where
310    S: Stream<Item = T> + Unpin,
311    T: Serialize,
312{
313    loop {
314        tokio::select! {
315            _ = sink.closed() => {
316                break
317            }
318            maybe_item = stream.next() => {
319                let Some(item) = maybe_item else {
320                    break
321                };
322                let msg = match SubscriptionMessage::new(sink.method_name(), sink.subscription_id(), &item) {
323                    Ok(msg) => msg,
324                    Err(err) => {
325                        tracing::error!(target: "rpc::reth", %err, "Failed to serialize subscription message");
326                        break
327                    }
328                };
329                if sink.send(msg).await.is_err() {
330                    break;
331                }
332            }
333        }
334    }
335}
336
337/// Buffers committed chain notifications and emits them when a new finalized block is received.
338async fn finalized_chain_notifications<N>(
339    sink: SubscriptionSink,
340    mut canon_stream: reth_chain_state::CanonStateNotificationStream<N>,
341    mut finalized_stream: reth_chain_state::ForkChoiceStream<SealedHeader<N::BlockHeader>>,
342) where
343    N: NodePrimitives,
344{
345    let mut buffered: Vec<CanonStateNotification<N>> = Vec::new();
346
347    loop {
348        tokio::select! {
349            _ = sink.closed() => {
350                break
351            }
352            maybe_canon = canon_stream.next() => {
353                let Some(notification) = maybe_canon else { break };
354                match &notification {
355                    CanonStateNotification::Commit { .. } => {
356                        buffered.push(notification);
357                    }
358                    CanonStateNotification::Reorg { old, new } => {
359                        let first_reverted = old.first().number();
360                        buffered.retain_mut(|notification| {
361                            let chain = notification.committed();
362                            if chain.first().number() >= first_reverted {
363                                return false
364                            }
365                            // Preserve the canonical prefix of a segment crossing the fork.
366                            if chain.tip().number() >= first_reverted {
367                                let (blocks, mut outcome, mut trie_data) = (*chain).clone().into_inner();
368                                outcome.revert_to(first_reverted - 1);
369                                trie_data.split_off(&first_reverted);
370                                *notification = CanonStateNotification::Commit {
371                                    new: Arc::new(Chain::new(
372                                        blocks.into_blocks().take_while(|b| b.number() < first_reverted),
373                                        outcome,
374                                        trie_data,
375                                    )),
376                                };
377                            }
378                            true
379                        });
380                        if !new.is_empty() {
381                            buffered.push(CanonStateNotification::Commit { new: new.clone() });
382                        }
383                    }
384                }
385            }
386            maybe_finalized = finalized_stream.next() => {
387                let Some(finalized_header) = maybe_finalized else { break };
388                let finalized_num = finalized_header.number();
389
390                let mut committed = Vec::new();
391                buffered.retain(|n| {
392                    if *n.committed().range().end() <= finalized_num {
393                        committed.push(n.clone());
394                        false
395                    } else {
396                        true
397                    }
398                });
399
400                if committed.is_empty() {
401                    continue;
402                }
403
404                committed.sort_by_key(|n| *n.committed().range().start());
405
406                let msg = match SubscriptionMessage::new(
407                    sink.method_name(),
408                    sink.subscription_id(),
409                    &committed,
410                ) {
411                    Ok(msg) => msg,
412                    Err(err) => {
413                        tracing::error!(target: "rpc::reth", %err, "Failed to serialize finalized chain notification");
414                        break
415                    }
416                };
417                if sink.send(msg).await.is_err() {
418                    break;
419                }
420            }
421        }
422    }
423}
424
425impl<Provider, EvmConfig> std::fmt::Debug for RethApi<Provider, EvmConfig> {
426    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
427        f.debug_struct("RethApi").finish_non_exhaustive()
428    }
429}
430
431impl<Provider, EvmConfig> Clone for RethApi<Provider, EvmConfig> {
432    fn clone(&self) -> Self {
433        Self { inner: Arc::clone(&self.inner) }
434    }
435}
436
437struct RethApiInner<Provider, EvmConfig> {
438    /// The provider that can interact with the chain.
439    provider: Provider,
440    /// The EVM configuration used to create block executors.
441    evm_config: EvmConfig,
442    /// Guard to restrict the number of concurrent block re-execution requests.
443    blocking_task_guard: BlockingTaskGuard,
444    /// The type that can spawn tasks which would otherwise block.
445    task_spawner: Runtime,
446}