1#![allow(clippy::get_first)]
2
3use std::str::FromStr;
4use std::sync::Arc;
5use std::time::Duration;
6
7use anyhow::{anyhow, Context, Result};
8use reqwest::{header::CONTENT_TYPE, Client as HttpClient, Method, StatusCode, Url};
9use tokio::sync::mpsc;
10
11pub mod evm;
12pub mod svm;
13
14#[derive(Debug, Clone, Copy)]
15pub struct ClientConfig {
16 pub max_num_retries: usize,
17 pub retry_backoff_ms: u64,
18 pub retry_base_ms: u64,
19 pub retry_ceiling_ms: u64,
20 pub http_req_timeout_millis: u64,
21}
22
23impl Default for ClientConfig {
24 fn default() -> Self {
25 Self {
26 max_num_retries: 9,
27 retry_backoff_ms: 1000,
28 retry_base_ms: 250,
29 retry_ceiling_ms: 2000,
30 http_req_timeout_millis: 40_000,
31 }
32 }
33}
34
35#[derive(Debug, Clone, Copy)]
36pub struct StreamConfig {
37 pub stop_on_head: bool,
38 pub head_poll_interval_millis: u64,
39 pub buffer_size: usize,
40}
41
42impl Default for StreamConfig {
43 fn default() -> Self {
44 Self {
45 stop_on_head: false,
46 head_poll_interval_millis: 1_000,
47 buffer_size: 10,
48 }
49 }
50}
51
52pub struct Client {
53 http_client: HttpClient,
54 url: Url,
55 max_num_retries: usize,
56 retry_backoff_ms: u64,
57 retry_base_ms: u64,
58 retry_ceiling_ms: u64,
59}
60
61static APP_USER_AGENT: &str = concat!("sqd-portal-client-rust/", env!("CARGO_PKG_VERSION"),);
62
63impl Client {
64 pub fn new(url: Url, config: ClientConfig) -> Self {
65 let http_client = HttpClient::builder()
66 .user_agent(APP_USER_AGENT)
67 .http1_only()
68 .gzip(true)
69 .timeout(Duration::from_millis(config.http_req_timeout_millis))
70 .build()
71 .unwrap();
72
73 Self {
74 http_client,
75 url,
76 max_num_retries: config.max_num_retries,
77 retry_backoff_ms: config.retry_backoff_ms,
78 retry_base_ms: config.retry_base_ms,
79 retry_ceiling_ms: config.retry_ceiling_ms,
80 }
81 }
82
83 pub async fn svm_arrow_finalized_query(
84 &self,
85 query: &svm::Query,
86 ) -> Result<Option<svm::ArrowResponse>> {
87 let query = simd_json::to_vec(query).context("serliaze query")?;
88 let query = bytes::Bytes::from(query);
89
90 let response = self.finalized_query(query).await.context("execute query")?;
91 let response = match response {
92 Some(r) => r,
93 None => return Ok(None),
94 };
95
96 let mut parser = svm::ArrowResponseParser::default();
97
98 let lines = response.split(|x| *x == b'\n');
99 let mut scratch = Vec::new();
100
101 for line in lines {
102 if line.is_empty() {
103 continue;
104 }
105
106 scratch.extend_from_slice(line);
107 let tape = simd_json::to_tape(&mut scratch).context("json to tape")?;
108 parser.parse_tape(&tape).context("parse tape")?;
109 scratch.clear();
110 }
111
112 Ok(Some(parser.finish()))
113 }
114
115 pub fn svm_arrow_finalized_stream(
116 self: Arc<Self>,
117 query: svm::Query,
118 config: StreamConfig,
119 ) -> mpsc::Receiver<Result<svm::ArrowResponse>> {
120 let (tx, rx) = mpsc::channel(config.buffer_size);
121
122 let mut query = query;
123 query.fields.block.number = true;
125
126 tokio::spawn(async move {
127 loop {
128 if let Some(tb) = query.to_block {
129 if tb < query.from_block {
130 break;
131 }
132 }
133
134 let res = match self
135 .svm_arrow_finalized_query(&query)
136 .await
137 .context("run query")
138 {
139 Ok(r) => r,
140 Err(e) => {
141 tx.send(Err(e)).await.ok();
142 return;
143 }
144 };
145 let res = match res {
146 Some(r) => r,
147 None => {
148 if config.stop_on_head {
149 break;
150 }
151 tokio::time::sleep(Duration::from_millis(config.head_poll_interval_millis))
152 .await;
153 log::debug!("waiting for block {}", query.from_block);
154 continue;
155 }
156 };
157
158 let next_block = match res.next_block().context("get next block from response") {
159 Ok(nb) => nb,
160 Err(e) => {
161 tx.send(Err(e)).await.ok();
162 return;
163 }
164 };
165
166 query.from_block = next_block;
167
168 if tx.send(Ok(res)).await.is_err() {
169 log::debug!("receiver is closed so quitting stream");
170 return;
171 }
172 }
173 });
174
175 rx
176 }
177
178 pub async fn evm_arrow_finalized_query(
179 &self,
180 query: &evm::Query,
181 ) -> Result<Option<evm::ArrowResponse>> {
182 let query = simd_json::to_vec(query).context("serialize query")?;
183 let query = bytes::Bytes::from(query);
184
185 let response = self.finalized_query(query).await.context("execute query")?;
186 let response = match response {
187 Some(r) => r,
188 None => return Ok(None),
189 };
190
191 let mut parser = evm::ArrowResponseParser::default();
192
193 let lines = response.split(|x| *x == b'\n');
194 let mut scratch = Vec::new();
195
196 for line in lines {
197 if line.is_empty() {
198 continue;
199 }
200
201 scratch.extend_from_slice(line);
202 let tape = simd_json::to_tape(&mut scratch).context("json to tape")?;
203 parser.parse_tape(&tape).context("parse tape")?;
204 scratch.clear();
205 }
206
207 Ok(Some(parser.finish()))
208 }
209
210 pub fn evm_arrow_finalized_stream(
211 self: Arc<Self>,
212 query: evm::Query,
213 config: StreamConfig,
214 ) -> mpsc::Receiver<Result<evm::ArrowResponse>> {
215 let (tx, rx) = mpsc::channel(config.buffer_size);
216
217 let mut query = query;
218 query.fields.block.number = true;
220
221 tokio::spawn(async move {
222 loop {
223 if let Some(tb) = query.to_block {
224 if tb < query.from_block {
225 break;
226 }
227 }
228
229 let res = match self
230 .evm_arrow_finalized_query(&query)
231 .await
232 .context("run query")
233 {
234 Ok(r) => r,
235 Err(e) => {
236 tx.send(Err(e)).await.ok();
237 return;
238 }
239 };
240 let res = match res {
241 Some(r) => r,
242 None => {
243 if config.stop_on_head {
244 break;
245 }
246 tokio::time::sleep(Duration::from_millis(config.head_poll_interval_millis))
247 .await;
248 log::debug!("waiting for block {}", query.from_block);
249 continue;
250 }
251 };
252
253 let next_block = match res.next_block().context("get next block from response") {
254 Ok(nb) => nb,
255 Err(e) => {
256 tx.send(Err(e)).await.ok();
257 return;
258 }
259 };
260
261 query.from_block = next_block;
262
263 if tx.send(Ok(res)).await.is_err() {
264 log::debug!("receiver is closed so quitting stream");
265 return;
266 }
267 }
268 });
269
270 rx
271 }
272
273 pub async fn finalized_height(&self) -> Result<u64> {
274 let res = self
275 .finalized_req(Method::GET, &["finalized-stream", "height"], None)
276 .await
277 .context("make req")?
278 .context("no response data")?;
279
280 let height = std::str::from_utf8(&res).context("check body is utf8")?;
281 let height = u64::from_str(height).context("parse height as number")?;
282
283 Ok(height)
284 }
285
286 async fn finalized_query(&self, query: bytes::Bytes) -> Result<Option<bytes::Bytes>> {
287 self.finalized_req(Method::POST, &["finalized-stream"], Some(query))
288 .await
289 }
290
291 async fn finalized_req(
292 &self,
293 method: Method,
294 url_segments: &[&str],
295 body: Option<bytes::Bytes>,
296 ) -> Result<Option<bytes::Bytes>> {
297 let mut base = self.retry_base_ms;
298
299 let mut err = anyhow!("");
300
301 for _ in 0..self.max_num_retries + 1 {
302 match self
303 .finalized_req_impl(method.clone(), url_segments, body.clone())
304 .await
305 {
306 Ok(res) => return Ok(res),
307 Err(e) => {
308 log::error!(
309 "failed to get data from server, retrying... The error was: {:?}",
310 e
311 );
312 err = err.context(format!("{:?}", e));
313 }
314 }
315
316 let base_ms = Duration::from_millis(base);
317 let jitter = Duration::from_millis(rand::random::<u64>() % self.retry_backoff_ms);
318
319 tokio::time::sleep(base_ms + jitter).await;
320
321 base = std::cmp::min(base + self.retry_backoff_ms, self.retry_ceiling_ms);
322 }
323
324 Err(err)
325 }
326
327 async fn finalized_req_impl(
328 &self,
329 method: Method,
330 url_segments: &[&str],
331 body: Option<bytes::Bytes>,
332 ) -> Result<Option<bytes::Bytes>> {
333 let mut url = self.url.clone();
334 let mut segments = url.path_segments_mut().ok().context("get path segments")?;
335 for s in url_segments {
336 segments.push(s);
337 }
338 std::mem::drop(segments);
339 let req = self.http_client.request(method, url);
340
341 let mut req = req.header(CONTENT_TYPE, "application/json");
342
343 if let Some(body) = body {
344 req = req.body(body);
345 }
346
347 let res = req.send().await.context("execute http req")?;
348
349 let status = res.status();
350 if !status.is_success() {
351 let text = res.text().await.context("read text to see error")?;
352
353 return Err(anyhow!(
354 "http response status code {}, err body: {}",
355 status,
356 text
357 ));
358 } else if status == StatusCode::NO_CONTENT {
359 return Ok(None);
360 }
361
362 res.bytes()
363 .await
364 .context("read response body bytes")
365 .map(Some)
366 }
367}
368
369#[cfg(test)]
370mod tests {
371 use super::*;
372
373 #[tokio::test(flavor = "multi_thread")]
374 #[ignore]
375 async fn continuous_stream_evm() {
376 let url = "https://portal.sqd.dev/datasets/ethereum-mainnet"
377 .parse()
378 .unwrap();
379 let client = Client::new(url, ClientConfig::default());
380
381 let height = client.finalized_height().await.unwrap();
382
383 let query = evm::Query {
384 from_block: height,
385 include_all_blocks: true,
386 fields: evm::Fields {
387 block: evm::BlockFields {
388 number: true,
389 ..Default::default()
390 },
391 ..Default::default()
392 },
393 ..Default::default()
395 };
396
397 let client = Arc::new(client);
398
399 let mut receiver = client.evm_arrow_finalized_stream(query, StreamConfig::default());
400
401 while let Some(arrow_data) = receiver.recv().await {
402 let arrow_data = arrow_data.unwrap();
403 let block_num = arrow_data
404 .blocks
405 .column_by_name("number")
406 .unwrap()
407 .as_any()
408 .downcast_ref::<arrow::array::UInt64Array>()
409 .unwrap();
410
411 for num in block_num.iter().flatten() {
412 dbg!(num);
413 }
414 }
415 }
416
417 #[tokio::test(flavor = "multi_thread")]
418 #[ignore]
419 async fn check_stream_finishes_properly_svm() {
420 let url = "https://portal.sqd.dev/datasets/solana-beta"
421 .parse()
422 .unwrap();
423 let client = Client::new(url, ClientConfig::default());
424
425 let query = svm::Query {
426 from_block: 317617480,
427 to_block: Some(317617500),
428 fields: svm::Fields {
429 transaction: svm::TransactionFields {
430 recent_blockhash: false,
431 ..svm::TransactionFields::all()
432 },
433 ..svm::Fields::all()
434 },
435 balances: vec![svm::BalanceRequest::default()],
436 include_all_blocks: true,
437 instructions: vec![svm::InstructionRequest::default()],
438 logs: vec![svm::LogRequest::default()],
439 rewards: vec![svm::RewardRequest::default()],
440 token_balances: vec![svm::TokenBalanceRequest::default()],
441 transactions: vec![svm::TransactionRequest::default()],
442 type_: Default::default(),
443 };
444
445 let client = Arc::new(client);
446
447 let mut receiver = client.svm_arrow_finalized_stream(query, StreamConfig::default());
448
449 while let Some(arrow_data) = receiver.recv().await {
450 let arrow_data = arrow_data.unwrap();
451 let tx_hash = arrow_data
452 .transactions
453 .column_by_name("block_slot")
454 .unwrap()
455 .as_any()
456 .downcast_ref::<arrow::array::UInt64Array>()
457 .unwrap();
458
459 for hash in tx_hash.iter().flatten() {
460 dbg!(hash.to_string());
461 }
462 }
463 }
464
465 #[tokio::test(flavor = "multi_thread")]
466 #[ignore]
467 async fn check_stream_finishes_properly() {
468 let url = "https://portal.sqd.dev/datasets/ethereum-mainnet"
469 .parse()
470 .unwrap();
471 let client = Client::new(url, ClientConfig::default());
472
473 let query = evm::Query {
474 from_block: 18123123,
475 to_block: Some(18123222),
476 logs: vec![evm::LogRequest::default()],
477 transactions: vec![evm::TransactionRequest::default()],
478 include_all_blocks: true,
479 fields: evm::Fields {
480 transaction: evm::TransactionFields {
481 value: true,
482 ..Default::default()
483 },
484 ..Default::default()
485 },
486 ..Default::default()
488 };
489
490 let client = Arc::new(client);
491
492 let mut receiver = client.evm_arrow_finalized_stream(query, StreamConfig::default());
493
494 while let Some(arrow_data) = receiver.recv().await {
495 let arrow_data = arrow_data.unwrap();
496 let tx_hash = arrow_data
497 .transactions
498 .column_by_name("value")
499 .unwrap()
500 .as_any()
501 .downcast_ref::<arrow::array::Decimal256Array>()
502 .unwrap();
503
504 for hash in tx_hash.iter().flatten() {
505 dbg!(hash.to_string());
506 }
507 }
508 }
509
510 #[tokio::test(flavor = "multi_thread")]
511 #[ignore]
512 async fn full_evm_query() {
513 env_logger::try_init().ok();
514
515 let url = "https://portal.sqd.dev/datasets/ethereum-mainnet"
516 .parse()
517 .unwrap();
518 let client = Client::new(
519 url,
520 ClientConfig {
521 max_num_retries: 0,
522 ..Default::default()
523 },
524 );
525
526 let query = evm::Query {
527 from_block: 18123123,
528 to_block: Some(18123200),
529 logs: vec![evm::LogRequest::default()],
530 transactions: vec![evm::TransactionRequest::default()],
531 traces: vec![evm::TraceRequest::default()],
532 include_all_blocks: true,
533 fields: evm::Fields::all(),
534 ..Default::default()
535 };
536
537 let client = Arc::new(client);
538
539 let mut receiver = client.evm_arrow_finalized_stream(query, StreamConfig::default());
540
541 while let Some(arrow_data) = receiver.recv().await {
542 let arrow_data = arrow_data.unwrap();
543
544 let tx_hash = arrow_data
545 .traces
546 .column_by_name("from")
547 .unwrap()
548 .as_any()
549 .downcast_ref::<arrow::array::BinaryArray>()
550 .unwrap();
551
552 for hash in tx_hash.iter().flatten() {
553 dbg!(faster_hex::hex_string(hash));
554 }
555 }
556 }
557
558 #[tokio::test(flavor = "multi_thread")]
559 #[ignore]
560 async fn dummy_stream() {
561 env_logger::try_init().ok();
562
563 let url = "https://portal.sqd.dev/datasets/zksync-mainnet"
564 .parse()
565 .unwrap();
566 let client = Client::new(
567 url,
568 ClientConfig {
569 max_num_retries: 0,
570 ..Default::default()
571 },
572 );
573
574 let query = evm::Query {
575 from_block: 12123123,
576 to_block: None,
577 transactions: vec![evm::TransactionRequest::default()],
578 fields: evm::Fields {
579 transaction: evm::TransactionFields {
580 value: true,
581 ..Default::default()
582 },
583 ..Default::default()
584 },
585 ..Default::default()
587 };
588
589 let client = Arc::new(client);
590
591 let mut receiver = client.evm_arrow_finalized_stream(query, StreamConfig::default());
592
593 while let Some(_arrow_data) = receiver.recv().await {
594 }
607 }
608
609 #[tokio::test(flavor = "multi_thread")]
610 #[ignore]
611 async fn dummy_svm() {
612 env_logger::try_init().ok();
613
614 let url = "https://portal.sqd.dev/datasets/solana-beta"
615 .parse()
616 .unwrap();
617 let client = Client::new(
618 url,
619 ClientConfig {
620 max_num_retries: 0,
621 ..Default::default()
622 },
623 );
624
625 let query = svm::Query {
626 from_block: 317617480,
627 to_block: Some(317617500),
628 fields: svm::Fields {
629 transaction: svm::TransactionFields {
630 recent_blockhash: false,
631 ..svm::TransactionFields::all()
632 },
633 ..svm::Fields::all()
634 },
635 balances: vec![svm::BalanceRequest::default()],
636 include_all_blocks: true,
637 instructions: vec![svm::InstructionRequest::default()],
638 logs: vec![svm::LogRequest::default()],
639 rewards: vec![svm::RewardRequest::default()],
640 token_balances: vec![svm::TokenBalanceRequest::default()],
641 transactions: vec![svm::TransactionRequest::default()],
642 type_: Default::default(),
643 };
644
645 let arrow_data = client
648 .svm_arrow_finalized_query(&query)
649 .await
650 .unwrap()
651 .unwrap();
652
653 let timestamp = arrow_data
654 .blocks
655 .column_by_name("parent_slot")
656 .unwrap()
657 .as_any()
658 .downcast_ref::<arrow::array::UInt64Array>()
659 .unwrap();
660
661 for t in timestamp.iter().flatten() {
662 dbg!(t);
663 }
664
665 }
667}