1#![allow(unused_imports)]
15use async_trait::async_trait;
16use derive_builder::Builder;
17use serde::{Deserialize, Serialize};
18use serde_json::Value;
19use std::{collections::HashMap, sync::Arc};
20
21use crate::alpha::websocket_streams::models;
22use crate::common::{
23 models::ParamBuildError,
24 utils::replace_websocket_streams_placeholders,
25 websocket::{WebsocketBase, WebsocketStream, WebsocketStreams, create_stream_handler},
26};
27use crate::models::StreamId;
28
29#[async_trait]
30pub trait Api: Send + Sync {
31 async fn aggregate_trade_stream(
32 &self,
33 params: AggregateTradeStreamParams,
34 ) -> anyhow::Result<Arc<WebsocketStream<models::AggregateTradeStreamResponse>>>;
35 async fn all_book_ticker_stream(
36 &self,
37 params: AllBookTickerStreamParams,
38 ) -> anyhow::Result<Arc<WebsocketStream<models::AllBookTickerStreamResponse>>>;
39 async fn all_mini_ticker_stream(
40 &self,
41 params: AllMiniTickerStreamParams,
42 ) -> anyhow::Result<Arc<WebsocketStream<models::AllMiniTickerStreamResponse>>>;
43 async fn all_ticker_stream(
44 &self,
45 params: AllTickerStreamParams,
46 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTickerStreamResponse>>>;
47 async fn all_tokens24h_ticker_stream(
48 &self,
49 params: AllTokens24hTickerStreamParams,
50 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTokens24hTickerStreamResponse>>>;
51 async fn book_ticker_stream(
52 &self,
53 params: BookTickerStreamParams,
54 ) -> anyhow::Result<Arc<WebsocketStream<models::BookTickerStreamResponse>>>;
55 async fn contract_kline_stream(
56 &self,
57 params: ContractKlineStreamParams,
58 ) -> anyhow::Result<Arc<WebsocketStream<models::ContractKlineStreamResponse>>>;
59 async fn full_depth_stream(
60 &self,
61 params: FullDepthStreamParams,
62 ) -> anyhow::Result<Arc<WebsocketStream<models::FullDepthStreamResponse>>>;
63 async fn kline_stream(
64 &self,
65 params: KlineStreamParams,
66 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineStreamResponse>>>;
67 async fn mini_ticker_stream(
68 &self,
69 params: MiniTickerStreamParams,
70 ) -> anyhow::Result<Arc<WebsocketStream<models::MiniTickerStreamResponse>>>;
71 async fn partial_depth_stream(
72 &self,
73 params: PartialDepthStreamParams,
74 ) -> anyhow::Result<Arc<WebsocketStream<models::PartialDepthStreamResponse>>>;
75 async fn ticker_stream(
76 &self,
77 params: TickerStreamParams,
78 ) -> anyhow::Result<Arc<WebsocketStream<models::TickerStreamResponse>>>;
79 async fn trade_stream(
80 &self,
81 params: TradeStreamParams,
82 ) -> anyhow::Result<Arc<WebsocketStream<models::TradeStreamResponse>>>;
83}
84
85pub struct ApiClient {
86 websocket_streams_base: Arc<WebsocketStreams>,
87}
88
89impl ApiClient {
90 pub fn new(websocket_streams_base: Arc<WebsocketStreams>) -> Self {
91 Self {
92 websocket_streams_base,
93 }
94 }
95}
96
97#[allow(non_camel_case_types)]
98#[derive(Debug, Clone, Serialize, Deserialize)]
99pub enum ContractKlineStreamIntervalEnum {
100 #[serde(rename = "1s")]
101 Interval1s,
102 #[serde(rename = "1m")]
103 Interval1m,
104 #[serde(rename = "5m")]
105 Interval5m,
106 #[serde(rename = "15m")]
107 Interval15m,
108 #[serde(rename = "1h")]
109 Interval1h,
110 #[serde(rename = "4h")]
111 Interval4h,
112 #[serde(rename = "1d")]
113 Interval1d,
114}
115
116impl ContractKlineStreamIntervalEnum {
117 #[must_use]
118 pub fn as_str(&self) -> &'static str {
119 match self {
120 Self::Interval1s => "1s",
121 Self::Interval1m => "1m",
122 Self::Interval5m => "5m",
123 Self::Interval15m => "15m",
124 Self::Interval1h => "1h",
125 Self::Interval4h => "4h",
126 Self::Interval1d => "1d",
127 }
128 }
129}
130
131impl std::str::FromStr for ContractKlineStreamIntervalEnum {
132 type Err = Box<dyn std::error::Error + Send + Sync>;
133
134 fn from_str(s: &str) -> Result<Self, Self::Err> {
135 match s {
136 "1s" => Ok(Self::Interval1s),
137 "1m" => Ok(Self::Interval1m),
138 "5m" => Ok(Self::Interval5m),
139 "15m" => Ok(Self::Interval15m),
140 "1h" => Ok(Self::Interval1h),
141 "4h" => Ok(Self::Interval4h),
142 "1d" => Ok(Self::Interval1d),
143 other => Err(format!("invalid ContractKlineStreamIntervalEnum: {}", other).into()),
144 }
145 }
146}
147
148#[allow(non_camel_case_types)]
149#[derive(Debug, Clone, Serialize, Deserialize)]
150pub enum FullDepthStreamIntervalEnum {
151 #[serde(rename = "0ms")]
152 Interval0ms,
153 #[serde(rename = "100ms")]
154 Interval100ms,
155 #[serde(rename = "500ms")]
156 Interval500ms,
157}
158
159impl FullDepthStreamIntervalEnum {
160 #[must_use]
161 pub fn as_str(&self) -> &'static str {
162 match self {
163 Self::Interval0ms => "0ms",
164 Self::Interval100ms => "100ms",
165 Self::Interval500ms => "500ms",
166 }
167 }
168}
169
170impl std::str::FromStr for FullDepthStreamIntervalEnum {
171 type Err = Box<dyn std::error::Error + Send + Sync>;
172
173 fn from_str(s: &str) -> Result<Self, Self::Err> {
174 match s {
175 "0ms" => Ok(Self::Interval0ms),
176 "100ms" => Ok(Self::Interval100ms),
177 "500ms" => Ok(Self::Interval500ms),
178 other => Err(format!("invalid FullDepthStreamIntervalEnum: {}", other).into()),
179 }
180 }
181}
182
183#[allow(non_camel_case_types)]
184#[derive(Debug, Clone, Serialize, Deserialize)]
185pub enum KlineStreamIntervalEnum {
186 #[serde(rename = "1m")]
187 Interval1m,
188 #[serde(rename = "3m")]
189 Interval3m,
190 #[serde(rename = "5m")]
191 Interval5m,
192 #[serde(rename = "15m")]
193 Interval15m,
194 #[serde(rename = "30m")]
195 Interval30m,
196 #[serde(rename = "1h")]
197 Interval1h,
198 #[serde(rename = "2h")]
199 Interval2h,
200 #[serde(rename = "4h")]
201 Interval4h,
202 #[serde(rename = "6h")]
203 Interval6h,
204 #[serde(rename = "8h")]
205 Interval8h,
206 #[serde(rename = "12h")]
207 Interval12h,
208 #[serde(rename = "1d")]
209 Interval1d,
210 #[serde(rename = "3d")]
211 Interval3d,
212 #[serde(rename = "1w")]
213 Interval1w,
214 #[serde(rename = "1M")]
215 Interval1M,
216}
217
218impl KlineStreamIntervalEnum {
219 #[must_use]
220 pub fn as_str(&self) -> &'static str {
221 match self {
222 Self::Interval1m => "1m",
223 Self::Interval3m => "3m",
224 Self::Interval5m => "5m",
225 Self::Interval15m => "15m",
226 Self::Interval30m => "30m",
227 Self::Interval1h => "1h",
228 Self::Interval2h => "2h",
229 Self::Interval4h => "4h",
230 Self::Interval6h => "6h",
231 Self::Interval8h => "8h",
232 Self::Interval12h => "12h",
233 Self::Interval1d => "1d",
234 Self::Interval3d => "3d",
235 Self::Interval1w => "1w",
236 Self::Interval1M => "1M",
237 }
238 }
239}
240
241impl std::str::FromStr for KlineStreamIntervalEnum {
242 type Err = Box<dyn std::error::Error + Send + Sync>;
243
244 fn from_str(s: &str) -> Result<Self, Self::Err> {
245 match s {
246 "1m" => Ok(Self::Interval1m),
247 "3m" => Ok(Self::Interval3m),
248 "5m" => Ok(Self::Interval5m),
249 "15m" => Ok(Self::Interval15m),
250 "30m" => Ok(Self::Interval30m),
251 "1h" => Ok(Self::Interval1h),
252 "2h" => Ok(Self::Interval2h),
253 "4h" => Ok(Self::Interval4h),
254 "6h" => Ok(Self::Interval6h),
255 "8h" => Ok(Self::Interval8h),
256 "12h" => Ok(Self::Interval12h),
257 "1d" => Ok(Self::Interval1d),
258 "3d" => Ok(Self::Interval3d),
259 "1w" => Ok(Self::Interval1w),
260 "1M" => Ok(Self::Interval1M),
261 other => Err(format!("invalid KlineStreamIntervalEnum: {}", other).into()),
262 }
263 }
264}
265
266#[allow(non_camel_case_types)]
267#[derive(Debug, Clone, Serialize, Deserialize)]
268pub enum PartialDepthStreamLevelsEnum {
269 #[serde(rename = "5")]
270 Levels5,
271 #[serde(rename = "10")]
272 Levels10,
273 #[serde(rename = "20")]
274 Levels20,
275}
276
277impl PartialDepthStreamLevelsEnum {
278 #[must_use]
279 pub fn as_str(&self) -> &'static str {
280 match self {
281 Self::Levels5 => "5",
282 Self::Levels10 => "10",
283 Self::Levels20 => "20",
284 }
285 }
286}
287
288impl std::str::FromStr for PartialDepthStreamLevelsEnum {
289 type Err = Box<dyn std::error::Error + Send + Sync>;
290
291 fn from_str(s: &str) -> Result<Self, Self::Err> {
292 match s {
293 "5" => Ok(Self::Levels5),
294 "10" => Ok(Self::Levels10),
295 "20" => Ok(Self::Levels20),
296 other => Err(format!("invalid PartialDepthStreamLevelsEnum: {}", other).into()),
297 }
298 }
299}
300
301#[allow(non_camel_case_types)]
302#[derive(Debug, Clone, Serialize, Deserialize)]
303pub enum PartialDepthStreamIntervalEnum {
304 #[serde(rename = "0ms")]
305 Interval0ms,
306 #[serde(rename = "100ms")]
307 Interval100ms,
308 #[serde(rename = "500ms")]
309 Interval500ms,
310}
311
312impl PartialDepthStreamIntervalEnum {
313 #[must_use]
314 pub fn as_str(&self) -> &'static str {
315 match self {
316 Self::Interval0ms => "0ms",
317 Self::Interval100ms => "100ms",
318 Self::Interval500ms => "500ms",
319 }
320 }
321}
322
323impl std::str::FromStr for PartialDepthStreamIntervalEnum {
324 type Err = Box<dyn std::error::Error + Send + Sync>;
325
326 fn from_str(s: &str) -> Result<Self, Self::Err> {
327 match s {
328 "0ms" => Ok(Self::Interval0ms),
329 "100ms" => Ok(Self::Interval100ms),
330 "500ms" => Ok(Self::Interval500ms),
331 other => Err(format!("invalid PartialDepthStreamIntervalEnum: {}", other).into()),
332 }
333 }
334}
335
336#[derive(Clone, Debug, Builder, Deserialize)]
341#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
342pub struct AggregateTradeStreamParams {
343 #[builder(setter(into))]
347 #[serde(rename = "symbol")]
348 pub symbol: String,
349 #[builder(setter(into), default)]
353 #[serde(rename = "id", default)]
354 pub id: Option<String>,
355}
356
357impl AggregateTradeStreamParams {
358 #[must_use]
365 pub fn builder(symbol: String) -> AggregateTradeStreamParamsBuilder {
366 AggregateTradeStreamParamsBuilder::default().symbol(symbol)
367 }
368}
369#[derive(Clone, Debug, Builder, Deserialize, Default)]
374#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
375pub struct AllBookTickerStreamParams {
376 #[builder(setter(into), default)]
380 #[serde(rename = "id", default)]
381 pub id: Option<String>,
382}
383
384impl AllBookTickerStreamParams {
385 #[must_use]
388 pub fn builder() -> AllBookTickerStreamParamsBuilder {
389 AllBookTickerStreamParamsBuilder::default()
390 }
391}
392#[derive(Clone, Debug, Builder, Deserialize, Default)]
397#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
398pub struct AllMiniTickerStreamParams {
399 #[builder(setter(into), default)]
403 #[serde(rename = "id", default)]
404 pub id: Option<String>,
405}
406
407impl AllMiniTickerStreamParams {
408 #[must_use]
411 pub fn builder() -> AllMiniTickerStreamParamsBuilder {
412 AllMiniTickerStreamParamsBuilder::default()
413 }
414}
415#[derive(Clone, Debug, Builder, Deserialize, Default)]
420#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
421pub struct AllTickerStreamParams {
422 #[builder(setter(into), default)]
426 #[serde(rename = "id", default)]
427 pub id: Option<String>,
428}
429
430impl AllTickerStreamParams {
431 #[must_use]
434 pub fn builder() -> AllTickerStreamParamsBuilder {
435 AllTickerStreamParamsBuilder::default()
436 }
437}
438#[derive(Clone, Debug, Builder, Deserialize, Default)]
443#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
444pub struct AllTokens24hTickerStreamParams {
445 #[builder(setter(into), default)]
449 #[serde(rename = "id", default)]
450 pub id: Option<String>,
451}
452
453impl AllTokens24hTickerStreamParams {
454 #[must_use]
457 pub fn builder() -> AllTokens24hTickerStreamParamsBuilder {
458 AllTokens24hTickerStreamParamsBuilder::default()
459 }
460}
461#[derive(Clone, Debug, Builder, Deserialize)]
466#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
467pub struct BookTickerStreamParams {
468 #[builder(setter(into))]
472 #[serde(rename = "symbol")]
473 pub symbol: String,
474 #[builder(setter(into), default)]
478 #[serde(rename = "id", default)]
479 pub id: Option<String>,
480}
481
482impl BookTickerStreamParams {
483 #[must_use]
490 pub fn builder(symbol: String) -> BookTickerStreamParamsBuilder {
491 BookTickerStreamParamsBuilder::default().symbol(symbol)
492 }
493}
494#[derive(Clone, Debug, Builder, Deserialize)]
499#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
500pub struct ContractKlineStreamParams {
501 #[builder(setter(into))]
505 #[serde(rename = "contractAddress")]
506 pub contract_address: String,
507 #[builder(setter(into))]
511 #[serde(rename = "chainId")]
512 pub chain_id: String,
513 #[builder(setter(into))]
517 #[serde(rename = "interval")]
518 pub interval: ContractKlineStreamIntervalEnum,
519 #[builder(setter(into), default)]
523 #[serde(rename = "id", default)]
524 pub id: Option<String>,
525}
526
527impl ContractKlineStreamParams {
528 #[must_use]
537 pub fn builder(
538 contract_address: String,
539 chain_id: String,
540 interval: ContractKlineStreamIntervalEnum,
541 ) -> ContractKlineStreamParamsBuilder {
542 ContractKlineStreamParamsBuilder::default()
543 .contract_address(contract_address)
544 .chain_id(chain_id)
545 .interval(interval)
546 }
547}
548#[derive(Clone, Debug, Builder, Deserialize)]
553#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
554pub struct FullDepthStreamParams {
555 #[builder(setter(into))]
559 #[serde(rename = "symbol")]
560 pub symbol: String,
561 #[builder(setter(into))]
565 #[serde(rename = "interval")]
566 pub interval: FullDepthStreamIntervalEnum,
567 #[builder(setter(into), default)]
571 #[serde(rename = "id", default)]
572 pub id: Option<String>,
573}
574
575impl FullDepthStreamParams {
576 #[must_use]
584 pub fn builder(
585 symbol: String,
586 interval: FullDepthStreamIntervalEnum,
587 ) -> FullDepthStreamParamsBuilder {
588 FullDepthStreamParamsBuilder::default()
589 .symbol(symbol)
590 .interval(interval)
591 }
592}
593#[derive(Clone, Debug, Builder, Deserialize)]
598#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
599pub struct KlineStreamParams {
600 #[builder(setter(into))]
604 #[serde(rename = "symbol")]
605 pub symbol: String,
606 #[builder(setter(into))]
610 #[serde(rename = "interval")]
611 pub interval: KlineStreamIntervalEnum,
612 #[builder(setter(into), default)]
616 #[serde(rename = "id", default)]
617 pub id: Option<String>,
618}
619
620impl KlineStreamParams {
621 #[must_use]
629 pub fn builder(symbol: String, interval: KlineStreamIntervalEnum) -> KlineStreamParamsBuilder {
630 KlineStreamParamsBuilder::default()
631 .symbol(symbol)
632 .interval(interval)
633 }
634}
635#[derive(Clone, Debug, Builder, Deserialize)]
640#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
641pub struct MiniTickerStreamParams {
642 #[builder(setter(into))]
646 #[serde(rename = "symbol")]
647 pub symbol: String,
648 #[builder(setter(into), default)]
652 #[serde(rename = "id", default)]
653 pub id: Option<String>,
654}
655
656impl MiniTickerStreamParams {
657 #[must_use]
664 pub fn builder(symbol: String) -> MiniTickerStreamParamsBuilder {
665 MiniTickerStreamParamsBuilder::default().symbol(symbol)
666 }
667}
668#[derive(Clone, Debug, Builder, Deserialize)]
673#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
674pub struct PartialDepthStreamParams {
675 #[builder(setter(into))]
679 #[serde(rename = "symbol")]
680 pub symbol: String,
681 #[builder(setter(into))]
685 #[serde(rename = "levels")]
686 pub levels: PartialDepthStreamLevelsEnum,
687 #[builder(setter(into))]
691 #[serde(rename = "interval")]
692 pub interval: PartialDepthStreamIntervalEnum,
693 #[builder(setter(into), default)]
697 #[serde(rename = "id", default)]
698 pub id: Option<String>,
699}
700
701impl PartialDepthStreamParams {
702 #[must_use]
711 pub fn builder(
712 symbol: String,
713 levels: PartialDepthStreamLevelsEnum,
714 interval: PartialDepthStreamIntervalEnum,
715 ) -> PartialDepthStreamParamsBuilder {
716 PartialDepthStreamParamsBuilder::default()
717 .symbol(symbol)
718 .levels(levels)
719 .interval(interval)
720 }
721}
722#[derive(Clone, Debug, Builder, Deserialize)]
727#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
728pub struct TickerStreamParams {
729 #[builder(setter(into))]
733 #[serde(rename = "symbol")]
734 pub symbol: String,
735 #[builder(setter(into), default)]
739 #[serde(rename = "id", default)]
740 pub id: Option<String>,
741}
742
743impl TickerStreamParams {
744 #[must_use]
751 pub fn builder(symbol: String) -> TickerStreamParamsBuilder {
752 TickerStreamParamsBuilder::default().symbol(symbol)
753 }
754}
755#[derive(Clone, Debug, Builder, Deserialize)]
760#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
761pub struct TradeStreamParams {
762 #[builder(setter(into))]
766 #[serde(rename = "symbol")]
767 pub symbol: String,
768 #[builder(setter(into), default)]
772 #[serde(rename = "id", default)]
773 pub id: Option<String>,
774}
775
776impl TradeStreamParams {
777 #[must_use]
784 pub fn builder(symbol: String) -> TradeStreamParamsBuilder {
785 TradeStreamParamsBuilder::default().symbol(symbol)
786 }
787}
788
789#[async_trait]
790impl Api for ApiClient {
791 async fn aggregate_trade_stream(
792 &self,
793 params: AggregateTradeStreamParams,
794 ) -> anyhow::Result<Arc<WebsocketStream<models::AggregateTradeStreamResponse>>> {
795 let AggregateTradeStreamParams { symbol, id } = params;
796
797 let pairs: &[(&str, Option<String>)] =
798 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
799
800 let vars: HashMap<_, _> = pairs
801 .iter()
802 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
803 .collect();
804
805 let id_opt: Option<String> = vars.get("id").cloned();
806
807 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
808
809 Ok(
810 create_stream_handler::<models::AggregateTradeStreamResponse>(
811 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
812 stream,
813 id_opt.map(|s| {
814 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
815 if let Ok(n) = s.parse::<u32>() {
816 return StreamId::Number(n);
817 }
818 }
819 StreamId::Str(s)
820 }),
821 None,
822 )
823 .await,
824 )
825 }
826
827 async fn all_book_ticker_stream(
828 &self,
829 params: AllBookTickerStreamParams,
830 ) -> anyhow::Result<Arc<WebsocketStream<models::AllBookTickerStreamResponse>>> {
831 let AllBookTickerStreamParams { id } = params;
832
833 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
834
835 let vars: HashMap<_, _> = pairs
836 .iter()
837 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
838 .collect();
839
840 let id_opt: Option<String> = vars.get("id").cloned();
841
842 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
843
844 Ok(
845 create_stream_handler::<models::AllBookTickerStreamResponse>(
846 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
847 stream,
848 id_opt.map(|s| {
849 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
850 if let Ok(n) = s.parse::<u32>() {
851 return StreamId::Number(n);
852 }
853 }
854 StreamId::Str(s)
855 }),
856 None,
857 )
858 .await,
859 )
860 }
861
862 async fn all_mini_ticker_stream(
863 &self,
864 params: AllMiniTickerStreamParams,
865 ) -> anyhow::Result<Arc<WebsocketStream<models::AllMiniTickerStreamResponse>>> {
866 let AllMiniTickerStreamParams { id } = params;
867
868 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
869
870 let vars: HashMap<_, _> = pairs
871 .iter()
872 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
873 .collect();
874
875 let id_opt: Option<String> = vars.get("id").cloned();
876
877 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
878
879 Ok(
880 create_stream_handler::<models::AllMiniTickerStreamResponse>(
881 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
882 stream,
883 id_opt.map(|s| {
884 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
885 if let Ok(n) = s.parse::<u32>() {
886 return StreamId::Number(n);
887 }
888 }
889 StreamId::Str(s)
890 }),
891 None,
892 )
893 .await,
894 )
895 }
896
897 async fn all_ticker_stream(
898 &self,
899 params: AllTickerStreamParams,
900 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTickerStreamResponse>>> {
901 let AllTickerStreamParams { id } = params;
902
903 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
904
905 let vars: HashMap<_, _> = pairs
906 .iter()
907 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
908 .collect();
909
910 let id_opt: Option<String> = vars.get("id").cloned();
911
912 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
913
914 Ok(create_stream_handler::<models::AllTickerStreamResponse>(
915 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
916 stream,
917 id_opt.map(|s| {
918 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
919 if let Ok(n) = s.parse::<u32>() {
920 return StreamId::Number(n);
921 }
922 }
923 StreamId::Str(s)
924 }),
925 None,
926 )
927 .await)
928 }
929
930 async fn all_tokens24h_ticker_stream(
931 &self,
932 params: AllTokens24hTickerStreamParams,
933 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTokens24hTickerStreamResponse>>> {
934 let AllTokens24hTickerStreamParams { id } = params;
935
936 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
937
938 let vars: HashMap<_, _> = pairs
939 .iter()
940 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
941 .collect();
942
943 let id_opt: Option<String> = vars.get("id").cloned();
944
945 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
946
947 Ok(
948 create_stream_handler::<models::AllTokens24hTickerStreamResponse>(
949 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
950 stream,
951 id_opt.map(|s| {
952 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
953 if let Ok(n) = s.parse::<u32>() {
954 return StreamId::Number(n);
955 }
956 }
957 StreamId::Str(s)
958 }),
959 None,
960 )
961 .await,
962 )
963 }
964
965 async fn book_ticker_stream(
966 &self,
967 params: BookTickerStreamParams,
968 ) -> anyhow::Result<Arc<WebsocketStream<models::BookTickerStreamResponse>>> {
969 let BookTickerStreamParams { symbol, id } = params;
970
971 let pairs: &[(&str, Option<String>)] =
972 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
973
974 let vars: HashMap<_, _> = pairs
975 .iter()
976 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
977 .collect();
978
979 let id_opt: Option<String> = vars.get("id").cloned();
980
981 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
982
983 Ok(create_stream_handler::<models::BookTickerStreamResponse>(
984 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
985 stream,
986 id_opt.map(|s| {
987 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
988 if let Ok(n) = s.parse::<u32>() {
989 return StreamId::Number(n);
990 }
991 }
992 StreamId::Str(s)
993 }),
994 None,
995 )
996 .await)
997 }
998
999 async fn contract_kline_stream(
1000 &self,
1001 params: ContractKlineStreamParams,
1002 ) -> anyhow::Result<Arc<WebsocketStream<models::ContractKlineStreamResponse>>> {
1003 let ContractKlineStreamParams {
1004 contract_address,
1005 chain_id,
1006 interval,
1007 id,
1008 } = params;
1009
1010 let pairs: &[(&str, Option<String>)] = &[
1011 ("contractAddress", Some(contract_address.clone())),
1012 ("chainId", Some(chain_id.clone())),
1013 ("interval", Some(interval.as_str().to_string())),
1014 ("id", id.clone()),
1015 ];
1016
1017 let vars: HashMap<_, _> = pairs
1018 .iter()
1019 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1020 .collect();
1021
1022 let id_opt: Option<String> = vars.get("id").cloned();
1023
1024 let stream = replace_websocket_streams_placeholders(
1025 "/came@<contractAddress>@<chainId>@kline_<interval>",
1026 &vars,
1027 );
1028
1029 Ok(
1030 create_stream_handler::<models::ContractKlineStreamResponse>(
1031 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1032 stream,
1033 id_opt.map(|s| {
1034 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1035 if let Ok(n) = s.parse::<u32>() {
1036 return StreamId::Number(n);
1037 }
1038 }
1039 StreamId::Str(s)
1040 }),
1041 None,
1042 )
1043 .await,
1044 )
1045 }
1046
1047 async fn full_depth_stream(
1048 &self,
1049 params: FullDepthStreamParams,
1050 ) -> anyhow::Result<Arc<WebsocketStream<models::FullDepthStreamResponse>>> {
1051 let FullDepthStreamParams {
1052 symbol,
1053 interval,
1054 id,
1055 } = params;
1056
1057 let pairs: &[(&str, Option<String>)] = &[
1058 ("symbol", Some(symbol.clone())),
1059 ("interval", Some(interval.as_str().to_string())),
1060 ("id", id.clone()),
1061 ];
1062
1063 let vars: HashMap<_, _> = pairs
1064 .iter()
1065 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1066 .collect();
1067
1068 let id_opt: Option<String> = vars.get("id").cloned();
1069
1070 let stream =
1071 replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
1072
1073 Ok(create_stream_handler::<models::FullDepthStreamResponse>(
1074 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1075 stream,
1076 id_opt.map(|s| {
1077 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1078 if let Ok(n) = s.parse::<u32>() {
1079 return StreamId::Number(n);
1080 }
1081 }
1082 StreamId::Str(s)
1083 }),
1084 None,
1085 )
1086 .await)
1087 }
1088
1089 async fn kline_stream(
1090 &self,
1091 params: KlineStreamParams,
1092 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineStreamResponse>>> {
1093 let KlineStreamParams {
1094 symbol,
1095 interval,
1096 id,
1097 } = params;
1098
1099 let pairs: &[(&str, Option<String>)] = &[
1100 ("symbol", Some(symbol.clone())),
1101 ("interval", Some(interval.as_str().to_string())),
1102 ("id", id.clone()),
1103 ];
1104
1105 let vars: HashMap<_, _> = pairs
1106 .iter()
1107 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1108 .collect();
1109
1110 let id_opt: Option<String> = vars.get("id").cloned();
1111
1112 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
1113
1114 Ok(create_stream_handler::<models::KlineStreamResponse>(
1115 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1116 stream,
1117 id_opt.map(|s| {
1118 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1119 if let Ok(n) = s.parse::<u32>() {
1120 return StreamId::Number(n);
1121 }
1122 }
1123 StreamId::Str(s)
1124 }),
1125 None,
1126 )
1127 .await)
1128 }
1129
1130 async fn mini_ticker_stream(
1131 &self,
1132 params: MiniTickerStreamParams,
1133 ) -> anyhow::Result<Arc<WebsocketStream<models::MiniTickerStreamResponse>>> {
1134 let MiniTickerStreamParams { symbol, id } = params;
1135
1136 let pairs: &[(&str, Option<String>)] =
1137 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1138
1139 let vars: HashMap<_, _> = pairs
1140 .iter()
1141 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1142 .collect();
1143
1144 let id_opt: Option<String> = vars.get("id").cloned();
1145
1146 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
1147
1148 Ok(create_stream_handler::<models::MiniTickerStreamResponse>(
1149 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1150 stream,
1151 id_opt.map(|s| {
1152 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1153 if let Ok(n) = s.parse::<u32>() {
1154 return StreamId::Number(n);
1155 }
1156 }
1157 StreamId::Str(s)
1158 }),
1159 None,
1160 )
1161 .await)
1162 }
1163
1164 async fn partial_depth_stream(
1165 &self,
1166 params: PartialDepthStreamParams,
1167 ) -> anyhow::Result<Arc<WebsocketStream<models::PartialDepthStreamResponse>>> {
1168 let PartialDepthStreamParams {
1169 symbol,
1170 levels,
1171 interval,
1172 id,
1173 } = params;
1174
1175 let pairs: &[(&str, Option<String>)] = &[
1176 ("symbol", Some(symbol.clone())),
1177 ("levels", Some(levels.as_str().to_string())),
1178 ("interval", Some(interval.as_str().to_string())),
1179 ("id", id.clone()),
1180 ];
1181
1182 let vars: HashMap<_, _> = pairs
1183 .iter()
1184 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1185 .collect();
1186
1187 let id_opt: Option<String> = vars.get("id").cloned();
1188
1189 let stream =
1190 replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
1191
1192 Ok(create_stream_handler::<models::PartialDepthStreamResponse>(
1193 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1194 stream,
1195 id_opt.map(|s| {
1196 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1197 if let Ok(n) = s.parse::<u32>() {
1198 return StreamId::Number(n);
1199 }
1200 }
1201 StreamId::Str(s)
1202 }),
1203 None,
1204 )
1205 .await)
1206 }
1207
1208 async fn ticker_stream(
1209 &self,
1210 params: TickerStreamParams,
1211 ) -> anyhow::Result<Arc<WebsocketStream<models::TickerStreamResponse>>> {
1212 let TickerStreamParams { symbol, id } = params;
1213
1214 let pairs: &[(&str, Option<String>)] =
1215 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1216
1217 let vars: HashMap<_, _> = pairs
1218 .iter()
1219 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1220 .collect();
1221
1222 let id_opt: Option<String> = vars.get("id").cloned();
1223
1224 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
1225
1226 Ok(create_stream_handler::<models::TickerStreamResponse>(
1227 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1228 stream,
1229 id_opt.map(|s| {
1230 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1231 if let Ok(n) = s.parse::<u32>() {
1232 return StreamId::Number(n);
1233 }
1234 }
1235 StreamId::Str(s)
1236 }),
1237 None,
1238 )
1239 .await)
1240 }
1241
1242 async fn trade_stream(
1243 &self,
1244 params: TradeStreamParams,
1245 ) -> anyhow::Result<Arc<WebsocketStream<models::TradeStreamResponse>>> {
1246 let TradeStreamParams { symbol, id } = params;
1247
1248 let pairs: &[(&str, Option<String>)] =
1249 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1250
1251 let vars: HashMap<_, _> = pairs
1252 .iter()
1253 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1254 .collect();
1255
1256 let id_opt: Option<String> = vars.get("id").cloned();
1257
1258 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
1259
1260 Ok(create_stream_handler::<models::TradeStreamResponse>(
1261 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1262 stream,
1263 id_opt.map(|s| {
1264 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1265 if let Ok(n) = s.parse::<u32>() {
1266 return StreamId::Number(n);
1267 }
1268 }
1269 StreamId::Str(s)
1270 }),
1271 None,
1272 )
1273 .await)
1274 }
1275}
1276
1277#[cfg(all(test, feature = "alpha"))]
1278mod tests {
1279 use super::*;
1280 use crate::TOKIO_SHARED_RT;
1281 use crate::{
1282 common::websocket::{WebsocketConnection, WebsocketHandler},
1283 config::ConfigurationWebsocketStreams,
1284 };
1285 use serde_json::json;
1286 use std::sync::atomic::{AtomicBool, Ordering};
1287 use tokio::task::yield_now;
1288
1289 async fn make_streams_base() -> (Arc<WebsocketStreams>, Arc<WebsocketConnection>) {
1290 let conn = WebsocketConnection::new("test");
1291 let config = ConfigurationWebsocketStreams::builder()
1292 .build()
1293 .expect("Failed to build configuration");
1294 let streams_base = WebsocketStreams::new(config, vec![conn.clone()], vec![]);
1295 conn.set_handler(streams_base.clone() as Arc<dyn WebsocketHandler>)
1296 .await;
1297 (streams_base, conn)
1298 }
1299
1300 #[test]
1301 fn aggregate_trade_stream_should_execute_successfully() {
1302 TOKIO_SHARED_RT.block_on(async {
1303 let (streams_base, _) = make_streams_base().await;
1304 let api = ApiClient::new(streams_base.clone());
1305
1306 let id = "test-id-123".to_string();
1307
1308 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string())
1309 .id(Some(id.clone()))
1310 .build()
1311 .unwrap();
1312
1313 let AggregateTradeStreamParams { symbol, id } = params.clone();
1314
1315 let pairs: &[(&str, Option<String>)] =
1316 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1317
1318 let vars: HashMap<_, _> = pairs
1319 .iter()
1320 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1321 .collect();
1322 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1323 let ws_stream = api
1324 .aggregate_trade_stream(params)
1325 .await
1326 .expect("aggregate_trade_stream should return a WebsocketStream");
1327
1328 assert!(
1329 streams_base.is_subscribed(&stream).await,
1330 "expected stream '{stream}' to be subscribed"
1331 );
1332 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1333 });
1334 }
1335
1336 #[test]
1337 fn aggregate_trade_stream_should_handle_incoming_message() {
1338 TOKIO_SHARED_RT.block_on(async {
1339 let (streams_base, conn) = make_streams_base().await;
1340 let api = ApiClient::new(streams_base.clone());
1341
1342 let id = "test-id-123".to_string();
1343
1344 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
1345
1346 let AggregateTradeStreamParams {
1347 symbol,id,
1348 } = params.clone();
1349
1350 let pairs: &[(&str, Option<String>)] = &[
1351 ("symbol",
1352 Some(symbol.clone())
1353 ),
1354 ("id",
1355 id.clone()
1356 ),
1357 ];
1358
1359 let vars: HashMap<_, _> = pairs
1360 .iter()
1361 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1362 .collect();
1363 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1364
1365 let ws_stream = api.aggregate_trade_stream(params).await.unwrap();
1366
1367 let called = Arc::new(AtomicBool::new(false));
1368 let called_with_message = called.clone();
1369 ws_stream.on_message(move |_payload: models::AggregateTradeStreamResponse| {
1370 called_with_message.store(true, Ordering::SeqCst);
1371 });
1372
1373 let payload: Value = serde_json::from_str(r#"{"e":"aggTrade","E":1771828569861,"T":1771828569702,"a":2684294,"f":2684294,"l":2684294,"m":false,"p":"0.08530000","q":"879.24000000","s":"ALPHA_474USDT"}"#).unwrap_or_else(|_| serde_json::json!({}));
1374 let msg = json!({
1375 "stream": stream,
1376 "data": payload,
1377 });
1378
1379 streams_base.on_message(msg.to_string(), conn.clone()).await;
1380 yield_now().await;
1381
1382 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1383 });
1384 }
1385
1386 #[test]
1387 fn aggregate_trade_stream_should_not_fire_after_unsubscribe() {
1388 TOKIO_SHARED_RT.block_on(async {
1389 let (streams_base, conn) = make_streams_base().await;
1390 let api = ApiClient::new(streams_base.clone());
1391
1392 let id = "test-id-123".to_string();
1393
1394 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
1395
1396 let AggregateTradeStreamParams {
1397 symbol,id,
1398 } = params.clone();
1399
1400 let pairs: &[(&str, Option<String>)] = &[
1401 ("symbol",
1402 Some(symbol.clone())
1403 ),
1404 ("id",
1405 id.clone()
1406 ),
1407 ];
1408
1409 let vars: HashMap<_, _> = pairs
1410 .iter()
1411 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1412 .collect();
1413 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1414
1415 let ws_stream = api.aggregate_trade_stream(params).await.unwrap();
1416
1417 let called = Arc::new(AtomicBool::new(false));
1418 let called_clone = called.clone();
1419 ws_stream.on_message(move |_payload: models::AggregateTradeStreamResponse| {
1420 called_clone.store(true, Ordering::SeqCst);
1421 });
1422
1423 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1424
1425 ws_stream.unsubscribe().await;
1426
1427 let payload: Value = serde_json::from_str(r#"{"e":"aggTrade","E":1771828569861,"T":1771828569702,"a":2684294,"f":2684294,"l":2684294,"m":false,"p":"0.08530000","q":"879.24000000","s":"ALPHA_474USDT"}"#).unwrap_or_else(|_| serde_json::json!({}));
1428 let msg = json!({
1429 "stream": stream,
1430 "data": payload,
1431 });
1432
1433 streams_base.on_message(msg.to_string(), conn.clone()).await;
1434
1435 yield_now().await;
1436
1437 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1438 });
1439 }
1440
1441 #[test]
1442 fn all_book_ticker_stream_should_execute_successfully() {
1443 TOKIO_SHARED_RT.block_on(async {
1444 let (streams_base, _) = make_streams_base().await;
1445 let api = ApiClient::new(streams_base.clone());
1446
1447 let id = "test-id-123".to_string();
1448
1449 let params = AllBookTickerStreamParams::builder()
1450 .id(Some(id.clone()))
1451 .build()
1452 .unwrap();
1453
1454 let AllBookTickerStreamParams { id } = params.clone();
1455
1456 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1457
1458 let vars: HashMap<_, _> = pairs
1459 .iter()
1460 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1461 .collect();
1462 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1463 let ws_stream = api
1464 .all_book_ticker_stream(params)
1465 .await
1466 .expect("all_book_ticker_stream should return a WebsocketStream");
1467
1468 assert!(
1469 streams_base.is_subscribed(&stream).await,
1470 "expected stream '{stream}' to be subscribed"
1471 );
1472 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1473 });
1474 }
1475
1476 #[test]
1477 fn all_book_ticker_stream_should_handle_incoming_message() {
1478 TOKIO_SHARED_RT.block_on(async {
1479 let (streams_base, conn) = make_streams_base().await;
1480 let api = ApiClient::new(streams_base.clone());
1481
1482 let id = "test-id-123".to_string();
1483
1484 let params = AllBookTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1485
1486 let AllBookTickerStreamParams {
1487 id,
1488 } = params.clone();
1489
1490 let pairs: &[(&str, Option<String>)] = &[
1491 ("id",
1492 id.clone()
1493 ),
1494 ];
1495
1496 let vars: HashMap<_, _> = pairs
1497 .iter()
1498 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1499 .collect();
1500 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1501
1502 let ws_stream = api.all_book_ticker_stream(params).await.unwrap();
1503
1504 let called = Arc::new(AtomicBool::new(false));
1505 let called_with_message = called.clone();
1506 ws_stream.on_message(move |_payload: models::AllBookTickerStreamResponse| {
1507 called_with_message.store(true, Ordering::SeqCst);
1508 });
1509
1510 let payload: Value = serde_json::from_str(r#"{"e":"bookTicker","E":1773108379067,"T":1773108379054,"u":6663207693,"s":"ALPHA_116USDT","b":"8.30000000","B":"0.91076900","a":"9.30000000","A":"2.00353200"}"#).unwrap_or_else(|_| serde_json::json!({}));
1511 let msg = json!({
1512 "stream": stream,
1513 "data": payload,
1514 });
1515
1516 streams_base.on_message(msg.to_string(), conn.clone()).await;
1517 yield_now().await;
1518
1519 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1520 });
1521 }
1522
1523 #[test]
1524 fn all_book_ticker_stream_should_not_fire_after_unsubscribe() {
1525 TOKIO_SHARED_RT.block_on(async {
1526 let (streams_base, conn) = make_streams_base().await;
1527 let api = ApiClient::new(streams_base.clone());
1528
1529 let id = "test-id-123".to_string();
1530
1531 let params = AllBookTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1532
1533 let AllBookTickerStreamParams {
1534 id,
1535 } = params.clone();
1536
1537 let pairs: &[(&str, Option<String>)] = &[
1538 ("id",
1539 id.clone()
1540 ),
1541 ];
1542
1543 let vars: HashMap<_, _> = pairs
1544 .iter()
1545 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1546 .collect();
1547 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1548
1549 let ws_stream = api.all_book_ticker_stream(params).await.unwrap();
1550
1551 let called = Arc::new(AtomicBool::new(false));
1552 let called_clone = called.clone();
1553 ws_stream.on_message(move |_payload: models::AllBookTickerStreamResponse| {
1554 called_clone.store(true, Ordering::SeqCst);
1555 });
1556
1557 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1558
1559 ws_stream.unsubscribe().await;
1560
1561 let payload: Value = serde_json::from_str(r#"{"e":"bookTicker","E":1773108379067,"T":1773108379054,"u":6663207693,"s":"ALPHA_116USDT","b":"8.30000000","B":"0.91076900","a":"9.30000000","A":"2.00353200"}"#).unwrap_or_else(|_| serde_json::json!({}));
1562 let msg = json!({
1563 "stream": stream,
1564 "data": payload,
1565 });
1566
1567 streams_base.on_message(msg.to_string(), conn.clone()).await;
1568
1569 yield_now().await;
1570
1571 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1572 });
1573 }
1574
1575 #[test]
1576 fn all_mini_ticker_stream_should_execute_successfully() {
1577 TOKIO_SHARED_RT.block_on(async {
1578 let (streams_base, _) = make_streams_base().await;
1579 let api = ApiClient::new(streams_base.clone());
1580
1581 let id = "test-id-123".to_string();
1582
1583 let params = AllMiniTickerStreamParams::builder()
1584 .id(Some(id.clone()))
1585 .build()
1586 .unwrap();
1587
1588 let AllMiniTickerStreamParams { id } = params.clone();
1589
1590 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1591
1592 let vars: HashMap<_, _> = pairs
1593 .iter()
1594 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1595 .collect();
1596 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1597 let ws_stream = api
1598 .all_mini_ticker_stream(params)
1599 .await
1600 .expect("all_mini_ticker_stream should return a WebsocketStream");
1601
1602 assert!(
1603 streams_base.is_subscribed(&stream).await,
1604 "expected stream '{stream}' to be subscribed"
1605 );
1606 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1607 });
1608 }
1609
1610 #[test]
1611 fn all_mini_ticker_stream_should_handle_incoming_message() {
1612 TOKIO_SHARED_RT.block_on(async {
1613 let (streams_base, conn) = make_streams_base().await;
1614 let api = ApiClient::new(streams_base.clone());
1615
1616 let id = "test-id-123".to_string();
1617
1618 let params = AllMiniTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1619
1620 let AllMiniTickerStreamParams {
1621 id,
1622 } = params.clone();
1623
1624 let pairs: &[(&str, Option<String>)] = &[
1625 ("id",
1626 id.clone()
1627 ),
1628 ];
1629
1630 let vars: HashMap<_, _> = pairs
1631 .iter()
1632 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1633 .collect();
1634 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1635
1636 let ws_stream = api.all_mini_ticker_stream(params).await.unwrap();
1637
1638 let called = Arc::new(AtomicBool::new(false));
1639 let called_with_message = called.clone();
1640 ws_stream.on_message(move |_payload: models::AllMiniTickerStreamResponse| {
1641 called_with_message.store(true, Ordering::SeqCst);
1642 });
1643
1644 let payload: Value = serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1773109449908,"s":"ALPHA_116USDT","c":"8.40000000","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64739.03978900","q":"543810.36005090"}"#).unwrap_or_else(|_| serde_json::json!({}));
1645 let msg = json!({
1646 "stream": stream,
1647 "data": payload,
1648 });
1649
1650 streams_base.on_message(msg.to_string(), conn.clone()).await;
1651 yield_now().await;
1652
1653 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1654 });
1655 }
1656
1657 #[test]
1658 fn all_mini_ticker_stream_should_not_fire_after_unsubscribe() {
1659 TOKIO_SHARED_RT.block_on(async {
1660 let (streams_base, conn) = make_streams_base().await;
1661 let api = ApiClient::new(streams_base.clone());
1662
1663 let id = "test-id-123".to_string();
1664
1665 let params = AllMiniTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1666
1667 let AllMiniTickerStreamParams {
1668 id,
1669 } = params.clone();
1670
1671 let pairs: &[(&str, Option<String>)] = &[
1672 ("id",
1673 id.clone()
1674 ),
1675 ];
1676
1677 let vars: HashMap<_, _> = pairs
1678 .iter()
1679 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1680 .collect();
1681 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1682
1683 let ws_stream = api.all_mini_ticker_stream(params).await.unwrap();
1684
1685 let called = Arc::new(AtomicBool::new(false));
1686 let called_clone = called.clone();
1687 ws_stream.on_message(move |_payload: models::AllMiniTickerStreamResponse| {
1688 called_clone.store(true, Ordering::SeqCst);
1689 });
1690
1691 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1692
1693 ws_stream.unsubscribe().await;
1694
1695 let payload: Value = serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1773109449908,"s":"ALPHA_116USDT","c":"8.40000000","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64739.03978900","q":"543810.36005090"}"#).unwrap_or_else(|_| serde_json::json!({}));
1696 let msg = json!({
1697 "stream": stream,
1698 "data": payload,
1699 });
1700
1701 streams_base.on_message(msg.to_string(), conn.clone()).await;
1702
1703 yield_now().await;
1704
1705 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1706 });
1707 }
1708
1709 #[test]
1710 fn all_ticker_stream_should_execute_successfully() {
1711 TOKIO_SHARED_RT.block_on(async {
1712 let (streams_base, _) = make_streams_base().await;
1713 let api = ApiClient::new(streams_base.clone());
1714
1715 let id = "test-id-123".to_string();
1716
1717 let params = AllTickerStreamParams::builder()
1718 .id(Some(id.clone()))
1719 .build()
1720 .unwrap();
1721
1722 let AllTickerStreamParams { id } = params.clone();
1723
1724 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1725
1726 let vars: HashMap<_, _> = pairs
1727 .iter()
1728 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1729 .collect();
1730 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1731 let ws_stream = api
1732 .all_ticker_stream(params)
1733 .await
1734 .expect("all_ticker_stream should return a WebsocketStream");
1735
1736 assert!(
1737 streams_base.is_subscribed(&stream).await,
1738 "expected stream '{stream}' to be subscribed"
1739 );
1740 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1741 });
1742 }
1743
1744 #[test]
1745 fn all_ticker_stream_should_handle_incoming_message() {
1746 TOKIO_SHARED_RT.block_on(async {
1747 let (streams_base, conn) = make_streams_base().await;
1748 let api = ApiClient::new(streams_base.clone());
1749
1750 let id = "test-id-123".to_string();
1751
1752 let params = AllTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1753
1754 let AllTickerStreamParams {
1755 id,
1756 } = params.clone();
1757
1758 let pairs: &[(&str, Option<String>)] = &[
1759 ("id",
1760 id.clone()
1761 ),
1762 ];
1763
1764 let vars: HashMap<_, _> = pairs
1765 .iter()
1766 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1767 .collect();
1768 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1769
1770 let ws_stream = api.all_ticker_stream(params).await.unwrap();
1771
1772 let called = Arc::new(AtomicBool::new(false));
1773 let called_with_message = called.clone();
1774 ws_stream.on_message(move |_payload: models::AllTickerStreamResponse| {
1775 called_with_message.store(true, Ordering::SeqCst);
1776 });
1777
1778 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1773109631569,"s":"ALPHA_116USDT","p":"0.00000000","P":"0.00","w":"8.40003127","c":"8.40000000","Q":"0.49293200","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64750.63418500","q":"543907.35169250","O":1773023220000,"C":1773109631555,"F":19847634,"L":19911287,"n":217505}"#).unwrap_or_else(|_| serde_json::json!({}));
1779 let msg = json!({
1780 "stream": stream,
1781 "data": payload,
1782 });
1783
1784 streams_base.on_message(msg.to_string(), conn.clone()).await;
1785 yield_now().await;
1786
1787 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1788 });
1789 }
1790
1791 #[test]
1792 fn all_ticker_stream_should_not_fire_after_unsubscribe() {
1793 TOKIO_SHARED_RT.block_on(async {
1794 let (streams_base, conn) = make_streams_base().await;
1795 let api = ApiClient::new(streams_base.clone());
1796
1797 let id = "test-id-123".to_string();
1798
1799 let params = AllTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1800
1801 let AllTickerStreamParams {
1802 id,
1803 } = params.clone();
1804
1805 let pairs: &[(&str, Option<String>)] = &[
1806 ("id",
1807 id.clone()
1808 ),
1809 ];
1810
1811 let vars: HashMap<_, _> = pairs
1812 .iter()
1813 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1814 .collect();
1815 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1816
1817 let ws_stream = api.all_ticker_stream(params).await.unwrap();
1818
1819 let called = Arc::new(AtomicBool::new(false));
1820 let called_clone = called.clone();
1821 ws_stream.on_message(move |_payload: models::AllTickerStreamResponse| {
1822 called_clone.store(true, Ordering::SeqCst);
1823 });
1824
1825 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1826
1827 ws_stream.unsubscribe().await;
1828
1829 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1773109631569,"s":"ALPHA_116USDT","p":"0.00000000","P":"0.00","w":"8.40003127","c":"8.40000000","Q":"0.49293200","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64750.63418500","q":"543907.35169250","O":1773023220000,"C":1773109631555,"F":19847634,"L":19911287,"n":217505}"#).unwrap_or_else(|_| serde_json::json!({}));
1830 let msg = json!({
1831 "stream": stream,
1832 "data": payload,
1833 });
1834
1835 streams_base.on_message(msg.to_string(), conn.clone()).await;
1836
1837 yield_now().await;
1838
1839 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1840 });
1841 }
1842
1843 #[test]
1844 fn all_tokens24h_ticker_stream_should_execute_successfully() {
1845 TOKIO_SHARED_RT.block_on(async {
1846 let (streams_base, _) = make_streams_base().await;
1847 let api = ApiClient::new(streams_base.clone());
1848
1849 let id = "test-id-123".to_string();
1850
1851 let params = AllTokens24hTickerStreamParams::builder()
1852 .id(Some(id.clone()))
1853 .build()
1854 .unwrap();
1855
1856 let AllTokens24hTickerStreamParams { id } = params.clone();
1857
1858 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1859
1860 let vars: HashMap<_, _> = pairs
1861 .iter()
1862 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1863 .collect();
1864 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1865 let ws_stream = api
1866 .all_tokens24h_ticker_stream(params)
1867 .await
1868 .expect("all_tokens24h_ticker_stream should return a WebsocketStream");
1869
1870 assert!(
1871 streams_base.is_subscribed(&stream).await,
1872 "expected stream '{stream}' to be subscribed"
1873 );
1874 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1875 });
1876 }
1877
1878 #[test]
1879 fn all_tokens24h_ticker_stream_should_handle_incoming_message() {
1880 TOKIO_SHARED_RT.block_on(async {
1881 let (streams_base, conn) = make_streams_base().await;
1882 let api = ApiClient::new(streams_base.clone());
1883
1884 let id = "test-id-123".to_string();
1885
1886 let params = AllTokens24hTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1887
1888 let AllTokens24hTickerStreamParams {
1889 id,
1890 } = params.clone();
1891
1892 let pairs: &[(&str, Option<String>)] = &[
1893 ("id",
1894 id.clone()
1895 ),
1896 ];
1897
1898 let vars: HashMap<_, _> = pairs
1899 .iter()
1900 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1901 .collect();
1902 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1903
1904 let ws_stream = api.all_tokens24h_ticker_stream(params).await.unwrap();
1905
1906 let called = Arc::new(AtomicBool::new(false));
1907 let called_with_message = called.clone();
1908 ws_stream.on_message(move |_payload: models::AllTokens24hTickerStreamResponse| {
1909 called_with_message.store(true, Ordering::SeqCst);
1910 });
1911
1912 let payload: Value = serde_json::from_str(r#"{"e":"tickerList","d":[{"ca":"0x8fce...@56","cnt24":10426,"fdv":"72834145.18711072","hc":"6822","liq":"1442996.98151177662835","mc":"13542034.24008014","p":"0.072849536538593663","pc24":"-2.32","s":"1","t":1771825733000,"vol24":"12930712.213702157742688594285"}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1913 let msg = json!({
1914 "stream": stream,
1915 "data": payload,
1916 });
1917
1918 streams_base.on_message(msg.to_string(), conn.clone()).await;
1919 yield_now().await;
1920
1921 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1922 });
1923 }
1924
1925 #[test]
1926 fn all_tokens24h_ticker_stream_should_not_fire_after_unsubscribe() {
1927 TOKIO_SHARED_RT.block_on(async {
1928 let (streams_base, conn) = make_streams_base().await;
1929 let api = ApiClient::new(streams_base.clone());
1930
1931 let id = "test-id-123".to_string();
1932
1933 let params = AllTokens24hTickerStreamParams::builder().id(Some(id.clone())).build().unwrap();
1934
1935 let AllTokens24hTickerStreamParams {
1936 id,
1937 } = params.clone();
1938
1939 let pairs: &[(&str, Option<String>)] = &[
1940 ("id",
1941 id.clone()
1942 ),
1943 ];
1944
1945 let vars: HashMap<_, _> = pairs
1946 .iter()
1947 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1948 .collect();
1949 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1950
1951 let ws_stream = api.all_tokens24h_ticker_stream(params).await.unwrap();
1952
1953 let called = Arc::new(AtomicBool::new(false));
1954 let called_clone = called.clone();
1955 ws_stream.on_message(move |_payload: models::AllTokens24hTickerStreamResponse| {
1956 called_clone.store(true, Ordering::SeqCst);
1957 });
1958
1959 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1960
1961 ws_stream.unsubscribe().await;
1962
1963 let payload: Value = serde_json::from_str(r#"{"e":"tickerList","d":[{"ca":"0x8fce...@56","cnt24":10426,"fdv":"72834145.18711072","hc":"6822","liq":"1442996.98151177662835","mc":"13542034.24008014","p":"0.072849536538593663","pc24":"-2.32","s":"1","t":1771825733000,"vol24":"12930712.213702157742688594285"}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1964 let msg = json!({
1965 "stream": stream,
1966 "data": payload,
1967 });
1968
1969 streams_base.on_message(msg.to_string(), conn.clone()).await;
1970
1971 yield_now().await;
1972
1973 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1974 });
1975 }
1976
1977 #[test]
1978 fn book_ticker_stream_should_execute_successfully() {
1979 TOKIO_SHARED_RT.block_on(async {
1980 let (streams_base, _) = make_streams_base().await;
1981 let api = ApiClient::new(streams_base.clone());
1982
1983 let id = "test-id-123".to_string();
1984
1985 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string())
1986 .id(Some(id.clone()))
1987 .build()
1988 .unwrap();
1989
1990 let BookTickerStreamParams { symbol, id } = params.clone();
1991
1992 let pairs: &[(&str, Option<String>)] =
1993 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1994
1995 let vars: HashMap<_, _> = pairs
1996 .iter()
1997 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1998 .collect();
1999 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2000 let ws_stream = api
2001 .book_ticker_stream(params)
2002 .await
2003 .expect("book_ticker_stream should return a WebsocketStream");
2004
2005 assert!(
2006 streams_base.is_subscribed(&stream).await,
2007 "expected stream '{stream}' to be subscribed"
2008 );
2009 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2010 });
2011 }
2012
2013 #[test]
2014 fn book_ticker_stream_should_handle_incoming_message() {
2015 TOKIO_SHARED_RT.block_on(async {
2016 let (streams_base, conn) = make_streams_base().await;
2017 let api = ApiClient::new(streams_base.clone());
2018
2019 let id = "test-id-123".to_string();
2020
2021 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
2022
2023 let BookTickerStreamParams {
2024 symbol,id,
2025 } = params.clone();
2026
2027 let pairs: &[(&str, Option<String>)] = &[
2028 ("symbol",
2029 Some(symbol.clone())
2030 ),
2031 ("id",
2032 id.clone()
2033 ),
2034 ];
2035
2036 let vars: HashMap<_, _> = pairs
2037 .iter()
2038 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2039 .collect();
2040 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2041
2042 let ws_stream = api.book_ticker_stream(params).await.unwrap();
2043
2044 let called = Arc::new(AtomicBool::new(false));
2045 let called_with_message = called.clone();
2046 ws_stream.on_message(move |_payload: models::BookTickerStreamResponse| {
2047 called_with_message.store(true, Ordering::SeqCst);
2048 });
2049
2050 let payload: Value = serde_json::from_str(r#"{"e":"bookTicker","E":1773108379067,"T":1773108379054,"u":6663207693,"s":"ALPHA_116USDT","b":"8.30000000","B":"0.91076900","a":"9.30000000","A":"2.00353200"}"#).unwrap_or_else(|_| serde_json::json!({}));
2051 let msg = json!({
2052 "stream": stream,
2053 "data": payload,
2054 });
2055
2056 streams_base.on_message(msg.to_string(), conn.clone()).await;
2057 yield_now().await;
2058
2059 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2060 });
2061 }
2062
2063 #[test]
2064 fn book_ticker_stream_should_not_fire_after_unsubscribe() {
2065 TOKIO_SHARED_RT.block_on(async {
2066 let (streams_base, conn) = make_streams_base().await;
2067 let api = ApiClient::new(streams_base.clone());
2068
2069 let id = "test-id-123".to_string();
2070
2071 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
2072
2073 let BookTickerStreamParams {
2074 symbol,id,
2075 } = params.clone();
2076
2077 let pairs: &[(&str, Option<String>)] = &[
2078 ("symbol",
2079 Some(symbol.clone())
2080 ),
2081 ("id",
2082 id.clone()
2083 ),
2084 ];
2085
2086 let vars: HashMap<_, _> = pairs
2087 .iter()
2088 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2089 .collect();
2090 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2091
2092 let ws_stream = api.book_ticker_stream(params).await.unwrap();
2093
2094 let called = Arc::new(AtomicBool::new(false));
2095 let called_clone = called.clone();
2096 ws_stream.on_message(move |_payload: models::BookTickerStreamResponse| {
2097 called_clone.store(true, Ordering::SeqCst);
2098 });
2099
2100 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2101
2102 ws_stream.unsubscribe().await;
2103
2104 let payload: Value = serde_json::from_str(r#"{"e":"bookTicker","E":1773108379067,"T":1773108379054,"u":6663207693,"s":"ALPHA_116USDT","b":"8.30000000","B":"0.91076900","a":"9.30000000","A":"2.00353200"}"#).unwrap_or_else(|_| serde_json::json!({}));
2105 let msg = json!({
2106 "stream": stream,
2107 "data": payload,
2108 });
2109
2110 streams_base.on_message(msg.to_string(), conn.clone()).await;
2111
2112 yield_now().await;
2113
2114 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2115 });
2116 }
2117
2118 #[test]
2119 fn contract_kline_stream_should_execute_successfully() {
2120 TOKIO_SHARED_RT.block_on(async {
2121 let (streams_base, _) = make_streams_base().await;
2122 let api = ApiClient::new(streams_base.clone());
2123
2124 let id = "test-id-123".to_string();
2125
2126 let params = ContractKlineStreamParams::builder(
2127 "G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),
2128 "CT_501".to_string(),
2129 ContractKlineStreamIntervalEnum::Interval1s,
2130 )
2131 .id(Some(id.clone()))
2132 .build()
2133 .unwrap();
2134
2135 let ContractKlineStreamParams {
2136 contract_address,
2137 chain_id,
2138 interval,
2139 id,
2140 } = params.clone();
2141
2142 let pairs: &[(&str, Option<String>)] = &[
2143 ("contractAddress", Some(contract_address.clone())),
2144 ("chainId", Some(chain_id.clone())),
2145 ("interval", Some(interval.as_str().to_string())),
2146 ("id", id.clone()),
2147 ];
2148
2149 let vars: HashMap<_, _> = pairs
2150 .iter()
2151 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2152 .collect();
2153 let stream = replace_websocket_streams_placeholders(
2154 "/came@<contractAddress>@<chainId>@kline_<interval>",
2155 &vars,
2156 );
2157 let ws_stream = api
2158 .contract_kline_stream(params)
2159 .await
2160 .expect("contract_kline_stream should return a WebsocketStream");
2161
2162 assert!(
2163 streams_base.is_subscribed(&stream).await,
2164 "expected stream '{stream}' to be subscribed"
2165 );
2166 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2167 });
2168 }
2169
2170 #[test]
2171 fn contract_kline_stream_should_handle_incoming_message() {
2172 TOKIO_SHARED_RT.block_on(async {
2173 let (streams_base, conn) = make_streams_base().await;
2174 let api = ApiClient::new(streams_base.clone());
2175
2176 let id = "test-id-123".to_string();
2177
2178 let params = ContractKlineStreamParams::builder("G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),"CT_501".to_string(),ContractKlineStreamIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2179
2180 let ContractKlineStreamParams {
2181 contract_address,chain_id,interval,id,
2182 } = params.clone();
2183
2184 let pairs: &[(&str, Option<String>)] = &[
2185 ("contractAddress",
2186 Some(contract_address.clone())
2187 ),
2188 ("chainId",
2189 Some(chain_id.clone())
2190 ),
2191 ("interval",
2192 Some(interval.as_str().to_string())
2193 ),
2194 ("id",
2195 id.clone()
2196 ),
2197 ];
2198
2199 let vars: HashMap<_, _> = pairs
2200 .iter()
2201 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2202 .collect();
2203 let stream = replace_websocket_streams_placeholders("/came@<contractAddress>@<chainId>@kline_<interval>", &vars);
2204
2205 let ws_stream = api.contract_kline_stream(params).await.unwrap();
2206
2207 let called = Arc::new(AtomicBool::new(false));
2208 let called_with_message = called.clone();
2209 ws_stream.on_message(move |_payload: models::ContractKlineStreamResponse| {
2210 called_with_message.store(true, Ordering::SeqCst);
2211 });
2212
2213 let payload: Value = serde_json::from_str(r#"{"ca":"G7vQW...@CT_501","e":"kline","k":{"o":"0.15584372258086979622","c":"0.15584957424118055058","h":"0.15584957424118055058","l":"0.15584372258086979622","v":"2.8584092999332775480","ot":1771828655000,"ct":1771828656000,"i":"1s"}}"#).unwrap_or_else(|_| serde_json::json!({}));
2214 let msg = json!({
2215 "stream": stream,
2216 "data": payload,
2217 });
2218
2219 streams_base.on_message(msg.to_string(), conn.clone()).await;
2220 yield_now().await;
2221
2222 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2223 });
2224 }
2225
2226 #[test]
2227 fn contract_kline_stream_should_not_fire_after_unsubscribe() {
2228 TOKIO_SHARED_RT.block_on(async {
2229 let (streams_base, conn) = make_streams_base().await;
2230 let api = ApiClient::new(streams_base.clone());
2231
2232 let id = "test-id-123".to_string();
2233
2234 let params = ContractKlineStreamParams::builder("G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),"CT_501".to_string(),ContractKlineStreamIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2235
2236 let ContractKlineStreamParams {
2237 contract_address,chain_id,interval,id,
2238 } = params.clone();
2239
2240 let pairs: &[(&str, Option<String>)] = &[
2241 ("contractAddress",
2242 Some(contract_address.clone())
2243 ),
2244 ("chainId",
2245 Some(chain_id.clone())
2246 ),
2247 ("interval",
2248 Some(interval.as_str().to_string())
2249 ),
2250 ("id",
2251 id.clone()
2252 ),
2253 ];
2254
2255 let vars: HashMap<_, _> = pairs
2256 .iter()
2257 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2258 .collect();
2259 let stream = replace_websocket_streams_placeholders("/came@<contractAddress>@<chainId>@kline_<interval>", &vars);
2260
2261 let ws_stream = api.contract_kline_stream(params).await.unwrap();
2262
2263 let called = Arc::new(AtomicBool::new(false));
2264 let called_clone = called.clone();
2265 ws_stream.on_message(move |_payload: models::ContractKlineStreamResponse| {
2266 called_clone.store(true, Ordering::SeqCst);
2267 });
2268
2269 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2270
2271 ws_stream.unsubscribe().await;
2272
2273 let payload: Value = serde_json::from_str(r#"{"ca":"G7vQW...@CT_501","e":"kline","k":{"o":"0.15584372258086979622","c":"0.15584957424118055058","h":"0.15584957424118055058","l":"0.15584372258086979622","v":"2.8584092999332775480","ot":1771828655000,"ct":1771828656000,"i":"1s"}}"#).unwrap_or_else(|_| serde_json::json!({}));
2274 let msg = json!({
2275 "stream": stream,
2276 "data": payload,
2277 });
2278
2279 streams_base.on_message(msg.to_string(), conn.clone()).await;
2280
2281 yield_now().await;
2282
2283 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2284 });
2285 }
2286
2287 #[test]
2288 fn full_depth_stream_should_execute_successfully() {
2289 TOKIO_SHARED_RT.block_on(async {
2290 let (streams_base, _) = make_streams_base().await;
2291 let api = ApiClient::new(streams_base.clone());
2292
2293 let id = "test-id-123".to_string();
2294
2295 let params = FullDepthStreamParams::builder(
2296 "alpha_116usdt".to_string(),
2297 FullDepthStreamIntervalEnum::Interval0ms,
2298 )
2299 .id(Some(id.clone()))
2300 .build()
2301 .unwrap();
2302
2303 let FullDepthStreamParams {
2304 symbol,
2305 interval,
2306 id,
2307 } = params.clone();
2308
2309 let pairs: &[(&str, Option<String>)] = &[
2310 ("symbol", Some(symbol.clone())),
2311 ("interval", Some(interval.as_str().to_string())),
2312 ("id", id.clone()),
2313 ];
2314
2315 let vars: HashMap<_, _> = pairs
2316 .iter()
2317 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2318 .collect();
2319 let stream =
2320 replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2321 let ws_stream = api
2322 .full_depth_stream(params)
2323 .await
2324 .expect("full_depth_stream should return a WebsocketStream");
2325
2326 assert!(
2327 streams_base.is_subscribed(&stream).await,
2328 "expected stream '{stream}' to be subscribed"
2329 );
2330 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2331 });
2332 }
2333
2334 #[test]
2335 fn full_depth_stream_should_handle_incoming_message() {
2336 TOKIO_SHARED_RT.block_on(async {
2337 let (streams_base, conn) = make_streams_base().await;
2338 let api = ApiClient::new(streams_base.clone());
2339
2340 let id = "test-id-123".to_string();
2341
2342 let params = FullDepthStreamParams::builder("alpha_116usdt".to_string(),FullDepthStreamIntervalEnum::Interval0ms,).id(Some(id.clone())).build().unwrap();
2343
2344 let FullDepthStreamParams {
2345 symbol,interval,id,
2346 } = params.clone();
2347
2348 let pairs: &[(&str, Option<String>)] = &[
2349 ("symbol",
2350 Some(symbol.clone())
2351 ),
2352 ("interval",
2353 Some(interval.as_str().to_string())
2354 ),
2355 ("id",
2356 id.clone()
2357 ),
2358 ];
2359
2360 let vars: HashMap<_, _> = pairs
2361 .iter()
2362 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2363 .collect();
2364 let stream = replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2365
2366 let ws_stream = api.full_depth_stream(params).await.unwrap();
2367
2368 let called = Arc::new(AtomicBool::new(false));
2369 let called_with_message = called.clone();
2370 ws_stream.on_message(move |_payload: models::FullDepthStreamResponse| {
2371 called_with_message.store(true, Ordering::SeqCst);
2372 });
2373
2374 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1771828614381,"T":1771828614258,"U":42648947917,"u":42648948229,"pu":42648947321,"s":"ALPHA_474USDT","b":[["0.16479999","0"],["0.16480000","18.68000000"]],"a":[["0.16490000","500.00000000"],["0.16500000","0"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2375 let msg = json!({
2376 "stream": stream,
2377 "data": payload,
2378 });
2379
2380 streams_base.on_message(msg.to_string(), conn.clone()).await;
2381 yield_now().await;
2382
2383 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2384 });
2385 }
2386
2387 #[test]
2388 fn full_depth_stream_should_not_fire_after_unsubscribe() {
2389 TOKIO_SHARED_RT.block_on(async {
2390 let (streams_base, conn) = make_streams_base().await;
2391 let api = ApiClient::new(streams_base.clone());
2392
2393 let id = "test-id-123".to_string();
2394
2395 let params = FullDepthStreamParams::builder("alpha_116usdt".to_string(),FullDepthStreamIntervalEnum::Interval0ms,).id(Some(id.clone())).build().unwrap();
2396
2397 let FullDepthStreamParams {
2398 symbol,interval,id,
2399 } = params.clone();
2400
2401 let pairs: &[(&str, Option<String>)] = &[
2402 ("symbol",
2403 Some(symbol.clone())
2404 ),
2405 ("interval",
2406 Some(interval.as_str().to_string())
2407 ),
2408 ("id",
2409 id.clone()
2410 ),
2411 ];
2412
2413 let vars: HashMap<_, _> = pairs
2414 .iter()
2415 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2416 .collect();
2417 let stream = replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2418
2419 let ws_stream = api.full_depth_stream(params).await.unwrap();
2420
2421 let called = Arc::new(AtomicBool::new(false));
2422 let called_clone = called.clone();
2423 ws_stream.on_message(move |_payload: models::FullDepthStreamResponse| {
2424 called_clone.store(true, Ordering::SeqCst);
2425 });
2426
2427 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2428
2429 ws_stream.unsubscribe().await;
2430
2431 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1771828614381,"T":1771828614258,"U":42648947917,"u":42648948229,"pu":42648947321,"s":"ALPHA_474USDT","b":[["0.16479999","0"],["0.16480000","18.68000000"]],"a":[["0.16490000","500.00000000"],["0.16500000","0"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2432 let msg = json!({
2433 "stream": stream,
2434 "data": payload,
2435 });
2436
2437 streams_base.on_message(msg.to_string(), conn.clone()).await;
2438
2439 yield_now().await;
2440
2441 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2442 });
2443 }
2444
2445 #[test]
2446 fn kline_stream_should_execute_successfully() {
2447 TOKIO_SHARED_RT.block_on(async {
2448 let (streams_base, _) = make_streams_base().await;
2449 let api = ApiClient::new(streams_base.clone());
2450
2451 let id = "test-id-123".to_string();
2452
2453 let params = KlineStreamParams::builder(
2454 "alpha_116usdt".to_string(),
2455 KlineStreamIntervalEnum::Interval1m,
2456 )
2457 .id(Some(id.clone()))
2458 .build()
2459 .unwrap();
2460
2461 let KlineStreamParams {
2462 symbol,
2463 interval,
2464 id,
2465 } = params.clone();
2466
2467 let pairs: &[(&str, Option<String>)] = &[
2468 ("symbol", Some(symbol.clone())),
2469 ("interval", Some(interval.as_str().to_string())),
2470 ("id", id.clone()),
2471 ];
2472
2473 let vars: HashMap<_, _> = pairs
2474 .iter()
2475 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2476 .collect();
2477 let stream =
2478 replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2479 let ws_stream = api
2480 .kline_stream(params)
2481 .await
2482 .expect("kline_stream should return a WebsocketStream");
2483
2484 assert!(
2485 streams_base.is_subscribed(&stream).await,
2486 "expected stream '{stream}' to be subscribed"
2487 );
2488 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2489 });
2490 }
2491
2492 #[test]
2493 fn kline_stream_should_handle_incoming_message() {
2494 TOKIO_SHARED_RT.block_on(async {
2495 let (streams_base, conn) = make_streams_base().await;
2496 let api = ApiClient::new(streams_base.clone());
2497
2498 let id = "test-id-123".to_string();
2499
2500 let params = KlineStreamParams::builder("alpha_116usdt".to_string(),KlineStreamIntervalEnum::Interval1m,).id(Some(id.clone())).build().unwrap();
2501
2502 let KlineStreamParams {
2503 symbol,interval,id,
2504 } = params.clone();
2505
2506 let pairs: &[(&str, Option<String>)] = &[
2507 ("symbol",
2508 Some(symbol.clone())
2509 ),
2510 ("interval",
2511 Some(interval.as_str().to_string())
2512 ),
2513 ("id",
2514 id.clone()
2515 ),
2516 ];
2517
2518 let vars: HashMap<_, _> = pairs
2519 .iter()
2520 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2521 .collect();
2522 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2523
2524 let ws_stream = api.kline_stream(params).await.unwrap();
2525
2526 let called = Arc::new(AtomicBool::new(false));
2527 let called_with_message = called.clone();
2528 ws_stream.on_message(move |_payload: models::KlineStreamResponse| {
2529 called_with_message.store(true, Ordering::SeqCst);
2530 });
2531
2532 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1773111055144,"s":"ALPHA_116USDT","k":{"t":1773111000000,"T":1773111059999,"s":"ALPHA_116USDT","i":"1m","f":19912696,"L":19912765,"o":"8.40000000","c":"8.40000000","h":"8.40000000","l":"8.40000000","v":"63.64876300","n":218,"x":false,"q":"534.64960920","V":"63.64876300","Q":"534.64960920","B":"0"}}"#).unwrap_or_else(|_| serde_json::json!({}));
2533 let msg = json!({
2534 "stream": stream,
2535 "data": payload,
2536 });
2537
2538 streams_base.on_message(msg.to_string(), conn.clone()).await;
2539 yield_now().await;
2540
2541 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2542 });
2543 }
2544
2545 #[test]
2546 fn kline_stream_should_not_fire_after_unsubscribe() {
2547 TOKIO_SHARED_RT.block_on(async {
2548 let (streams_base, conn) = make_streams_base().await;
2549 let api = ApiClient::new(streams_base.clone());
2550
2551 let id = "test-id-123".to_string();
2552
2553 let params = KlineStreamParams::builder("alpha_116usdt".to_string(),KlineStreamIntervalEnum::Interval1m,).id(Some(id.clone())).build().unwrap();
2554
2555 let KlineStreamParams {
2556 symbol,interval,id,
2557 } = params.clone();
2558
2559 let pairs: &[(&str, Option<String>)] = &[
2560 ("symbol",
2561 Some(symbol.clone())
2562 ),
2563 ("interval",
2564 Some(interval.as_str().to_string())
2565 ),
2566 ("id",
2567 id.clone()
2568 ),
2569 ];
2570
2571 let vars: HashMap<_, _> = pairs
2572 .iter()
2573 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2574 .collect();
2575 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2576
2577 let ws_stream = api.kline_stream(params).await.unwrap();
2578
2579 let called = Arc::new(AtomicBool::new(false));
2580 let called_clone = called.clone();
2581 ws_stream.on_message(move |_payload: models::KlineStreamResponse| {
2582 called_clone.store(true, Ordering::SeqCst);
2583 });
2584
2585 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2586
2587 ws_stream.unsubscribe().await;
2588
2589 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1773111055144,"s":"ALPHA_116USDT","k":{"t":1773111000000,"T":1773111059999,"s":"ALPHA_116USDT","i":"1m","f":19912696,"L":19912765,"o":"8.40000000","c":"8.40000000","h":"8.40000000","l":"8.40000000","v":"63.64876300","n":218,"x":false,"q":"534.64960920","V":"63.64876300","Q":"534.64960920","B":"0"}}"#).unwrap_or_else(|_| serde_json::json!({}));
2590 let msg = json!({
2591 "stream": stream,
2592 "data": payload,
2593 });
2594
2595 streams_base.on_message(msg.to_string(), conn.clone()).await;
2596
2597 yield_now().await;
2598
2599 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2600 });
2601 }
2602
2603 #[test]
2604 fn mini_ticker_stream_should_execute_successfully() {
2605 TOKIO_SHARED_RT.block_on(async {
2606 let (streams_base, _) = make_streams_base().await;
2607 let api = ApiClient::new(streams_base.clone());
2608
2609 let id = "test-id-123".to_string();
2610
2611 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string())
2612 .id(Some(id.clone()))
2613 .build()
2614 .unwrap();
2615
2616 let MiniTickerStreamParams { symbol, id } = params.clone();
2617
2618 let pairs: &[(&str, Option<String>)] =
2619 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2620
2621 let vars: HashMap<_, _> = pairs
2622 .iter()
2623 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2624 .collect();
2625 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2626 let ws_stream = api
2627 .mini_ticker_stream(params)
2628 .await
2629 .expect("mini_ticker_stream should return a WebsocketStream");
2630
2631 assert!(
2632 streams_base.is_subscribed(&stream).await,
2633 "expected stream '{stream}' to be subscribed"
2634 );
2635 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2636 });
2637 }
2638
2639 #[test]
2640 fn mini_ticker_stream_should_handle_incoming_message() {
2641 TOKIO_SHARED_RT.block_on(async {
2642 let (streams_base, conn) = make_streams_base().await;
2643 let api = ApiClient::new(streams_base.clone());
2644
2645 let id = "test-id-123".to_string();
2646
2647 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
2648
2649 let MiniTickerStreamParams {
2650 symbol,id,
2651 } = params.clone();
2652
2653 let pairs: &[(&str, Option<String>)] = &[
2654 ("symbol",
2655 Some(symbol.clone())
2656 ),
2657 ("id",
2658 id.clone()
2659 ),
2660 ];
2661
2662 let vars: HashMap<_, _> = pairs
2663 .iter()
2664 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2665 .collect();
2666 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2667
2668 let ws_stream = api.mini_ticker_stream(params).await.unwrap();
2669
2670 let called = Arc::new(AtomicBool::new(false));
2671 let called_with_message = called.clone();
2672 ws_stream.on_message(move |_payload: models::MiniTickerStreamResponse| {
2673 called_with_message.store(true, Ordering::SeqCst);
2674 });
2675
2676 let payload: Value = serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1773109449908,"s":"ALPHA_116USDT","c":"8.40000000","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64739.03978900","q":"543810.36005090"}"#).unwrap_or_else(|_| serde_json::json!({}));
2677 let msg = json!({
2678 "stream": stream,
2679 "data": payload,
2680 });
2681
2682 streams_base.on_message(msg.to_string(), conn.clone()).await;
2683 yield_now().await;
2684
2685 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2686 });
2687 }
2688
2689 #[test]
2690 fn mini_ticker_stream_should_not_fire_after_unsubscribe() {
2691 TOKIO_SHARED_RT.block_on(async {
2692 let (streams_base, conn) = make_streams_base().await;
2693 let api = ApiClient::new(streams_base.clone());
2694
2695 let id = "test-id-123".to_string();
2696
2697 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
2698
2699 let MiniTickerStreamParams {
2700 symbol,id,
2701 } = params.clone();
2702
2703 let pairs: &[(&str, Option<String>)] = &[
2704 ("symbol",
2705 Some(symbol.clone())
2706 ),
2707 ("id",
2708 id.clone()
2709 ),
2710 ];
2711
2712 let vars: HashMap<_, _> = pairs
2713 .iter()
2714 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2715 .collect();
2716 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2717
2718 let ws_stream = api.mini_ticker_stream(params).await.unwrap();
2719
2720 let called = Arc::new(AtomicBool::new(false));
2721 let called_clone = called.clone();
2722 ws_stream.on_message(move |_payload: models::MiniTickerStreamResponse| {
2723 called_clone.store(true, Ordering::SeqCst);
2724 });
2725
2726 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2727
2728 ws_stream.unsubscribe().await;
2729
2730 let payload: Value = serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1773109449908,"s":"ALPHA_116USDT","c":"8.40000000","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64739.03978900","q":"543810.36005090"}"#).unwrap_or_else(|_| serde_json::json!({}));
2731 let msg = json!({
2732 "stream": stream,
2733 "data": payload,
2734 });
2735
2736 streams_base.on_message(msg.to_string(), conn.clone()).await;
2737
2738 yield_now().await;
2739
2740 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2741 });
2742 }
2743
2744 #[test]
2745 fn partial_depth_stream_should_execute_successfully() {
2746 TOKIO_SHARED_RT.block_on(async {
2747 let (streams_base, _) = make_streams_base().await;
2748 let api = ApiClient::new(streams_base.clone());
2749
2750 let id = "test-id-123".to_string();
2751
2752 let params = PartialDepthStreamParams::builder(
2753 "alpha_116usdt".to_string(),
2754 PartialDepthStreamLevelsEnum::Levels5,
2755 PartialDepthStreamIntervalEnum::Interval0ms,
2756 )
2757 .id(Some(id.clone()))
2758 .build()
2759 .unwrap();
2760
2761 let PartialDepthStreamParams {
2762 symbol,
2763 levels,
2764 interval,
2765 id,
2766 } = params.clone();
2767
2768 let pairs: &[(&str, Option<String>)] = &[
2769 ("symbol", Some(symbol.clone())),
2770 ("levels", Some(levels.as_str().to_string())),
2771 ("interval", Some(interval.as_str().to_string())),
2772 ("id", id.clone()),
2773 ];
2774
2775 let vars: HashMap<_, _> = pairs
2776 .iter()
2777 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2778 .collect();
2779 let stream =
2780 replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2781 let ws_stream = api
2782 .partial_depth_stream(params)
2783 .await
2784 .expect("partial_depth_stream should return a WebsocketStream");
2785
2786 assert!(
2787 streams_base.is_subscribed(&stream).await,
2788 "expected stream '{stream}' to be subscribed"
2789 );
2790 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2791 });
2792 }
2793
2794 #[test]
2795 fn partial_depth_stream_should_handle_incoming_message() {
2796 TOKIO_SHARED_RT.block_on(async {
2797 let (streams_base, conn) = make_streams_base().await;
2798 let api = ApiClient::new(streams_base.clone());
2799
2800 let id = "test-id-123".to_string();
2801
2802 let params = PartialDepthStreamParams::builder("alpha_116usdt".to_string(),PartialDepthStreamLevelsEnum::Levels5,PartialDepthStreamIntervalEnum::Interval0ms,).id(Some(id.clone())).build().unwrap();
2803
2804 let PartialDepthStreamParams {
2805 symbol,levels,interval,id,
2806 } = params.clone();
2807
2808 let pairs: &[(&str, Option<String>)] = &[
2809 ("symbol",
2810 Some(symbol.clone())
2811 ),
2812 ("levels",
2813 Some(levels.as_str().to_string())
2814 ),
2815 ("interval",
2816 Some(interval.as_str().to_string())
2817 ),
2818 ("id",
2819 id.clone()
2820 ),
2821 ];
2822
2823 let vars: HashMap<_, _> = pairs
2824 .iter()
2825 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2826 .collect();
2827 let stream = replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2828
2829 let ws_stream = api.partial_depth_stream(params).await.unwrap();
2830
2831 let called = Arc::new(AtomicBool::new(false));
2832 let called_with_message = called.clone();
2833 ws_stream.on_message(move |_payload: models::PartialDepthStreamResponse| {
2834 called_with_message.store(true, Ordering::SeqCst);
2835 });
2836
2837 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1773110222878,"T":1773110222767,"U":6663367135,"u":6663367147,"pu":6663366971,"s":"ALPHA_116USDT","b":[["0.16479999","0"],["0.16480000","18.68000000"]],"a":[["0.16490000","500.00000000"],["0.16500000","0"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2838 let msg = json!({
2839 "stream": stream,
2840 "data": payload,
2841 });
2842
2843 streams_base.on_message(msg.to_string(), conn.clone()).await;
2844 yield_now().await;
2845
2846 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2847 });
2848 }
2849
2850 #[test]
2851 fn partial_depth_stream_should_not_fire_after_unsubscribe() {
2852 TOKIO_SHARED_RT.block_on(async {
2853 let (streams_base, conn) = make_streams_base().await;
2854 let api = ApiClient::new(streams_base.clone());
2855
2856 let id = "test-id-123".to_string();
2857
2858 let params = PartialDepthStreamParams::builder("alpha_116usdt".to_string(),PartialDepthStreamLevelsEnum::Levels5,PartialDepthStreamIntervalEnum::Interval0ms,).id(Some(id.clone())).build().unwrap();
2859
2860 let PartialDepthStreamParams {
2861 symbol,levels,interval,id,
2862 } = params.clone();
2863
2864 let pairs: &[(&str, Option<String>)] = &[
2865 ("symbol",
2866 Some(symbol.clone())
2867 ),
2868 ("levels",
2869 Some(levels.as_str().to_string())
2870 ),
2871 ("interval",
2872 Some(interval.as_str().to_string())
2873 ),
2874 ("id",
2875 id.clone()
2876 ),
2877 ];
2878
2879 let vars: HashMap<_, _> = pairs
2880 .iter()
2881 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2882 .collect();
2883 let stream = replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2884
2885 let ws_stream = api.partial_depth_stream(params).await.unwrap();
2886
2887 let called = Arc::new(AtomicBool::new(false));
2888 let called_clone = called.clone();
2889 ws_stream.on_message(move |_payload: models::PartialDepthStreamResponse| {
2890 called_clone.store(true, Ordering::SeqCst);
2891 });
2892
2893 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2894
2895 ws_stream.unsubscribe().await;
2896
2897 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1773110222878,"T":1773110222767,"U":6663367135,"u":6663367147,"pu":6663366971,"s":"ALPHA_116USDT","b":[["0.16479999","0"],["0.16480000","18.68000000"]],"a":[["0.16490000","500.00000000"],["0.16500000","0"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2898 let msg = json!({
2899 "stream": stream,
2900 "data": payload,
2901 });
2902
2903 streams_base.on_message(msg.to_string(), conn.clone()).await;
2904
2905 yield_now().await;
2906
2907 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2908 });
2909 }
2910
2911 #[test]
2912 fn ticker_stream_should_execute_successfully() {
2913 TOKIO_SHARED_RT.block_on(async {
2914 let (streams_base, _) = make_streams_base().await;
2915 let api = ApiClient::new(streams_base.clone());
2916
2917 let id = "test-id-123".to_string();
2918
2919 let params = TickerStreamParams::builder("alpha_116usdt".to_string())
2920 .id(Some(id.clone()))
2921 .build()
2922 .unwrap();
2923
2924 let TickerStreamParams { symbol, id } = params.clone();
2925
2926 let pairs: &[(&str, Option<String>)] =
2927 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2928
2929 let vars: HashMap<_, _> = pairs
2930 .iter()
2931 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2932 .collect();
2933 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
2934 let ws_stream = api
2935 .ticker_stream(params)
2936 .await
2937 .expect("ticker_stream should return a WebsocketStream");
2938
2939 assert!(
2940 streams_base.is_subscribed(&stream).await,
2941 "expected stream '{stream}' to be subscribed"
2942 );
2943 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2944 });
2945 }
2946
2947 #[test]
2948 fn ticker_stream_should_handle_incoming_message() {
2949 TOKIO_SHARED_RT.block_on(async {
2950 let (streams_base, conn) = make_streams_base().await;
2951 let api = ApiClient::new(streams_base.clone());
2952
2953 let id = "test-id-123".to_string();
2954
2955 let params = TickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
2956
2957 let TickerStreamParams {
2958 symbol,id,
2959 } = params.clone();
2960
2961 let pairs: &[(&str, Option<String>)] = &[
2962 ("symbol",
2963 Some(symbol.clone())
2964 ),
2965 ("id",
2966 id.clone()
2967 ),
2968 ];
2969
2970 let vars: HashMap<_, _> = pairs
2971 .iter()
2972 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2973 .collect();
2974 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
2975
2976 let ws_stream = api.ticker_stream(params).await.unwrap();
2977
2978 let called = Arc::new(AtomicBool::new(false));
2979 let called_with_message = called.clone();
2980 ws_stream.on_message(move |_payload: models::TickerStreamResponse| {
2981 called_with_message.store(true, Ordering::SeqCst);
2982 });
2983
2984 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1773109631569,"s":"ALPHA_116USDT","p":"0.00000000","P":"0.00","w":"8.40003127","c":"8.40000000","Q":"0.49293200","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64750.63418500","q":"543907.35169250","O":1773023220000,"C":1773109631555,"F":19847634,"L":19911287,"n":217505}"#).unwrap_or_else(|_| serde_json::json!({}));
2985 let msg = json!({
2986 "stream": stream,
2987 "data": payload,
2988 });
2989
2990 streams_base.on_message(msg.to_string(), conn.clone()).await;
2991 yield_now().await;
2992
2993 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2994 });
2995 }
2996
2997 #[test]
2998 fn ticker_stream_should_not_fire_after_unsubscribe() {
2999 TOKIO_SHARED_RT.block_on(async {
3000 let (streams_base, conn) = make_streams_base().await;
3001 let api = ApiClient::new(streams_base.clone());
3002
3003 let id = "test-id-123".to_string();
3004
3005 let params = TickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
3006
3007 let TickerStreamParams {
3008 symbol,id,
3009 } = params.clone();
3010
3011 let pairs: &[(&str, Option<String>)] = &[
3012 ("symbol",
3013 Some(symbol.clone())
3014 ),
3015 ("id",
3016 id.clone()
3017 ),
3018 ];
3019
3020 let vars: HashMap<_, _> = pairs
3021 .iter()
3022 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3023 .collect();
3024 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
3025
3026 let ws_stream = api.ticker_stream(params).await.unwrap();
3027
3028 let called = Arc::new(AtomicBool::new(false));
3029 let called_clone = called.clone();
3030 ws_stream.on_message(move |_payload: models::TickerStreamResponse| {
3031 called_clone.store(true, Ordering::SeqCst);
3032 });
3033
3034 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3035
3036 ws_stream.unsubscribe().await;
3037
3038 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1773109631569,"s":"ALPHA_116USDT","p":"0.00000000","P":"0.00","w":"8.40003127","c":"8.40000000","Q":"0.49293200","o":"8.40000000","h":"8.50000000","l":"8.30000000","v":"64750.63418500","q":"543907.35169250","O":1773023220000,"C":1773109631555,"F":19847634,"L":19911287,"n":217505}"#).unwrap_or_else(|_| serde_json::json!({}));
3039 let msg = json!({
3040 "stream": stream,
3041 "data": payload,
3042 });
3043
3044 streams_base.on_message(msg.to_string(), conn.clone()).await;
3045
3046 yield_now().await;
3047
3048 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3049 });
3050 }
3051
3052 #[test]
3053 fn trade_stream_should_execute_successfully() {
3054 TOKIO_SHARED_RT.block_on(async {
3055 let (streams_base, _) = make_streams_base().await;
3056 let api = ApiClient::new(streams_base.clone());
3057
3058 let id = "test-id-123".to_string();
3059
3060 let params = TradeStreamParams::builder("alpha_116usdt".to_string())
3061 .id(Some(id.clone()))
3062 .build()
3063 .unwrap();
3064
3065 let TradeStreamParams { symbol, id } = params.clone();
3066
3067 let pairs: &[(&str, Option<String>)] =
3068 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3069
3070 let vars: HashMap<_, _> = pairs
3071 .iter()
3072 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3073 .collect();
3074 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3075 let ws_stream = api
3076 .trade_stream(params)
3077 .await
3078 .expect("trade_stream should return a WebsocketStream");
3079
3080 assert!(
3081 streams_base.is_subscribed(&stream).await,
3082 "expected stream '{stream}' to be subscribed"
3083 );
3084 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3085 });
3086 }
3087
3088 #[test]
3089 fn trade_stream_should_handle_incoming_message() {
3090 TOKIO_SHARED_RT.block_on(async {
3091 let (streams_base, conn) = make_streams_base().await;
3092 let api = ApiClient::new(streams_base.clone());
3093
3094 let id = "test-id-123".to_string();
3095
3096 let params = TradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
3097
3098 let TradeStreamParams {
3099 symbol,id,
3100 } = params.clone();
3101
3102 let pairs: &[(&str, Option<String>)] = &[
3103 ("symbol",
3104 Some(symbol.clone())
3105 ),
3106 ("id",
3107 id.clone()
3108 ),
3109 ];
3110
3111 let vars: HashMap<_, _> = pairs
3112 .iter()
3113 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3114 .collect();
3115 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3116
3117 let ws_stream = api.trade_stream(params).await.unwrap();
3118
3119 let called = Arc::new(AtomicBool::new(false));
3120 let called_with_message = called.clone();
3121 ws_stream.on_message(move |_payload: models::TradeStreamResponse| {
3122 called_with_message.store(true, Ordering::SeqCst);
3123 });
3124
3125 let payload: Value = serde_json::from_str(r#"{"e":"trade","E":1773110023891,"T":1773110023877,"s":"ALPHA_116USDT","t":19911650,"p":"8.40000000","q":"0.97915700","m":false}"#).unwrap_or_else(|_| serde_json::json!({}));
3126 let msg = json!({
3127 "stream": stream,
3128 "data": payload,
3129 });
3130
3131 streams_base.on_message(msg.to_string(), conn.clone()).await;
3132 yield_now().await;
3133
3134 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3135 });
3136 }
3137
3138 #[test]
3139 fn trade_stream_should_not_fire_after_unsubscribe() {
3140 TOKIO_SHARED_RT.block_on(async {
3141 let (streams_base, conn) = make_streams_base().await;
3142 let api = ApiClient::new(streams_base.clone());
3143
3144 let id = "test-id-123".to_string();
3145
3146 let params = TradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id.clone())).build().unwrap();
3147
3148 let TradeStreamParams {
3149 symbol,id,
3150 } = params.clone();
3151
3152 let pairs: &[(&str, Option<String>)] = &[
3153 ("symbol",
3154 Some(symbol.clone())
3155 ),
3156 ("id",
3157 id.clone()
3158 ),
3159 ];
3160
3161 let vars: HashMap<_, _> = pairs
3162 .iter()
3163 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3164 .collect();
3165 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3166
3167 let ws_stream = api.trade_stream(params).await.unwrap();
3168
3169 let called = Arc::new(AtomicBool::new(false));
3170 let called_clone = called.clone();
3171 ws_stream.on_message(move |_payload: models::TradeStreamResponse| {
3172 called_clone.store(true, Ordering::SeqCst);
3173 });
3174
3175 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3176
3177 ws_stream.unsubscribe().await;
3178
3179 let payload: Value = serde_json::from_str(r#"{"e":"trade","E":1773110023891,"T":1773110023877,"s":"ALPHA_116USDT","t":19911650,"p":"8.40000000","q":"0.97915700","m":false}"#).unwrap_or_else(|_| serde_json::json!({}));
3180 let msg = json!({
3181 "stream": stream,
3182 "data": payload,
3183 });
3184
3185 streams_base.on_message(msg.to_string(), conn.clone()).await;
3186
3187 yield_now().await;
3188
3189 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3190 });
3191 }
3192}