Skip to content

Commit efc5491

Browse files
cleanup jrpc interface (#45)
1 parent 23a4eb7 commit efc5491

4 files changed

Lines changed: 177 additions & 227 deletions

File tree

Cargo.lock

Lines changed: 3 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

be/src/sync.rs

Lines changed: 12 additions & 87 deletions
Original file line numberDiff line numberDiff line change
@@ -92,23 +92,12 @@ impl RemoteConfig {
9292
pub async fn test(url: &str, chain: u64) -> Result<(), shared::Error> {
9393
let parsed: Url = url.parse().wrap_err("unable to parse rpc url")?;
9494
let jrpc_client = jrpc::Client::new(parsed.as_str());
95-
let resp = jrpc_client
96-
.send_one(serde_json::json!({
97-
"id": "1",
98-
"jsonrpc": "2.0",
99-
"method": "eth_chainId",
100-
"params": [],
101-
}))
102-
.await;
103-
match resp {
95+
match jrpc_client.chain_id().await {
10496
Err(e) => Err(shared::Error::User(format!("rpc error {e}"))),
105-
Ok(resp) => match resp.to::<U64>() {
106-
Ok(id) if id.to::<u64>() == chain => Ok(()),
107-
Ok(id) => Err(shared::Error::User(format!(
108-
"expected chain {chain} got {id}",
109-
))),
110-
Err(e) => Err(shared::Error::User(format!("rpc error {e}"))),
111-
},
97+
Ok(resp) if resp.to::<u64>() == chain => Ok(()),
98+
Ok(id) => Err(shared::Error::User(format!(
99+
"expected chain {chain} got {id}",
100+
))),
112101
}
113102
}
114103

@@ -206,8 +195,8 @@ impl Downloader {
206195
return Ok(());
207196
}
208197
let block = match self.start_block {
209-
Some(n) => remote_block(&self.jrpc_client, U64::from(n)).await?,
210-
None => remote_block_latest(&self.jrpc_client).await?,
198+
Some(n) => self.jrpc_client.block(U64::from(n).to_string()).await?,
199+
None => self.jrpc_client.block("latest".to_string()).await?,
211200
};
212201
tracing::info!("initializing blocks table at: {}", block.number);
213202
let mut pg = self.be_pool.get().await.wrap_err("getting pg")?;
@@ -296,7 +285,7 @@ impl Downloader {
296285
*/
297286
#[tracing::instrument(level="info" skip_all, fields(from, to, blocks, txs, logs))]
298287
async fn download(&mut self, batch_size: u16) -> Result<u64, Error> {
299-
let latest = remote_block_latest(&self.jrpc_client).await?;
288+
let latest = self.jrpc_client.block("latest".to_string()).await?;
300289
let _ = self.broadcaster.json_updates.send(serde_json::json!({
301290
"new_block": "remote",
302291
"chain": self.chain.0,
@@ -320,8 +309,8 @@ impl Downloader {
320309
.record("from", from)
321310
.record("to", to);
322311
let (mut blocks, mut logs) = (
323-
download_blocks(&self.jrpc_client, from, to).await?,
324-
download_logs(&self.jrpc_client, from, to).await?,
312+
self.jrpc_client.blocks(from, to).await?,
313+
self.jrpc_client.logs(from, to).await?,
325314
);
326315
add_timestamp(&mut blocks, &mut logs);
327316
validate_blocks(from, to, &blocks)?;
@@ -365,8 +354,8 @@ pub async fn sync_one(
365354
chain: u64,
366355
n: u64,
367356
) -> Result<u64, Error> {
368-
let mut blocks = download_blocks(client, n, n).await?;
369-
let mut logs = download_logs(client, n, n).await?;
357+
let mut blocks = client.blocks(n, n).await?;
358+
let mut logs = client.logs(n, n).await?;
370359
add_timestamp(&mut blocks, &mut logs);
371360
validate_blocks(n, n, &blocks)?;
372361

@@ -376,70 +365,6 @@ pub async fn sync_one(
376365
Ok(num_logs)
377366
}
378367

379-
async fn remote_block(client: &jrpc::Client, n: U64) -> Result<jrpc::Block, Error> {
380-
Ok(client
381-
.send_one(serde_json::json!({
382-
"id": "1",
383-
"jsonrpc": "2.0",
384-
"method": "eth_getBlockByNumber",
385-
"params": [n, true],
386-
}))
387-
.await?
388-
.to()?)
389-
}
390-
391-
async fn remote_block_latest(client: &jrpc::Client) -> Result<jrpc::Block, Error> {
392-
Ok(client
393-
.send_one(serde_json::json!({
394-
"id": "1",
395-
"jsonrpc": "2.0",
396-
"method": "eth_getBlockByNumber",
397-
"params": ["latest", true],
398-
}))
399-
.await?
400-
.to()?)
401-
}
402-
#[tracing::instrument(level="info" skip_all, fields(from, to))]
403-
async fn download_blocks(
404-
client: &jrpc::Client,
405-
from: u64,
406-
to: u64,
407-
) -> Result<Vec<jrpc::Block>, Error> {
408-
Ok(client
409-
.send(
410-
(from..=to)
411-
.map(|n| {
412-
serde_json::json!({
413-
"id": "1",
414-
"jsonrpc": "2.0",
415-
"method": "eth_getBlockByNumber",
416-
"params": [U64::from(n), true],
417-
})
418-
})
419-
.collect(),
420-
)
421-
.await?
422-
.into_iter()
423-
.map(|resp| resp.to::<jrpc::Block>())
424-
.collect::<Result<Vec<_>, _>>()?
425-
.into_iter()
426-
.sorted_by(|a, b| a.number.cmp(&b.number))
427-
.collect())
428-
}
429-
430-
#[tracing::instrument(level="info" skip_all, fields(from, to))]
431-
async fn download_logs(client: &jrpc::Client, from: u64, to: u64) -> Result<Vec<jrpc::Log>, Error> {
432-
Ok(client
433-
.send_one(serde_json::json!({
434-
"id": "1",
435-
"jsonrpc": "2.0",
436-
"method": "eth_getLogs",
437-
"params": [{"fromBlock": U64::from(from), "toBlock": U64::from(to)}],
438-
}))
439-
.await?
440-
.to()?)
441-
}
442-
443368
fn validate_logs(blocks: &[jrpc::Block], logs: &[jrpc::Log]) -> Result<(), Error> {
444369
let mut logs_by_block: HashMap<U64, Vec<&jrpc::Log>> = HashMap::new();
445370
for log in logs {

shared/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ edition = "2021"
77
test = ["dep:target-triple", "dep:flate2", "dep:tar", "dep:rand"]
88

99
[dependencies]
10+
itertools = "0.13.0"
1011
alloy = { version = "0.8.3", features = ["rpc-types-eth"] }
1112
axum = { version = "0.7.5" }
1213
tokio = { version = "1", features = ["macros", "fs", "process"] }

0 commit comments

Comments
 (0)