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 })
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 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(¶ms).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(¶ms).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", ¶ms).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", ¶ms).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", ¶ms)
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", ¶ms)
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(¶ms).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(¶ms).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}