Skip to main content

reth_consensus_debug_client/providers/
etherscan.rs

1use crate::PayloadProvider;
2use alloy_eips::BlockNumberOrTag;
3use alloy_json_rpc::{Response, ResponsePayload};
4use reqwest::Client;
5use reth_node_api::ExecutionPayload;
6use reth_tracing::tracing::{debug, warn};
7use serde::{de::DeserializeOwned, Serialize};
8use std::{sync::Arc, time::Duration};
9use tokio::{sync::mpsc, time::interval};
10
11/// Block provider that fetches new blocks from Etherscan API.
12#[derive(derive_more::Debug)]
13pub struct EtherscanBlockProvider<RpcBlock, ExecutionData> {
14    http_client: Client,
15    base_url: String,
16    api_key: String,
17    chain_id: u64,
18    interval: Duration,
19    #[debug(skip)]
20    convert: Arc<dyn Fn(RpcBlock) -> ExecutionData + Send + Sync>,
21}
22
23impl<RpcBlock, ExecutionData> EtherscanBlockProvider<RpcBlock, ExecutionData>
24where
25    RpcBlock: Serialize + DeserializeOwned,
26{
27    /// Create a new Etherscan block provider with the given base URL and API key.
28    pub fn new(
29        base_url: String,
30        api_key: String,
31        chain_id: u64,
32        convert: impl Fn(RpcBlock) -> ExecutionData + Send + Sync + 'static,
33    ) -> Self {
34        Self {
35            http_client: Client::new(),
36            base_url,
37            api_key,
38            chain_id,
39            interval: Duration::from_secs(3),
40            convert: Arc::new(convert),
41        }
42    }
43
44    /// Sets the interval at which the provider fetches new blocks.
45    pub const fn with_interval(mut self, interval: Duration) -> Self {
46        self.interval = interval;
47        self
48    }
49
50    /// Load block using Etherscan API. Note: only `BlockNumberOrTag::Latest`,
51    /// `BlockNumberOrTag::Earliest`, `BlockNumberOrTag::Pending`, `BlockNumberOrTag::Number(u64)`
52    /// are supported.
53    pub async fn load_payload(
54        &self,
55        block_number_or_tag: BlockNumberOrTag,
56    ) -> eyre::Result<ExecutionData> {
57        let tag = match block_number_or_tag {
58            BlockNumberOrTag::Number(num) => format!("{num:#x}"),
59            tag => tag.to_string(),
60        };
61
62        let mut req = self.http_client.get(&self.base_url).query(&[
63            ("module", "proxy"),
64            ("action", "eth_getBlockByNumber"),
65            ("tag", &tag),
66            ("boolean", "true"),
67            ("apikey", &self.api_key),
68        ]);
69
70        if !self.base_url.contains("chainid=") {
71            // only append chainid if not part of the base url already
72            req = req.query(&[("chainid", &self.chain_id.to_string())]);
73        }
74
75        let resp = req.send().await?.text().await?;
76
77        debug!(target: "etherscan", %resp, "fetched block from etherscan");
78
79        let resp: Response<RpcBlock> = serde_json::from_str(&resp).inspect_err(|err| {
80            warn!(target: "etherscan", "Failed to parse block response from etherscan: {}", err);
81        })?;
82
83        let payload = resp.payload;
84        match payload {
85            ResponsePayload::Success(block) => Ok((self.convert)(block)),
86            ResponsePayload::Failure(err) => Err(eyre::eyre!("Failed to get block: {err}")),
87        }
88    }
89}
90
91// Implemented manually so cloning doesn't require `RpcBlock: Clone` and `ExecutionData: Clone`,
92// which are only used by the shared conversion function.
93impl<RpcBlock, ExecutionData> Clone for EtherscanBlockProvider<RpcBlock, ExecutionData> {
94    fn clone(&self) -> Self {
95        Self {
96            http_client: self.http_client.clone(),
97            base_url: self.base_url.clone(),
98            api_key: self.api_key.clone(),
99            chain_id: self.chain_id,
100            interval: self.interval,
101            convert: Arc::clone(&self.convert),
102        }
103    }
104}
105
106impl<RpcBlock, ExecutionData> PayloadProvider for EtherscanBlockProvider<RpcBlock, ExecutionData>
107where
108    RpcBlock: Serialize + DeserializeOwned + 'static,
109    ExecutionData: ExecutionPayload,
110{
111    type ExecutionData = ExecutionData;
112
113    async fn subscribe_payloads(&self, tx: mpsc::Sender<Self::ExecutionData>) {
114        let mut last_block_number: Option<u64> = None;
115        let mut interval = interval(self.interval);
116        loop {
117            interval.tick().await;
118            let payload = match self.load_payload(BlockNumberOrTag::Latest).await {
119                Ok(payload) => payload,
120                Err(err) => {
121                    warn!(
122                        target: "consensus::debug-client",
123                        %err,
124                        "Failed to fetch a block from Etherscan",
125                    );
126                    continue
127                }
128            };
129            let block_number = payload.block_number();
130            if Some(block_number) == last_block_number {
131                continue;
132            }
133
134            if tx.send(payload).await.is_err() {
135                // Channel closed.
136                break;
137            }
138
139            last_block_number = Some(block_number);
140        }
141    }
142
143    async fn get_payload(&self, block_number: u64) -> eyre::Result<Self::ExecutionData> {
144        self.load_payload(BlockNumberOrTag::Number(block_number)).await
145    }
146}