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#[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 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 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) = ¶ms.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) = ¶ms.order_by {
99 pairs.append_pair("order_by", v);
100 }
101 if let Some(v) = ¶ms.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 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 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 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 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 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 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 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 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#[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 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(¶ms).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(¶ms).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", ¶ms).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", ¶ms).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", ¶ms)
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", ¶ms)
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(¶ms).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(¶ms).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}