reth_consensus_debug_client/providers/
etherscan.rs1use 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#[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 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 pub const fn with_interval(mut self, interval: Duration) -> Self {
46 self.interval = interval;
47 self
48 }
49
50 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 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
91impl<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 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}