Skip to main content

sqd_portal_client/
lib.rs

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        // we need this to iterate
124        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        // we need this to iterate
219        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            // fields: evm::Fields::all(),
394            ..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            // fields: evm::Fields::all(),
487            ..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            // fields: evm::Fields::all(),
586            ..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            // let arrow_data = arrow_data.unwrap();
595            // let tx_hash = arrow_data
596            //     .transactions
597            //     .column_by_name("value")
598            //     .unwrap()
599            //     .as_any()
600            //     .downcast_ref::<arrow::array::Decimal256Array>()
601            //     .unwrap();
602
603            // for hash in tx_hash.iter().flatten() {
604            //     dbg!(hash.to_string());
605            // }
606        }
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        // dbg!(&query);
646
647        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        // dbg!(arrow_data);
666    }
667}