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