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
27pub struct RethApi<Provider, EvmConfig> {
31 inner: Arc<RethApiInner<Provider, EvmConfig>>,
32}
33
34impl<Provider, EvmConfig> RethApi<Provider, EvmConfig> {
37 pub fn provider(&self) -> &Provider {
39 &self.inner.provider
40 }
41
42 pub fn evm_config(&self) -> &EvmConfig {
44 &self.inner.evm_config
45 }
46
47 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 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 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 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 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 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 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 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 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 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 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
307async 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
337async 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 ¬ification {
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 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 provider: Provider,
440 evm_config: EvmConfig,
442 blocking_task_guard: BlockingTaskGuard,
444 task_spawner: Runtime,
446}