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        })
332        .unwrap()
333    }
334
335    fn webhook_params() -> CreateStreamParams {
336        CreateStreamParams {
337            name: "test-stream".to_string(),
338            region: StreamRegion::UsaEast,
339            network: "ethereum-mainnet".to_string(),
340            dataset: StreamDataset::Block,
341            start_range: 17000000,
342            end_range: -1,
343            destination_attributes: DestinationAttributes::Webhook(WebhookAttributes {
344                url: "https://example.com/webhook".to_string(),
345                max_retry: 3,
346                retry_interval_sec: 1,
347                post_timeout_sec: 10,
348                compression: Some("none".to_string()),
349                security_token: None,
350            }),
351            plan: Some("growth_plan".to_string()),
352            threshold_fetch_buffer: Some(1000),
353            dataset_batch_size: 1,
354            max_batch_size: None,
355            max_buffer_range_size: None,
356            max_buffer_processing_workers: None,
357            keep_distance_from_tip: None,
358            filter_function: None,
359            filter_language: None,
360            address_book_config: None,
361            include_stream_metadata: None,
362            product_type: None,
363            status: None,
364            notification_email: None,
365            charge_min_cap: None,
366            fix_block_reorgs: None,
367            elastic_batch_enabled: true,
368            extra_destinations: None,
369        }
370    }
371
372    fn stream_response_json() -> serde_json::Value {
373        serde_json::json!({
374            "id": "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab",
375            "name": "test-stream",
376            "status": "active",
377            "created_at": "2026-03-19T12:00:00Z",
378            "updated_at": "2026-03-19T12:00:00Z",
379            "sequence": 0,
380            "network": "ethereum-mainnet",
381            "dataset": "block",
382            "region": "usa_east",
383            "destination": "webhook",
384            "destination_attributes": {
385                "url": "https://example.com/webhook",
386                "max_retry": 3,
387                "retry_interval_sec": 1,
388                "post_timeout_sec": 10,
389                "compression": "none"
390            },
391            "start_range": 17000000,
392            "end_range": -1
393        })
394    }
395
396    #[tokio::test]
397    async fn create_stream_success() {
398        let server = MockServer::start().await;
399
400        Mock::given(method("POST"))
401            .and(path("/streams"))
402            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
403            .mount(&server)
404            .await;
405
406        let sdk = make_sdk(format!("{}/", server.uri()));
407        let resp = sdk.streams.create_stream(&webhook_params()).await.unwrap();
408
409        assert_eq!(resp.id, "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab");
410        assert_eq!(resp.name, "test-stream");
411        assert_eq!(resp.status, "active");
412        assert_eq!(resp.network, "ethereum-mainnet");
413        assert_eq!(resp.dataset, "block");
414        // Verify the full destination_attributes round-trip on the response
415        // side. Without this assertion, serde's flatten+Option silently
416        // swallows malformed destination_attributes as None.
417        match resp.destination_attributes {
418            Some(DestinationAttributes::Webhook(attrs)) => {
419                assert_eq!(attrs.url, "https://example.com/webhook");
420                assert_eq!(attrs.max_retry, 3);
421            }
422            other => panic!("expected Webhook destination, got {other:?}"),
423        }
424    }
425
426    #[tokio::test]
427    async fn create_stream_sends_typed_webhook_destination() {
428        use wiremock::matchers::body_partial_json;
429        let server = MockServer::start().await;
430
431        Mock::given(method("POST"))
432            .and(path("/streams"))
433            .and(body_partial_json(serde_json::json!({
434                "destination": "webhook",
435                "destination_attributes": {
436                    "url": "https://example.com/webhook",
437                    "max_retry": 3,
438                    "retry_interval_sec": 1,
439                    "post_timeout_sec": 10,
440                    "compression": "none"
441                }
442            })))
443            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
444            .mount(&server)
445            .await;
446
447        let sdk = make_sdk(format!("{}/", server.uri()));
448        let params = CreateStreamParams {
449            destination_attributes: DestinationAttributes::Webhook(WebhookAttributes {
450                url: "https://example.com/webhook".to_string(),
451                max_retry: 3,
452                retry_interval_sec: 1,
453                post_timeout_sec: 10,
454                compression: Some("none".to_string()),
455                security_token: None,
456            }),
457            ..webhook_params()
458        };
459        sdk.streams.create_stream(&params).await.unwrap();
460    }
461
462    #[tokio::test]
463    async fn create_stream_sends_extra_destinations() {
464        use wiremock::matchers::body_partial_json;
465        let server = MockServer::start().await;
466
467        Mock::given(method("POST"))
468            .and(path("/streams"))
469            .and(body_partial_json(serde_json::json!({
470                "extra_destinations": [
471                    {
472                        "destination": "webhook",
473                        "destination_attributes": {
474                            "url": "https://example.com/extra-hook",
475                            "max_retry": 5,
476                            "retry_interval_sec": 2,
477                            "post_timeout_sec": 15,
478                            "compression": "none"
479                        }
480                    },
481                    {
482                        "destination": "s3",
483                        "destination_attributes": {
484                            "endpoint": "s3.example.com",
485                            "access_key": "AKIA",
486                            "secret_key": "secret",
487                            "bucket": "my-bucket",
488                            "object_prefix": "streams/",
489                            "compression": "gzip",
490                            "file_type": ".json",
491                            "max_retry": 3,
492                            "retry_interval_sec": 1
493                        }
494                    }
495                ]
496            })))
497            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
498            .mount(&server)
499            .await;
500
501        let sdk = make_sdk(format!("{}/", server.uri()));
502        let params = CreateStreamParams {
503            extra_destinations: Some(vec![
504                DestinationAttributes::Webhook(WebhookAttributes {
505                    url: "https://example.com/extra-hook".to_string(),
506                    max_retry: 5,
507                    retry_interval_sec: 2,
508                    post_timeout_sec: 15,
509                    compression: Some("none".to_string()),
510                    security_token: None,
511                }),
512                DestinationAttributes::S3(S3Attributes {
513                    endpoint: "s3.example.com".to_string(),
514                    access_key: "AKIA".to_string(),
515                    secret_key: "secret".to_string(),
516                    bucket: "my-bucket".to_string(),
517                    object_prefix: "streams/".to_string(),
518                    compression: "gzip".to_string(),
519                    file_type: ".json".to_string(),
520                    max_retry: 3,
521                    retry_interval_sec: 1,
522                    use_ssl: None,
523                }),
524            ]),
525            ..webhook_params()
526        };
527        sdk.streams.create_stream(&params).await.unwrap();
528    }
529
530    #[tokio::test]
531    async fn update_stream_sends_extra_destinations() {
532        use wiremock::matchers::body_partial_json;
533        let server = MockServer::start().await;
534
535        Mock::given(method("PATCH"))
536            .and(path("/streams/test-id"))
537            .and(body_partial_json(serde_json::json!({
538                "extra_destinations": [
539                    {
540                        "destination": "webhook",
541                        "destination_attributes": {
542                            "url": "https://example.com/patched-hook",
543                            "max_retry": 1,
544                            "retry_interval_sec": 1,
545                            "post_timeout_sec": 5,
546                            "compression": "none"
547                        }
548                    }
549                ]
550            })))
551            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
552            .mount(&server)
553            .await;
554
555        let sdk = make_sdk(format!("{}/", server.uri()));
556        let params = UpdateStreamParams {
557            extra_destinations: Some(vec![DestinationAttributes::Webhook(WebhookAttributes {
558                url: "https://example.com/patched-hook".to_string(),
559                max_retry: 1,
560                retry_interval_sec: 1,
561                post_timeout_sec: 5,
562                compression: Some("none".to_string()),
563                security_token: None,
564            })]),
565            ..Default::default()
566        };
567        sdk.streams.update_stream("test-id", &params).await.unwrap();
568    }
569
570    #[tokio::test]
571    async fn get_stream_parses_extra_destinations() {
572        let server = MockServer::start().await;
573        let mut body = stream_response_json();
574        body["extra_destinations"] = serde_json::json!([
575            {
576                "destination": "webhook",
577                "destination_attributes": {
578                    "url": "https://example.com/extra",
579                    "max_retry": 2,
580                    "retry_interval_sec": 1,
581                    "post_timeout_sec": 10,
582                    "compression": "none"
583                }
584            }
585        ]);
586        Mock::given(method("GET"))
587            .and(path("/streams/test-id"))
588            .respond_with(ResponseTemplate::new(200).set_body_json(body))
589            .mount(&server)
590            .await;
591
592        let sdk = make_sdk(format!("{}/", server.uri()));
593        let resp = sdk.streams.get_stream("test-id").await.unwrap();
594        let extras = resp.extra_destinations.expect("extra_destinations present");
595        assert_eq!(extras.len(), 1);
596        match &extras[0] {
597            DestinationAttributes::Webhook(w) => {
598                assert_eq!(w.url, "https://example.com/extra");
599                assert_eq!(w.max_retry, 2);
600            }
601            other => panic!("expected Webhook, got {other:?}"),
602        }
603    }
604
605    #[tokio::test]
606    async fn create_stream_api_error() {
607        let server = MockServer::start().await;
608
609        Mock::given(method("POST"))
610            .and(path("/streams"))
611            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
612            .mount(&server)
613            .await;
614
615        let sdk = make_sdk(format!("{}/", server.uri()));
616        let err = sdk
617            .streams
618            .create_stream(&webhook_params())
619            .await
620            .unwrap_err();
621
622        assert!(matches!(err, SdkError::Api { .. }));
623    }
624
625    #[tokio::test]
626    async fn create_stream_server_error() {
627        let server = MockServer::start().await;
628
629        Mock::given(method("POST"))
630            .and(path("/streams"))
631            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
632            .mount(&server)
633            .await;
634
635        let sdk = make_sdk(format!("{}/", server.uri()));
636        let err = sdk
637            .streams
638            .create_stream(&webhook_params())
639            .await
640            .unwrap_err();
641
642        assert!(matches!(err, SdkError::Api { .. }));
643    }
644
645    #[tokio::test]
646    async fn list_streams_success() {
647        let server = MockServer::start().await;
648        let response = serde_json::json!({
649            "data": [stream_response_json()],
650            "pageInfo": { "limit": 100, "offset": 0, "total": 1 }
651        });
652        Mock::given(method("GET"))
653            .and(path("/streams"))
654            .respond_with(ResponseTemplate::new(200).set_body_json(response))
655            .mount(&server)
656            .await;
657        let sdk = make_sdk(format!("{}/", server.uri()));
658        let resp = sdk
659            .streams
660            .list_streams(&ListStreamsParams::default())
661            .await
662            .unwrap();
663        assert_eq!(resp.data.len(), 1);
664        assert_eq!(resp.page_info.total, 1);
665    }
666
667    #[tokio::test]
668    async fn list_streams_api_error() {
669        let server = MockServer::start().await;
670        Mock::given(method("GET"))
671            .and(path("/streams"))
672            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
673            .mount(&server)
674            .await;
675        let sdk = make_sdk(format!("{}/", server.uri()));
676        let err = sdk
677            .streams
678            .list_streams(&ListStreamsParams::default())
679            .await
680            .unwrap_err();
681        assert!(matches!(err, SdkError::Api { .. }));
682    }
683
684    #[tokio::test]
685    async fn list_streams_server_error() {
686        let server = MockServer::start().await;
687        Mock::given(method("GET"))
688            .and(path("/streams"))
689            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
690            .mount(&server)
691            .await;
692        let sdk = make_sdk(format!("{}/", server.uri()));
693        let err = sdk
694            .streams
695            .list_streams(&ListStreamsParams::default())
696            .await
697            .unwrap_err();
698        assert!(matches!(err, SdkError::Api { .. }));
699    }
700
701    #[tokio::test]
702    async fn get_stream_success() {
703        let server = MockServer::start().await;
704        Mock::given(method("GET"))
705            .and(path("/streams/test-id"))
706            .respond_with(ResponseTemplate::new(200).set_body_json(stream_response_json()))
707            .mount(&server)
708            .await;
709        let sdk = make_sdk(format!("{}/", server.uri()));
710        let resp = sdk.streams.get_stream("test-id").await.unwrap();
711        assert_eq!(resp.id, "7d3c1a22-4f9e-4b1e-8b3d-1234567890ab");
712    }
713
714    #[tokio::test]
715    async fn get_stream_not_found() {
716        let server = MockServer::start().await;
717        Mock::given(method("GET"))
718            .and(path("/streams/test-id"))
719            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
720            .mount(&server)
721            .await;
722        let sdk = make_sdk(format!("{}/", server.uri()));
723        let err = sdk.streams.get_stream("test-id").await.unwrap_err();
724        assert!(matches!(err, SdkError::Api { .. }));
725    }
726
727    #[tokio::test]
728    async fn get_stream_server_error() {
729        let server = MockServer::start().await;
730        Mock::given(method("GET"))
731            .and(path("/streams/test-id"))
732            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
733            .mount(&server)
734            .await;
735        let sdk = make_sdk(format!("{}/", server.uri()));
736        let err = sdk.streams.get_stream("test-id").await.unwrap_err();
737        assert!(matches!(err, SdkError::Api { .. }));
738    }
739
740    #[tokio::test]
741    async fn update_stream_success() {
742        let server = MockServer::start().await;
743        let mut updated = stream_response_json();
744        updated["name"] = serde_json::json!("updated-name");
745        Mock::given(method("PATCH"))
746            .and(path("/streams/test-id"))
747            .respond_with(ResponseTemplate::new(200).set_body_json(updated))
748            .mount(&server)
749            .await;
750        let sdk = make_sdk(format!("{}/", server.uri()));
751        let params = UpdateStreamParams {
752            name: Some("updated-name".to_string()),
753            ..Default::default()
754        };
755        let resp = sdk.streams.update_stream("test-id", &params).await.unwrap();
756        assert_eq!(resp.name, "updated-name");
757    }
758
759    #[tokio::test]
760    async fn update_stream_api_error() {
761        let server = MockServer::start().await;
762        Mock::given(method("PATCH"))
763            .and(path("/streams/test-id"))
764            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
765            .mount(&server)
766            .await;
767        let sdk = make_sdk(format!("{}/", server.uri()));
768        let params = UpdateStreamParams::default();
769        let err = sdk
770            .streams
771            .update_stream("test-id", &params)
772            .await
773            .unwrap_err();
774        assert!(matches!(err, SdkError::Api { .. }));
775    }
776
777    #[tokio::test]
778    async fn update_stream_server_error() {
779        let server = MockServer::start().await;
780        Mock::given(method("PATCH"))
781            .and(path("/streams/test-id"))
782            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
783            .mount(&server)
784            .await;
785        let sdk = make_sdk(format!("{}/", server.uri()));
786        let params = UpdateStreamParams::default();
787        let err = sdk
788            .streams
789            .update_stream("test-id", &params)
790            .await
791            .unwrap_err();
792        assert!(matches!(err, SdkError::Api { .. }));
793    }
794
795    #[tokio::test]
796    async fn delete_stream_success() {
797        let server = MockServer::start().await;
798        Mock::given(method("DELETE"))
799            .and(path("/streams/test-id"))
800            .respond_with(ResponseTemplate::new(200))
801            .mount(&server)
802            .await;
803        let sdk = make_sdk(format!("{}/", server.uri()));
804        sdk.streams.delete_stream("test-id").await.unwrap();
805    }
806
807    #[tokio::test]
808    async fn delete_stream_not_found() {
809        let server = MockServer::start().await;
810        Mock::given(method("DELETE"))
811            .and(path("/streams/test-id"))
812            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
813            .mount(&server)
814            .await;
815        let sdk = make_sdk(format!("{}/", server.uri()));
816        let err = sdk.streams.delete_stream("test-id").await.unwrap_err();
817        assert!(matches!(err, SdkError::Api { .. }));
818    }
819
820    #[tokio::test]
821    async fn delete_stream_server_error() {
822        let server = MockServer::start().await;
823        Mock::given(method("DELETE"))
824            .and(path("/streams/test-id"))
825            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
826            .mount(&server)
827            .await;
828        let sdk = make_sdk(format!("{}/", server.uri()));
829        let err = sdk.streams.delete_stream("test-id").await.unwrap_err();
830        assert!(matches!(err, SdkError::Api { .. }));
831    }
832
833    #[tokio::test]
834    async fn delete_all_streams_success() {
835        let server = MockServer::start().await;
836        Mock::given(method("DELETE"))
837            .and(path("/streams"))
838            .respond_with(ResponseTemplate::new(204))
839            .mount(&server)
840            .await;
841        let sdk = make_sdk(format!("{}/", server.uri()));
842        sdk.streams.delete_all_streams().await.unwrap();
843    }
844
845    #[tokio::test]
846    async fn delete_all_streams_not_found() {
847        let server = MockServer::start().await;
848        Mock::given(method("DELETE"))
849            .and(path("/streams"))
850            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
851            .mount(&server)
852            .await;
853        let sdk = make_sdk(format!("{}/", server.uri()));
854        let err = sdk.streams.delete_all_streams().await.unwrap_err();
855        assert!(matches!(err, SdkError::Api { .. }));
856    }
857
858    #[tokio::test]
859    async fn activate_stream_success() {
860        let server = MockServer::start().await;
861        Mock::given(method("POST"))
862            .and(path("/streams/test-id/activate"))
863            .respond_with(ResponseTemplate::new(201))
864            .mount(&server)
865            .await;
866        let sdk = make_sdk(format!("{}/", server.uri()));
867        sdk.streams.activate_stream("test-id").await.unwrap();
868    }
869
870    #[tokio::test]
871    async fn activate_stream_not_found() {
872        let server = MockServer::start().await;
873        Mock::given(method("POST"))
874            .and(path("/streams/test-id/activate"))
875            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
876            .mount(&server)
877            .await;
878        let sdk = make_sdk(format!("{}/", server.uri()));
879        let err = sdk.streams.activate_stream("test-id").await.unwrap_err();
880        assert!(matches!(err, SdkError::Api { .. }));
881    }
882
883    #[tokio::test]
884    async fn activate_stream_server_error() {
885        let server = MockServer::start().await;
886        Mock::given(method("POST"))
887            .and(path("/streams/test-id/activate"))
888            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
889            .mount(&server)
890            .await;
891        let sdk = make_sdk(format!("{}/", server.uri()));
892        let err = sdk.streams.activate_stream("test-id").await.unwrap_err();
893        assert!(matches!(err, SdkError::Api { .. }));
894    }
895
896    #[tokio::test]
897    async fn pause_stream_success() {
898        let server = MockServer::start().await;
899        Mock::given(method("POST"))
900            .and(path("/streams/test-id/pause"))
901            .respond_with(ResponseTemplate::new(201))
902            .mount(&server)
903            .await;
904        let sdk = make_sdk(format!("{}/", server.uri()));
905        sdk.streams.pause_stream("test-id").await.unwrap();
906    }
907
908    #[tokio::test]
909    async fn pause_stream_not_found() {
910        let server = MockServer::start().await;
911        Mock::given(method("POST"))
912            .and(path("/streams/test-id/pause"))
913            .respond_with(ResponseTemplate::new(404).set_body_string("Not Found"))
914            .mount(&server)
915            .await;
916        let sdk = make_sdk(format!("{}/", server.uri()));
917        let err = sdk.streams.pause_stream("test-id").await.unwrap_err();
918        assert!(matches!(err, SdkError::Api { .. }));
919    }
920
921    #[tokio::test]
922    async fn pause_stream_server_error() {
923        let server = MockServer::start().await;
924        Mock::given(method("POST"))
925            .and(path("/streams/test-id/pause"))
926            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
927            .mount(&server)
928            .await;
929        let sdk = make_sdk(format!("{}/", server.uri()));
930        let err = sdk.streams.pause_stream("test-id").await.unwrap_err();
931        assert!(matches!(err, SdkError::Api { .. }));
932    }
933
934    #[tokio::test]
935    async fn test_filter_success() {
936        let server = MockServer::start().await;
937        let response = serde_json::json!({ "result": {"hash": "0xabc"}, "logs": [] });
938        Mock::given(method("POST"))
939            .and(path("/streams/test_filter"))
940            .respond_with(ResponseTemplate::new(201).set_body_json(response))
941            .mount(&server)
942            .await;
943        let sdk = make_sdk(format!("{}/", server.uri()));
944        let params = TestFilterParams {
945            network: "ethereum-mainnet".to_string(),
946            dataset: StreamDataset::Block,
947            block: "17811625".to_string(),
948            filter_function: "ZnVuY3Rpb24gbWFpbihkYXRhKSB7IHJldHVybiBkYXRhOyB9".to_string(),
949            filter_language: None,
950            address_book_config: None,
951        };
952        let resp = sdk.streams.test_filter(&params).await.unwrap();
953        assert!(resp.logs.is_empty());
954    }
955
956    #[tokio::test]
957    async fn test_filter_api_error() {
958        let server = MockServer::start().await;
959        Mock::given(method("POST"))
960            .and(path("/streams/test_filter"))
961            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
962            .mount(&server)
963            .await;
964        let sdk = make_sdk(format!("{}/", server.uri()));
965        let params = TestFilterParams {
966            network: "ethereum-mainnet".to_string(),
967            dataset: StreamDataset::Block,
968            block: "17811625".to_string(),
969            filter_function: "ZnVuY3Rpb24gbWFpbihkYXRhKSB7IHJldHVybiBkYXRhOyB9".to_string(),
970            filter_language: None,
971            address_book_config: None,
972        };
973        let err = sdk.streams.test_filter(&params).await.unwrap_err();
974        assert!(matches!(err, SdkError::Api { .. }));
975    }
976
977    #[tokio::test]
978    async fn get_enabled_count_success() {
979        let server = MockServer::start().await;
980        Mock::given(method("GET"))
981            .and(path("/streams/enabled_count"))
982            .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"total": 3})))
983            .mount(&server)
984            .await;
985        let sdk = make_sdk(format!("{}/", server.uri()));
986        let resp = sdk.streams.get_enabled_count(None).await.unwrap();
987        assert_eq!(resp.total, 3);
988    }
989
990    #[tokio::test]
991    async fn get_enabled_count_api_error() {
992        let server = MockServer::start().await;
993        Mock::given(method("GET"))
994            .and(path("/streams/enabled_count"))
995            .respond_with(ResponseTemplate::new(400).set_body_string("Bad Request"))
996            .mount(&server)
997            .await;
998        let sdk = make_sdk(format!("{}/", server.uri()));
999        let err = sdk.streams.get_enabled_count(None).await.unwrap_err();
1000        assert!(matches!(err, SdkError::Api { .. }));
1001    }
1002
1003    #[tokio::test]
1004    async fn get_enabled_count_server_error() {
1005        let server = MockServer::start().await;
1006        Mock::given(method("GET"))
1007            .and(path("/streams/enabled_count"))
1008            .respond_with(ResponseTemplate::new(500).set_body_string("Internal Server Error"))
1009            .mount(&server)
1010            .await;
1011        let sdk = make_sdk(format!("{}/", server.uri()));
1012        let err = sdk.streams.get_enabled_count(None).await.unwrap_err();
1013        assert!(matches!(err, SdkError::Api { .. }));
1014    }
1015}