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