Skip to main content

quicknode_sdk/streams/
mod.rs

1pub mod stream;
2
3pub use stream::{
4    AddressBookConfig, AzureAttributes, CreateStreamParams, DestinationAttributes,
5    EnabledCountResponse, FilterLanguage, KafkaAttributes, ListStreamsParams, ListStreamsResponse,
6    PageInfo, PostgresAttributes, ProductType, S3Attributes, Stream, StreamDataset,
7    StreamDestination, StreamMetadataLocation, StreamRegion, StreamStatus, TestFilterParams,
8    TestFilterResponse, UpdateStreamParams, WebhookAttributes,
9};
10
11use crate::{config::StreamsConfig, errors::SdkError, SdkConfig};
12
13const STREAMS_BASE_URL: &str = "https://api.quicknode.com/streams/rest/v1/";
14
15pub(crate) struct ResolvedStreamsConfig {
16    pub(crate) base_url: reqwest::Url,
17}
18
19impl ResolvedStreamsConfig {
20    pub(crate) fn from_config(config: Option<&StreamsConfig>) -> Result<Self, SdkError> {
21        let url_str = config
22            .and_then(|s| s.base_url.as_deref())
23            .unwrap_or(STREAMS_BASE_URL);
24        let mut base_url = reqwest::Url::parse(url_str)?;
25        if !base_url.path().ends_with('/') {
26            base_url.set_path(&format!("{}/", base_url.path()));
27        }
28        Ok(Self { base_url })
29    }
30}
31
32// ── Client ─────────────────────────────────────────────────────────────────
33
34/// Client for the Quicknode Streams REST API. Create, list, update, and control
35/// blockchain data streams that deliver filtered on-chain events to configured
36/// destinations.
37#[derive(Debug, Clone)]
38pub struct StreamsApiClient {
39    config: SdkConfig,
40}
41
42impl StreamsApiClient {
43    pub fn new(config: SdkConfig) -> Self {
44        Self { config }
45    }
46
47    /// Creates a new Stream on a given blockchain network and dataset, delivering
48    /// batches to the configured destination. Start from a specific block for
49    /// backfills or from the tip for real-time streaming, and optionally attach
50    /// a base64-encoded JavaScript filter to transform data before delivery.
51    /// The stream can be created in an active or paused state and supports
52    /// reorg handling, distance-from-tip, elastic batching, notification emails,
53    /// and extra destinations for multi-destination delivery.
54    pub async fn create_stream(&self, params: &CreateStreamParams) -> Result<Stream, SdkError> {
55        let url = self.config.streams().base_url.join("streams")?;
56        let resp = self
57            .config
58            .http_client()
59            .post(url)
60            .json(params)
61            .send()
62            .await
63            .map_err(SdkError::Http)?;
64
65        let status = resp.status();
66        let body = resp.text().await.map_err(SdkError::Http)?;
67
68        if !status.is_success() {
69            return Err(SdkError::Api { status, body });
70        }
71        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
72    }
73
74    /// Returns a paginated list of streams on the account. Each stream includes
75    /// its full configuration — identifiers, timestamps, network and dataset,
76    /// filter, block range, destination settings, and operational status — and
77    /// surfaces advanced features such as elastic batching and extra
78    /// destinations, where batches must be delivered to every configured
79    /// destination before the stream advances. Supports pagination via
80    /// `offset`/`limit` and sorting via `order_by`/`order_direction`, and can
81    /// filter by stream type.
82    pub async fn list_streams(
83        &self,
84        params: &ListStreamsParams,
85    ) -> Result<ListStreamsResponse, SdkError> {
86        let mut url = self.config.streams().base_url.join("streams")?;
87        {
88            let mut pairs = url.query_pairs_mut();
89            if let Some(t) = &params.stream_type {
90                pairs.append_pair("type", t);
91            }
92            if let Some(v) = params.offset {
93                pairs.append_pair("offset", &v.to_string());
94            }
95            if let Some(v) = params.limit {
96                pairs.append_pair("limit", &v.to_string());
97            }
98            if let Some(v) = &params.order_by {
99                pairs.append_pair("order_by", v);
100            }
101            if let Some(v) = &params.order_direction {
102                pairs.append_pair("order_direction", v);
103            }
104        }
105        let resp = self
106            .config
107            .http_client()
108            .get(url)
109            .send()
110            .await
111            .map_err(SdkError::Http)?;
112        let status = resp.status();
113        let body = resp.text().await.map_err(SdkError::Http)?;
114        if !status.is_success() {
115            return Err(SdkError::Api { status, body });
116        }
117        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
118    }
119
120    /// Removes every stream on the account. Takes no filters and cannot be
121    /// undone.
122    pub async fn delete_all_streams(&self) -> Result<(), SdkError> {
123        let url = self.config.streams().base_url.join("streams")?;
124        let resp = self
125            .config
126            .http_client()
127            .delete(url)
128            .send()
129            .await
130            .map_err(SdkError::Http)?;
131        let status = resp.status();
132        if !status.is_success() {
133            let body = resp.text().await.map_err(SdkError::Http)?;
134            return Err(SdkError::Api { status, body });
135        }
136        Ok(())
137    }
138
139    /// Returns a single stream by ID, including its full configuration and
140    /// current status.
141    pub async fn get_stream(&self, id: &str) -> Result<Stream, SdkError> {
142        let url = self
143            .config
144            .streams()
145            .base_url
146            .join(&format!("streams/{id}"))?;
147        let resp = self
148            .config
149            .http_client()
150            .get(url)
151            .send()
152            .await
153            .map_err(SdkError::Http)?;
154        let status = resp.status();
155        let body = resp.text().await.map_err(SdkError::Http)?;
156        if !status.is_success() {
157            return Err(SdkError::Api { status, body });
158        }
159        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
160    }
161
162    /// Updates an existing stream's configuration. Only fields present on
163    /// `params` are modified; omitted fields are left unchanged.
164    pub async fn update_stream(
165        &self,
166        id: &str,
167        params: &UpdateStreamParams,
168    ) -> Result<Stream, SdkError> {
169        let url = self
170            .config
171            .streams()
172            .base_url
173            .join(&format!("streams/{id}"))?;
174        let resp = self
175            .config
176            .http_client()
177            .patch(url)
178            .json(params)
179            .send()
180            .await
181            .map_err(SdkError::Http)?;
182        let status = resp.status();
183        let body = resp.text().await.map_err(SdkError::Http)?;
184        if !status.is_success() {
185            return Err(SdkError::Api { status, body });
186        }
187        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
188    }
189
190    /// Deletes a single stream by ID.
191    pub async fn delete_stream(&self, id: &str) -> Result<(), SdkError> {
192        let url = self
193            .config
194            .streams()
195            .base_url
196            .join(&format!("streams/{id}"))?;
197        let resp = self
198            .config
199            .http_client()
200            .delete(url)
201            .send()
202            .await
203            .map_err(SdkError::Http)?;
204        let status = resp.status();
205        if !status.is_success() {
206            let body = resp.text().await.map_err(SdkError::Http)?;
207            return Err(SdkError::Api { status, body });
208        }
209        Ok(())
210    }
211
212    /// Activates a stream by ID, resuming delivery from its current position.
213    pub async fn activate_stream(&self, id: &str) -> Result<(), SdkError> {
214        let url = self
215            .config
216            .streams()
217            .base_url
218            .join(&format!("streams/{id}/activate"))?;
219        let resp = self
220            .config
221            .http_client()
222            .post(url)
223            .send()
224            .await
225            .map_err(SdkError::Http)?;
226        let status = resp.status();
227        if !status.is_success() {
228            let body = resp.text().await.map_err(SdkError::Http)?;
229            return Err(SdkError::Api { status, body });
230        }
231        Ok(())
232    }
233
234    /// Pauses a stream by ID, halting delivery until it is activated again.
235    pub async fn pause_stream(&self, id: &str) -> Result<(), SdkError> {
236        let url = self
237            .config
238            .streams()
239            .base_url
240            .join(&format!("streams/{id}/pause"))?;
241        let resp = self
242            .config
243            .http_client()
244            .post(url)
245            .send()
246            .await
247            .map_err(SdkError::Http)?;
248        let status = resp.status();
249        if !status.is_success() {
250            let body = resp.text().await.map_err(SdkError::Http)?;
251            return Err(SdkError::Api { status, body });
252        }
253        Ok(())
254    }
255
256    /// Runs a filter function against a specified block on a given network and
257    /// dataset, returning the filter's output so it can be validated before
258    /// being attached to a live stream.
259    pub async fn test_filter(
260        &self,
261        params: &TestFilterParams,
262    ) -> Result<TestFilterResponse, SdkError> {
263        let url = self.config.streams().base_url.join("streams/test_filter")?;
264        let resp = self
265            .config
266            .http_client()
267            .post(url)
268            .json(params)
269            .send()
270            .await
271            .map_err(SdkError::Http)?;
272        let status = resp.status();
273        let body = resp.text().await.map_err(SdkError::Http)?;
274        if !status.is_success() {
275            return Err(SdkError::Api { status, body });
276        }
277        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
278    }
279
280    /// Returns the total count of currently enabled (active) streams on the
281    /// account, optionally filtered by stream type.
282    pub async fn get_enabled_count(
283        &self,
284        stream_type: Option<&str>,
285    ) -> Result<EnabledCountResponse, SdkError> {
286        let mut url = self
287            .config
288            .streams()
289            .base_url
290            .join("streams/enabled_count")?;
291        if let Some(t) = stream_type {
292            url.query_pairs_mut().append_pair("type", t);
293        }
294        let resp = self
295            .config
296            .http_client()
297            .get(url)
298            .send()
299            .await
300            .map_err(SdkError::Http)?;
301        let status = resp.status();
302        let body = resp.text().await.map_err(SdkError::Http)?;
303        if !status.is_success() {
304            return Err(SdkError::Api { status, body });
305        }
306        serde_json::from_str(&body).map_err(|source| SdkError::Decode { source, body })
307    }
308}
309
310// ── Tests ──────────────────────────────────────────────────────────────────
311
312#[cfg(test)]
313#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
314mod tests {
315    use super::*;
316    use crate::{QuicknodeSdk, SdkFullConfig, StreamsConfig};
317    use wiremock::matchers::{method, path};
318    use wiremock::{Mock, MockServer, ResponseTemplate};
319
320    fn make_sdk(base_url: String) -> QuicknodeSdk {
321        QuicknodeSdk::new(&SdkFullConfig {
322            api_key: "test-key".to_string(),
323            http: None,
324            admin: None,
325            streams: Some(StreamsConfig {
326                base_url: Some(base_url),
327            }),
328            webhooks: None,
329            kvstore: None,
330            sql: None,
331            rpc: None,
332        })
333        .unwrap()
334    }
335
336    fn webhook_params() -> CreateStreamParams {
337        CreateStreamParams {
338            name: "test-stream".to_string(),
339            region: StreamRegion::UsaEast,
340            network: "ethereum-mainnet".to_string(),
341            dataset: StreamDataset::Block,
342            start_range: 17000000,
343            end_range: -1,
344            destination_attributes: DestinationAttributes::Webhook(WebhookAttributes {
345                url: "https://example.com/webhook".to_string(),
346                max_retry: 3,
347                retry_interval_sec: 1,
348                post_timeout_sec: 10,
349                compression: Some("none".to_string()),
350                security_token: None,
351            }),
352            plan: Some("growth_plan".to_string()),
353            threshold_fetch_buffer: Some(1000),
354            dataset_batch_size: 1,
355            max_batch_size: None,
356            max_buffer_range_size: None,
357            max_buffer_processing_workers: None,
358            keep_distance_from_tip: None,
359            filter_function: None,
360            filter_language: None,
361            address_book_config: None,
362            include_stream_metadata: None,
363            product_type: None,
364            status: None,
365            notification_email: None,
366            charge_min_cap: None,
367            fix_block_reorgs: None,
368            elastic_batch_enabled: true,
369            extra_destinations: None,
370        }
371    }
372
373    fn stream_response_json() -> serde_json::Value {
374        serde_json::json!({
375            "id": "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab",
376            "name": "test-stream",
377            "status": "active",
378            "created_at": "2026-03-19T12:00:00Z",
379            "updated_at": "2026-03-19T12:00:00Z",
380            "sequence": 0,
381            "network": "ethereum-mainnet",
382            "dataset": "block",
383            "region": "usa_east",
384            "destination": "webhook",
385            "destination_attributes": {
386                "url": "https://example.com/webhook",
387                "max_retry": 3,
388                "retry_interval_sec": 1,
389                "post_timeout_sec": 10,
390                "compression": "none"
391            },
392            "start_range": 17000000,
393            "end_range": -1
394        })
395    }
396
397    #[tokio::test]
398    async fn create_stream_success() {
399        let server = MockServer::start().await;
400
401        Mock::given(method("POST"))
402            .and(path("/streams"))
403            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
404            .mount(&server)
405            .await;
406
407        let sdk = make_sdk(format!("{}/", server.uri()));
408        let resp = sdk.streams.create_stream(&webhook_params()).await.unwrap();
409
410        assert_eq!(resp.id, "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab");
411        assert_eq!(resp.name, "test-stream");
412        assert_eq!(resp.status, "active");
413        assert_eq!(resp.network, "ethereum-mainnet");
414        assert_eq!(resp.dataset, "block");
415        // Verify the full destination_attributes round-trip on the response
416        // side. Without this assertion, serde's flatten+Option silently
417        // swallows malformed destination_attributes as None.
418        match resp.destination_attributes {
419            Some(DestinationAttributes::Webhook(attrs)) => {
420                assert_eq!(attrs.url, "https://example.com/webhook");
421                assert_eq!(attrs.max_retry, 3);
422            }
423            other => panic!("expected Webhook destination, got {other:?}"),
424        }
425    }
426
427    #[tokio::test]
428    async fn create_stream_sends_typed_webhook_destination() {
429        use wiremock::matchers::body_partial_json;
430        let server = MockServer::start().await;
431
432        Mock::given(method("POST"))
433            .and(path("/streams"))
434            .and(body_partial_json(serde_json::json!({
435                "destination": "webhook",
436                "destination_attributes": {
437                    "url": "https://example.com/webhook",
438                    "max_retry": 3,
439                    "retry_interval_sec": 1,
440                    "post_timeout_sec": 10,
441                    "compression": "none"
442                }
443            })))
444            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
445            .mount(&server)
446            .await;
447
448        let sdk = make_sdk(format!("{}/", server.uri()));
449        let params = CreateStreamParams {
450            destination_attributes: DestinationAttributes::Webhook(WebhookAttributes {
451                url: "https://example.com/webhook".to_string(),
452                max_retry: 3,
453                retry_interval_sec: 1,
454                post_timeout_sec: 10,
455                compression: Some("none".to_string()),
456                security_token: None,
457            }),
458            ..webhook_params()
459        };
460        sdk.streams.create_stream(&params).await.unwrap();
461    }
462
463    #[tokio::test]
464    async fn create_stream_sends_extra_destinations() {
465        use wiremock::matchers::body_partial_json;
466        let server = MockServer::start().await;
467
468        Mock::given(method("POST"))
469            .and(path("/streams"))
470            .and(body_partial_json(serde_json::json!({
471                "extra_destinations": [
472                    {
473                        "destination": "webhook",
474                        "destination_attributes": {
475                            "url": "https://example.com/extra-hook",
476                            "max_retry": 5,
477                            "retry_interval_sec": 2,
478                            "post_timeout_sec": 15,
479                            "compression": "none"
480                        }
481                    },
482                    {
483                        "destination": "s3",
484                        "destination_attributes": {
485                            "endpoint": "s3.example.com",
486                            "access_key": "AKIA",
487                            "secret_key": "secret",
488                            "bucket": "my-bucket",
489                            "object_prefix": "streams/",
490                            "compression": "gzip",
491                            "file_type": ".json",
492                            "max_retry": 3,
493                            "retry_interval_sec": 1
494                        }
495                    }
496                ]
497            })))
498            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
499            .mount(&server)
500            .await;
501
502        let sdk = make_sdk(format!("{}/", server.uri()));
503        let params = CreateStreamParams {
504            extra_destinations: Some(vec![
505                DestinationAttributes::Webhook(WebhookAttributes {
506                    url: "https://example.com/extra-hook".to_string(),
507                    max_retry: 5,
508                    retry_interval_sec: 2,
509                    post_timeout_sec: 15,
510                    compression: Some("none".to_string()),
511                    security_token: None,
512                }),
513                DestinationAttributes::S3(S3Attributes {
514                    endpoint: "s3.example.com".to_string(),
515                    access_key: "AKIA".to_string(),
516                    secret_key: "secret".to_string(),
517                    bucket: "my-bucket".to_string(),
518                    object_prefix: "streams/".to_string(),
519                    compression: "gzip".to_string(),
520                    file_type: ".json".to_string(),
521                    max_retry: 3,
522                    retry_interval_sec: 1,
523                    use_ssl: None,
524                }),
525            ]),
526            ..webhook_params()
527        };
528        sdk.streams.create_stream(&params).await.unwrap();
529    }
530
531    #[tokio::test]
532    async fn update_stream_sends_extra_destinations() {
533        use wiremock::matchers::body_partial_json;
534        let server = MockServer::start().await;
535
536        Mock::given(method("PATCH"))
537            .and(path("/streams/test-id"))
538            .and(body_partial_json(serde_json::json!({
539                "extra_destinations": [
540                    {
541                        "destination": "webhook",
542                        "destination_attributes": {
543                            "url": "https://example.com/patched-hook",
544                            "max_retry": 1,
545                            "retry_interval_sec": 1,
546                            "post_timeout_sec": 5,
547                            "compression": "none"
548                        }
549                    }
550                ]
551            })))
552            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
553            .mount(&server)
554            .await;
555
556        let sdk = make_sdk(format!("{}/", server.uri()));
557        let params = UpdateStreamParams {
558            extra_destinations: Some(vec![DestinationAttributes::Webhook(WebhookAttributes {
559                url: "https://example.com/patched-hook".to_string(),
560                max_retry: 1,
561                retry_interval_sec: 1,
562                post_timeout_sec: 5,
563                compression: Some("none".to_string()),
564                security_token: None,
565            })]),
566            ..Default::default()
567        };
568        sdk.streams.update_stream("test-id", &params).await.unwrap();
569    }
570
571    #[tokio::test]
572    async fn get_stream_parses_extra_destinations() {
573        let server = MockServer::start().await;
574        let mut body = stream_response_json();
575        body["extra_destinations"] = serde_json::json!([
576            {
577                "destination": "webhook",
578                "destination_attributes": {
579                    "url": "https://example.com/extra",
580                    "max_retry": 2,
581                    "retry_interval_sec": 1,
582                    "post_timeout_sec": 10,
583                    "compression": "none"
584                }
585            }
586        ]);
587        Mock::given(method("GET"))
588            .and(path("/streams/test-id"))
589            .respond_with(ResponseTemplate::new(200).set_body_json(body))
590            .mount(&server)
591            .await;
592
593        let sdk = make_sdk(format!("{}/", server.uri()));
594        let resp = sdk.streams.get_stream("test-id").await.unwrap();
595        let extras = resp.extra_destinations.expect("extra_destinations present");
596        assert_eq!(extras.len(), 1);
597        match &extras[0] {
598            DestinationAttributes::Webhook(w) => {
599                assert_eq!(w.url, "https://example.com/extra");
600                assert_eq!(w.max_retry, 2);
601            }
602            other => panic!("expected Webhook, got {other:?}"),
603        }
604    }
605
606    #[tokio::test]
607    async fn create_stream_api_error() {
608        let server = MockServer::start().await;
609
610        Mock::given(method("POST"))
611            .and(path("/streams"))
612            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
613            .mount(&server)
614            .await;
615
616        let sdk = make_sdk(format!("{}/", server.uri()));
617        let err = sdk
618            .streams
619            .create_stream(&webhook_params())
620            .await
621            .unwrap_err();
622
623        assert!(matches!(err, SdkError::Api { .. }));
624    }
625
626    #[tokio::test]
627    async fn create_stream_server_error() {
628        let server = MockServer::start().await;
629
630        Mock::given(method("POST"))
631            .and(path("/streams"))
632            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
633            .mount(&server)
634            .await;
635
636        let sdk = make_sdk(format!("{}/", server.uri()));
637        let err = sdk
638            .streams
639            .create_stream(&webhook_params())
640            .await
641            .unwrap_err();
642
643        assert!(matches!(err, SdkError::Api { .. }));
644    }
645
646    #[tokio::test]
647    async fn list_streams_success() {
648        let server = MockServer::start().await;
649        let response = serde_json::json!({
650            "data": [stream_response_json()],
651            "pageInfo": { "limit": 100, "offset": 0, "total": 1 }
652        });
653        Mock::given(method("GET"))
654            .and(path("/streams"))
655            .respond_with(ResponseTemplate::new(200).set_body_json(response))
656            .mount(&server)
657            .await;
658        let sdk = make_sdk(format!("{}/", server.uri()));
659        let resp = sdk
660            .streams
661            .list_streams(&ListStreamsParams::default())
662            .await
663            .unwrap();
664        assert_eq!(resp.data.len(), 1);
665        assert_eq!(resp.page_info.total, 1);
666    }
667
668    #[tokio::test]
669    async fn list_streams_api_error() {
670        let server = MockServer::start().await;
671        Mock::given(method("GET"))
672            .and(path("/streams"))
673            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
674            .mount(&server)
675            .await;
676        let sdk = make_sdk(format!("{}/", server.uri()));
677        let err = sdk
678            .streams
679            .list_streams(&ListStreamsParams::default())
680            .await
681            .unwrap_err();
682        assert!(matches!(err, SdkError::Api { .. }));
683    }
684
685    #[tokio::test]
686    async fn list_streams_server_error() {
687        let server = MockServer::start().await;
688        Mock::given(method("GET"))
689            .and(path("/streams"))
690            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
691            .mount(&server)
692            .await;
693        let sdk = make_sdk(format!("{}/", server.uri()));
694        let err = sdk
695            .streams
696            .list_streams(&ListStreamsParams::default())
697            .await
698            .unwrap_err();
699        assert!(matches!(err, SdkError::Api { .. }));
700    }
701
702    #[tokio::test]
703    async fn get_stream_success() {
704        let server = MockServer::start().await;
705        Mock::given(method("GET"))
706            .and(path("/streams/test-id"))
707            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
708            .mount(&server)
709            .await;
710        let sdk = make_sdk(format!("{}/", server.uri()));
711        let resp = sdk.streams.get_stream("test-id").await.unwrap();
712        assert_eq!(resp.id, "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab");
713    }
714
715    #[tokio::test]
716    async fn get_stream_not_found() {
717        let server = MockServer::start().await;
718        Mock::given(method("GET"))
719            .and(path("/streams/test-id"))
720            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
721            .mount(&server)
722            .await;
723        let sdk = make_sdk(format!("{}/", server.uri()));
724        let err = sdk.streams.get_stream("test-id").await.unwrap_err();
725        assert!(matches!(err, SdkError::Api { .. }));
726    }
727
728    #[tokio::test]
729    async fn get_stream_server_error() {
730        let server = MockServer::start().await;
731        Mock::given(method("GET"))
732            .and(path("/streams/test-id"))
733            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
734            .mount(&server)
735            .await;
736        let sdk = make_sdk(format!("{}/", server.uri()));
737        let err = sdk.streams.get_stream("test-id").await.unwrap_err();
738        assert!(matches!(err, SdkError::Api { .. }));
739    }
740
741    #[tokio::test]
742    async fn update_stream_success() {
743        let server = MockServer::start().await;
744        let mut updated = stream_response_json();
745        updated["name"] = serde_json::json!("updated-name");
746        Mock::given(method("PATCH"))
747            .and(path("/streams/test-id"))
748            .respond_with(ResponseTemplate::new(200).set_body_json(updated))
749            .mount(&server)
750            .await;
751        let sdk = make_sdk(format!("{}/", server.uri()));
752        let params = UpdateStreamParams {
753            name: Some("updated-name".to_string()),
754            ..Default::default()
755        };
756        let resp = sdk.streams.update_stream("test-id", &params).await.unwrap();
757        assert_eq!(resp.name, "updated-name");
758    }
759
760    #[tokio::test]
761    async fn update_stream_api_error() {
762        let server = MockServer::start().await;
763        Mock::given(method("PATCH"))
764            .and(path("/streams/test-id"))
765            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
766            .mount(&server)
767            .await;
768        let sdk = make_sdk(format!("{}/", server.uri()));
769        let params = UpdateStreamParams::default();
770        let err = sdk
771            .streams
772            .update_stream("test-id", &params)
773            .await
774            .unwrap_err();
775        assert!(matches!(err, SdkError::Api { .. }));
776    }
777
778    #[tokio::test]
779    async fn update_stream_server_error() {
780        let server = MockServer::start().await;
781        Mock::given(method("PATCH"))
782            .and(path("/streams/test-id"))
783            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
784            .mount(&server)
785            .await;
786        let sdk = make_sdk(format!("{}/", server.uri()));
787        let params = UpdateStreamParams::default();
788        let err = sdk
789            .streams
790            .update_stream("test-id", &params)
791            .await
792            .unwrap_err();
793        assert!(matches!(err, SdkError::Api { .. }));
794    }
795
796    #[tokio::test]
797    async fn delete_stream_success() {
798        let server = MockServer::start().await;
799        Mock::given(method("DELETE"))
800            .and(path("/streams/test-id"))
801            .respond_with(ResponseTemplate::new(200))
802            .mount(&server)
803            .await;
804        let sdk = make_sdk(format!("{}/", server.uri()));
805        sdk.streams.delete_stream("test-id").await.unwrap();
806    }
807
808    #[tokio::test]
809    async fn delete_stream_not_found() {
810        let server = MockServer::start().await;
811        Mock::given(method("DELETE"))
812            .and(path("/streams/test-id"))
813            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
814            .mount(&server)
815            .await;
816        let sdk = make_sdk(format!("{}/", server.uri()));
817        let err = sdk.streams.delete_stream("test-id").await.unwrap_err();
818        assert!(matches!(err, SdkError::Api { .. }));
819    }
820
821    #[tokio::test]
822    async fn delete_stream_server_error() {
823        let server = MockServer::start().await;
824        Mock::given(method("DELETE"))
825            .and(path("/streams/test-id"))
826            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
827            .mount(&server)
828            .await;
829        let sdk = make_sdk(format!("{}/", server.uri()));
830        let err = sdk.streams.delete_stream("test-id").await.unwrap_err();
831        assert!(matches!(err, SdkError::Api { .. }));
832    }
833
834    #[tokio::test]
835    async fn delete_all_streams_success() {
836        let server = MockServer::start().await;
837        Mock::given(method("DELETE"))
838            .and(path("/streams"))
839            .respond_with(ResponseTemplate::new(204))
840            .mount(&server)
841            .await;
842        let sdk = make_sdk(format!("{}/", server.uri()));
843        sdk.streams.delete_all_streams().await.unwrap();
844    }
845
846    #[tokio::test]
847    async fn delete_all_streams_not_found() {
848        let server = MockServer::start().await;
849        Mock::given(method("DELETE"))
850            .and(path("/streams"))
851            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
852            .mount(&server)
853            .await;
854        let sdk = make_sdk(format!("{}/", server.uri()));
855        let err = sdk.streams.delete_all_streams().await.unwrap_err();
856        assert!(matches!(err, SdkError::Api { .. }));
857    }
858
859    #[tokio::test]
860    async fn activate_stream_success() {
861        let server = MockServer::start().await;
862        Mock::given(method("POST"))
863            .and(path("/streams/test-id/activate"))
864            .respond_with(ResponseTemplate::new(201))
865            .mount(&server)
866            .await;
867        let sdk = make_sdk(format!("{}/", server.uri()));
868        sdk.streams.activate_stream("test-id").await.unwrap();
869    }
870
871    #[tokio::test]
872    async fn activate_stream_not_found() {
873        let server = MockServer::start().await;
874        Mock::given(method("POST"))
875            .and(path("/streams/test-id/activate"))
876            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
877            .mount(&server)
878            .await;
879        let sdk = make_sdk(format!("{}/", server.uri()));
880        let err = sdk.streams.activate_stream("test-id").await.unwrap_err();
881        assert!(matches!(err, SdkError::Api { .. }));
882    }
883
884    #[tokio::test]
885    async fn activate_stream_server_error() {
886        let server = MockServer::start().await;
887        Mock::given(method("POST"))
888            .and(path("/streams/test-id/activate"))
889            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
890            .mount(&server)
891            .await;
892        let sdk = make_sdk(format!("{}/", server.uri()));
893        let err = sdk.streams.activate_stream("test-id").await.unwrap_err();
894        assert!(matches!(err, SdkError::Api { .. }));
895    }
896
897    #[tokio::test]
898    async fn pause_stream_success() {
899        let server = MockServer::start().await;
900        Mock::given(method("POST"))
901            .and(path("/streams/test-id/pause"))
902            .respond_with(ResponseTemplate::new(201))
903            .mount(&server)
904            .await;
905        let sdk = make_sdk(format!("{}/", server.uri()));
906        sdk.streams.pause_stream("test-id").await.unwrap();
907    }
908
909    #[tokio::test]
910    async fn pause_stream_not_found() {
911        let server = MockServer::start().await;
912        Mock::given(method("POST"))
913            .and(path("/streams/test-id/pause"))
914            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
915            .mount(&server)
916            .await;
917        let sdk = make_sdk(format!("{}/", server.uri()));
918        let err = sdk.streams.pause_stream("test-id").await.unwrap_err();
919        assert!(matches!(err, SdkError::Api { .. }));
920    }
921
922    #[tokio::test]
923    async fn pause_stream_server_error() {
924        let server = MockServer::start().await;
925        Mock::given(method("POST"))
926            .and(path("/streams/test-id/pause"))
927            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
928            .mount(&server)
929            .await;
930        let sdk = make_sdk(format!("{}/", server.uri()));
931        let err = sdk.streams.pause_stream("test-id").await.unwrap_err();
932        assert!(matches!(err, SdkError::Api { .. }));
933    }
934
935    #[tokio::test]
936    async fn test_filter_success() {
937        let server = MockServer::start().await;
938        let response = serde_json::json!({ "result": {"hash": "0xabc"}, "logs": [] });
939        Mock::given(method("POST"))
940            .and(path("/streams/test_filter"))
941            .respond_with(ResponseTemplate::new(201).set_body_json(response))
942            .mount(&server)
943            .await;
944        let sdk = make_sdk(format!("{}/", server.uri()));
945        let params = TestFilterParams {
946            network: "ethereum-mainnet".to_string(),
947            dataset: StreamDataset::Block,
948            block: "17811625".to_string(),
949            filter_function: "ZnVuY3Rpb24gbWFpbihkYXRhKSB7IHJldHVybiBkYXRhOyB9".to_string(),
950            filter_language: None,
951            address_book_config: None,
952        };
953        let resp = sdk.streams.test_filter(&params).await.unwrap();
954        assert!(resp.logs.is_empty());
955    }
956
957    #[tokio::test]
958    async fn test_filter_api_error() {
959        let server = MockServer::start().await;
960        Mock::given(method("POST"))
961            .and(path("/streams/test_filter"))
962            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
963            .mount(&server)
964            .await;
965        let sdk = make_sdk(format!("{}/", server.uri()));
966        let params = TestFilterParams {
967            network: "ethereum-mainnet".to_string(),
968            dataset: StreamDataset::Block,
969            block: "17811625".to_string(),
970            filter_function: "ZnVuY3Rpb24gbWFpbihkYXRhKSB7IHJldHVybiBkYXRhOyB9".to_string(),
971            filter_language: None,
972            address_book_config: None,
973        };
974        let err = sdk.streams.test_filter(&params).await.unwrap_err();
975        assert!(matches!(err, SdkError::Api { .. }));
976    }
977
978    #[tokio::test]
979    async fn get_enabled_count_success() {
980        let server = MockServer::start().await;
981        Mock::given(method("GET"))
982            .and(path("/streams/enabled_count"))
983            .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"total": 3})))
984            .mount(&server)
985            .await;
986        let sdk = make_sdk(format!("{}/", server.uri()));
987        let resp = sdk.streams.get_enabled_count(None).await.unwrap();
988        assert_eq!(resp.total, 3);
989    }
990
991    #[tokio::test]
992    async fn get_enabled_count_api_error() {
993        let server = MockServer::start().await;
994        Mock::given(method("GET"))
995            .and(path("/streams/enabled_count"))
996            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
997            .mount(&server)
998            .await;
999        let sdk = make_sdk(format!("{}/", server.uri()));
1000        let err = sdk.streams.get_enabled_count(None).await.unwrap_err();
1001        assert!(matches!(err, SdkError::Api { .. }));
1002    }
1003
1004    #[tokio::test]
1005    async fn get_enabled_count_server_error() {
1006        let server = MockServer::start().await;
1007        Mock::given(method("GET"))
1008            .and(path("/streams/enabled_count"))
1009            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
1010            .mount(&server)
1011            .await;
1012        let sdk = make_sdk(format!("{}/", server.uri()));
1013        let err = sdk.streams.get_enabled_count(None).await.unwrap_err();
1014        assert!(matches!(err, SdkError::Api { .. }));
1015    }
1016}