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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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<u32>,
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())),
799 ("id", id.map(|v| v.to_string())),
800 ];
801
802 let vars: HashMap<_, _> = pairs
803 .iter()
804 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
805 .collect();
806
807 let id_opt: Option<String> = vars.get("id").cloned();
808
809 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
810
811 Ok(
812 create_stream_handler::<models::AggregateTradeStreamResponse>(
813 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
814 stream,
815 id_opt.map(|s| {
816 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
817 if let Ok(n) = s.parse::<u32>() {
818 return StreamId::Number(n);
819 }
820 }
821 StreamId::Str(s)
822 }),
823 None,
824 )
825 .await,
826 )
827 }
828
829 async fn all_book_ticker_stream(
830 &self,
831 params: AllBookTickerStreamParams,
832 ) -> anyhow::Result<Arc<WebsocketStream<models::AllBookTickerStreamResponse>>> {
833 let AllBookTickerStreamParams { id } = params;
834
835 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
836
837 let vars: HashMap<_, _> = pairs
838 .iter()
839 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
840 .collect();
841
842 let id_opt: Option<String> = vars.get("id").cloned();
843
844 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
845
846 Ok(
847 create_stream_handler::<models::AllBookTickerStreamResponse>(
848 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
849 stream,
850 id_opt.map(|s| {
851 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
852 if let Ok(n) = s.parse::<u32>() {
853 return StreamId::Number(n);
854 }
855 }
856 StreamId::Str(s)
857 }),
858 None,
859 )
860 .await,
861 )
862 }
863
864 async fn all_mini_ticker_stream(
865 &self,
866 params: AllMiniTickerStreamParams,
867 ) -> anyhow::Result<Arc<WebsocketStream<models::AllMiniTickerStreamResponse>>> {
868 let AllMiniTickerStreamParams { id } = params;
869
870 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
871
872 let vars: HashMap<_, _> = pairs
873 .iter()
874 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
875 .collect();
876
877 let id_opt: Option<String> = vars.get("id").cloned();
878
879 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
880
881 Ok(
882 create_stream_handler::<models::AllMiniTickerStreamResponse>(
883 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
884 stream,
885 id_opt.map(|s| {
886 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
887 if let Ok(n) = s.parse::<u32>() {
888 return StreamId::Number(n);
889 }
890 }
891 StreamId::Str(s)
892 }),
893 None,
894 )
895 .await,
896 )
897 }
898
899 async fn all_ticker_stream(
900 &self,
901 params: AllTickerStreamParams,
902 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTickerStreamResponse>>> {
903 let AllTickerStreamParams { id } = params;
904
905 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
906
907 let vars: HashMap<_, _> = pairs
908 .iter()
909 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
910 .collect();
911
912 let id_opt: Option<String> = vars.get("id").cloned();
913
914 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
915
916 Ok(create_stream_handler::<models::AllTickerStreamResponse>(
917 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
918 stream,
919 id_opt.map(|s| {
920 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
921 if let Ok(n) = s.parse::<u32>() {
922 return StreamId::Number(n);
923 }
924 }
925 StreamId::Str(s)
926 }),
927 None,
928 )
929 .await)
930 }
931
932 async fn all_tokens24h_ticker_stream(
933 &self,
934 params: AllTokens24hTickerStreamParams,
935 ) -> anyhow::Result<Arc<WebsocketStream<models::AllTokens24hTickerStreamResponse>>> {
936 let AllTokens24hTickerStreamParams { id } = params;
937
938 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
939
940 let vars: HashMap<_, _> = pairs
941 .iter()
942 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
943 .collect();
944
945 let id_opt: Option<String> = vars.get("id").cloned();
946
947 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
948
949 Ok(
950 create_stream_handler::<models::AllTokens24hTickerStreamResponse>(
951 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
952 stream,
953 id_opt.map(|s| {
954 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
955 if let Ok(n) = s.parse::<u32>() {
956 return StreamId::Number(n);
957 }
958 }
959 StreamId::Str(s)
960 }),
961 None,
962 )
963 .await,
964 )
965 }
966
967 async fn book_ticker_stream(
968 &self,
969 params: BookTickerStreamParams,
970 ) -> anyhow::Result<Arc<WebsocketStream<models::BookTickerStreamResponse>>> {
971 let BookTickerStreamParams { symbol, id } = params;
972
973 let pairs: &[(&str, Option<String>)] = &[
974 ("symbol", Some(symbol.clone())),
975 ("id", id.map(|v| v.to_string())),
976 ];
977
978 let vars: HashMap<_, _> = pairs
979 .iter()
980 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
981 .collect();
982
983 let id_opt: Option<String> = vars.get("id").cloned();
984
985 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
986
987 Ok(create_stream_handler::<models::BookTickerStreamResponse>(
988 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
989 stream,
990 id_opt.map(|s| {
991 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
992 if let Ok(n) = s.parse::<u32>() {
993 return StreamId::Number(n);
994 }
995 }
996 StreamId::Str(s)
997 }),
998 None,
999 )
1000 .await)
1001 }
1002
1003 async fn contract_kline_stream(
1004 &self,
1005 params: ContractKlineStreamParams,
1006 ) -> anyhow::Result<Arc<WebsocketStream<models::ContractKlineStreamResponse>>> {
1007 let ContractKlineStreamParams {
1008 contract_address,
1009 chain_id,
1010 interval,
1011 id,
1012 } = params;
1013
1014 let pairs: &[(&str, Option<String>)] = &[
1015 ("contractAddress", Some(contract_address.clone())),
1016 ("chainId", Some(chain_id.clone())),
1017 ("interval", Some(interval.as_str().to_string())),
1018 ("id", id.map(|v| v.to_string())),
1019 ];
1020
1021 let vars: HashMap<_, _> = pairs
1022 .iter()
1023 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1024 .collect();
1025
1026 let id_opt: Option<String> = vars.get("id").cloned();
1027
1028 let stream = replace_websocket_streams_placeholders(
1029 "/came@<contractAddress>@<chainId>@kline_<interval>",
1030 &vars,
1031 );
1032
1033 Ok(
1034 create_stream_handler::<models::ContractKlineStreamResponse>(
1035 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1036 stream,
1037 id_opt.map(|s| {
1038 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1039 if let Ok(n) = s.parse::<u32>() {
1040 return StreamId::Number(n);
1041 }
1042 }
1043 StreamId::Str(s)
1044 }),
1045 None,
1046 )
1047 .await,
1048 )
1049 }
1050
1051 async fn full_depth_stream(
1052 &self,
1053 params: FullDepthStreamParams,
1054 ) -> anyhow::Result<Arc<WebsocketStream<models::FullDepthStreamResponse>>> {
1055 let FullDepthStreamParams {
1056 symbol,
1057 interval,
1058 id,
1059 } = params;
1060
1061 let pairs: &[(&str, Option<String>)] = &[
1062 ("symbol", Some(symbol.clone())),
1063 ("interval", Some(interval.as_str().to_string())),
1064 ("id", id.map(|v| v.to_string())),
1065 ];
1066
1067 let vars: HashMap<_, _> = pairs
1068 .iter()
1069 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1070 .collect();
1071
1072 let id_opt: Option<String> = vars.get("id").cloned();
1073
1074 let stream =
1075 replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
1076
1077 Ok(create_stream_handler::<models::FullDepthStreamResponse>(
1078 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1079 stream,
1080 id_opt.map(|s| {
1081 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1082 if let Ok(n) = s.parse::<u32>() {
1083 return StreamId::Number(n);
1084 }
1085 }
1086 StreamId::Str(s)
1087 }),
1088 None,
1089 )
1090 .await)
1091 }
1092
1093 async fn kline_stream(
1094 &self,
1095 params: KlineStreamParams,
1096 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineStreamResponse>>> {
1097 let KlineStreamParams {
1098 symbol,
1099 interval,
1100 id,
1101 } = params;
1102
1103 let pairs: &[(&str, Option<String>)] = &[
1104 ("symbol", Some(symbol.clone())),
1105 ("interval", Some(interval.as_str().to_string())),
1106 ("id", id.map(|v| v.to_string())),
1107 ];
1108
1109 let vars: HashMap<_, _> = pairs
1110 .iter()
1111 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1112 .collect();
1113
1114 let id_opt: Option<String> = vars.get("id").cloned();
1115
1116 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
1117
1118 Ok(create_stream_handler::<models::KlineStreamResponse>(
1119 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1120 stream,
1121 id_opt.map(|s| {
1122 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1123 if let Ok(n) = s.parse::<u32>() {
1124 return StreamId::Number(n);
1125 }
1126 }
1127 StreamId::Str(s)
1128 }),
1129 None,
1130 )
1131 .await)
1132 }
1133
1134 async fn mini_ticker_stream(
1135 &self,
1136 params: MiniTickerStreamParams,
1137 ) -> anyhow::Result<Arc<WebsocketStream<models::MiniTickerStreamResponse>>> {
1138 let MiniTickerStreamParams { symbol, id } = params;
1139
1140 let pairs: &[(&str, Option<String>)] = &[
1141 ("symbol", Some(symbol.clone())),
1142 ("id", id.map(|v| v.to_string())),
1143 ];
1144
1145 let vars: HashMap<_, _> = pairs
1146 .iter()
1147 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1148 .collect();
1149
1150 let id_opt: Option<String> = vars.get("id").cloned();
1151
1152 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
1153
1154 Ok(create_stream_handler::<models::MiniTickerStreamResponse>(
1155 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1156 stream,
1157 id_opt.map(|s| {
1158 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1159 if let Ok(n) = s.parse::<u32>() {
1160 return StreamId::Number(n);
1161 }
1162 }
1163 StreamId::Str(s)
1164 }),
1165 None,
1166 )
1167 .await)
1168 }
1169
1170 async fn partial_depth_stream(
1171 &self,
1172 params: PartialDepthStreamParams,
1173 ) -> anyhow::Result<Arc<WebsocketStream<models::PartialDepthStreamResponse>>> {
1174 let PartialDepthStreamParams {
1175 symbol,
1176 levels,
1177 interval,
1178 id,
1179 } = params;
1180
1181 let pairs: &[(&str, Option<String>)] = &[
1182 ("symbol", Some(symbol.clone())),
1183 ("levels", Some(levels.as_str().to_string())),
1184 ("interval", Some(interval.as_str().to_string())),
1185 ("id", id.map(|v| v.to_string())),
1186 ];
1187
1188 let vars: HashMap<_, _> = pairs
1189 .iter()
1190 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1191 .collect();
1192
1193 let id_opt: Option<String> = vars.get("id").cloned();
1194
1195 let stream =
1196 replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
1197
1198 Ok(create_stream_handler::<models::PartialDepthStreamResponse>(
1199 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1200 stream,
1201 id_opt.map(|s| {
1202 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1203 if let Ok(n) = s.parse::<u32>() {
1204 return StreamId::Number(n);
1205 }
1206 }
1207 StreamId::Str(s)
1208 }),
1209 None,
1210 )
1211 .await)
1212 }
1213
1214 async fn ticker_stream(
1215 &self,
1216 params: TickerStreamParams,
1217 ) -> anyhow::Result<Arc<WebsocketStream<models::TickerStreamResponse>>> {
1218 let TickerStreamParams { symbol, id } = params;
1219
1220 let pairs: &[(&str, Option<String>)] = &[
1221 ("symbol", Some(symbol.clone())),
1222 ("id", id.map(|v| v.to_string())),
1223 ];
1224
1225 let vars: HashMap<_, _> = pairs
1226 .iter()
1227 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1228 .collect();
1229
1230 let id_opt: Option<String> = vars.get("id").cloned();
1231
1232 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
1233
1234 Ok(create_stream_handler::<models::TickerStreamResponse>(
1235 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1236 stream,
1237 id_opt.map(|s| {
1238 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1239 if let Ok(n) = s.parse::<u32>() {
1240 return StreamId::Number(n);
1241 }
1242 }
1243 StreamId::Str(s)
1244 }),
1245 None,
1246 )
1247 .await)
1248 }
1249
1250 async fn trade_stream(
1251 &self,
1252 params: TradeStreamParams,
1253 ) -> anyhow::Result<Arc<WebsocketStream<models::TradeStreamResponse>>> {
1254 let TradeStreamParams { symbol, id } = params;
1255
1256 let pairs: &[(&str, Option<String>)] = &[
1257 ("symbol", Some(symbol.clone())),
1258 ("id", id.map(|v| v.to_string())),
1259 ];
1260
1261 let vars: HashMap<_, _> = pairs
1262 .iter()
1263 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1264 .collect();
1265
1266 let id_opt: Option<String> = vars.get("id").cloned();
1267
1268 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
1269
1270 Ok(create_stream_handler::<models::TradeStreamResponse>(
1271 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1272 stream,
1273 id_opt.map(|s| {
1274 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1275 if let Ok(n) = s.parse::<u32>() {
1276 return StreamId::Number(n);
1277 }
1278 }
1279 StreamId::Str(s)
1280 }),
1281 None,
1282 )
1283 .await)
1284 }
1285}
1286
1287#[cfg(all(test, feature = "alpha"))]
1288mod tests {
1289 use super::*;
1290 use crate::TOKIO_SHARED_RT;
1291 use crate::{
1292 common::websocket::{WebsocketConnection, WebsocketHandler},
1293 config::ConfigurationWebsocketStreams,
1294 };
1295 use serde_json::json;
1296 use std::sync::atomic::{AtomicBool, Ordering};
1297 use tokio::task::yield_now;
1298
1299 async fn make_streams_base() -> (Arc<WebsocketStreams>, Arc<WebsocketConnection>) {
1300 let conn = WebsocketConnection::new("test");
1301 let config = ConfigurationWebsocketStreams::builder()
1302 .build()
1303 .expect("Failed to build configuration");
1304 let streams_base = WebsocketStreams::new(config, vec![conn.clone()], vec![]);
1305 conn.set_handler(streams_base.clone() as Arc<dyn WebsocketHandler>)
1306 .await;
1307 (streams_base, conn)
1308 }
1309
1310 #[test]
1311 fn aggregate_trade_stream_should_execute_successfully() {
1312 TOKIO_SHARED_RT.block_on(async {
1313 let (streams_base, _) = make_streams_base().await;
1314 let api = ApiClient::new(streams_base.clone());
1315
1316 let id = 123456u32;
1317
1318 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string())
1319 .id(Some(id))
1320 .build()
1321 .unwrap();
1322
1323 let AggregateTradeStreamParams { symbol, id } = params.clone();
1324
1325 let pairs: &[(&str, Option<String>)] = &[
1326 ("symbol", Some(symbol.clone())),
1327 ("id", id.map(|v| v.to_string())),
1328 ];
1329
1330 let vars: HashMap<_, _> = pairs
1331 .iter()
1332 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1333 .collect();
1334 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1335 let ws_stream = api
1336 .aggregate_trade_stream(params)
1337 .await
1338 .expect("aggregate_trade_stream should return a WebsocketStream");
1339
1340 assert!(
1341 streams_base.is_subscribed(&stream).await,
1342 "expected stream '{stream}' to be subscribed"
1343 );
1344 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
1345 });
1346 }
1347
1348 #[test]
1349 fn aggregate_trade_stream_should_handle_incoming_message() {
1350 TOKIO_SHARED_RT.block_on(async {
1351 let (streams_base, conn) = make_streams_base().await;
1352 let api = ApiClient::new(streams_base.clone());
1353
1354 let id = 123456u32;
1355
1356 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
1357
1358 let AggregateTradeStreamParams {
1359 symbol,id,
1360 } = params.clone();
1361
1362 let pairs: &[(&str, Option<String>)] = &[
1363 ("symbol",
1364 Some(symbol.clone())
1365 ),
1366 ("id",
1367 id.map(|v| v.to_string())
1368 ),
1369 ];
1370
1371 let vars: HashMap<_, _> = pairs
1372 .iter()
1373 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1374 .collect();
1375 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1376
1377 let ws_stream = api.aggregate_trade_stream(params).await.unwrap();
1378
1379 let called = Arc::new(AtomicBool::new(false));
1380 let called_with_message = called.clone();
1381 ws_stream.on_message(move |_payload: models::AggregateTradeStreamResponse| {
1382 called_with_message.store(true, Ordering::SeqCst);
1383 });
1384
1385 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!({}));
1386 let msg = json!({
1387 "stream": stream,
1388 "data": payload,
1389 });
1390
1391 streams_base.on_message(msg.to_string(), conn.clone()).await;
1392 yield_now().await;
1393
1394 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1395 });
1396 }
1397
1398 #[test]
1399 fn aggregate_trade_stream_should_not_fire_after_unsubscribe() {
1400 TOKIO_SHARED_RT.block_on(async {
1401 let (streams_base, conn) = make_streams_base().await;
1402 let api = ApiClient::new(streams_base.clone());
1403
1404 let id = 123456u32;
1405
1406 let params = AggregateTradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
1407
1408 let AggregateTradeStreamParams {
1409 symbol,id,
1410 } = params.clone();
1411
1412 let pairs: &[(&str, Option<String>)] = &[
1413 ("symbol",
1414 Some(symbol.clone())
1415 ),
1416 ("id",
1417 id.map(|v| v.to_string())
1418 ),
1419 ];
1420
1421 let vars: HashMap<_, _> = pairs
1422 .iter()
1423 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1424 .collect();
1425 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1426
1427 let ws_stream = api.aggregate_trade_stream(params).await.unwrap();
1428
1429 let called = Arc::new(AtomicBool::new(false));
1430 let called_clone = called.clone();
1431 ws_stream.on_message(move |_payload: models::AggregateTradeStreamResponse| {
1432 called_clone.store(true, Ordering::SeqCst);
1433 });
1434
1435 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1436
1437 ws_stream.unsubscribe().await;
1438
1439 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!({}));
1440 let msg = json!({
1441 "stream": stream,
1442 "data": payload,
1443 });
1444
1445 streams_base.on_message(msg.to_string(), conn.clone()).await;
1446
1447 yield_now().await;
1448
1449 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1450 });
1451 }
1452
1453 #[test]
1454 fn all_book_ticker_stream_should_execute_successfully() {
1455 TOKIO_SHARED_RT.block_on(async {
1456 let (streams_base, _) = make_streams_base().await;
1457 let api = ApiClient::new(streams_base.clone());
1458
1459 let id = 123456u32;
1460
1461 let params = AllBookTickerStreamParams::builder()
1462 .id(Some(id))
1463 .build()
1464 .unwrap();
1465
1466 let AllBookTickerStreamParams { id } = params.clone();
1467
1468 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
1469
1470 let vars: HashMap<_, _> = pairs
1471 .iter()
1472 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1473 .collect();
1474 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1475 let ws_stream = api
1476 .all_book_ticker_stream(params)
1477 .await
1478 .expect("all_book_ticker_stream should return a WebsocketStream");
1479
1480 assert!(
1481 streams_base.is_subscribed(&stream).await,
1482 "expected stream '{stream}' to be subscribed"
1483 );
1484 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
1485 });
1486 }
1487
1488 #[test]
1489 fn all_book_ticker_stream_should_handle_incoming_message() {
1490 TOKIO_SHARED_RT.block_on(async {
1491 let (streams_base, conn) = make_streams_base().await;
1492 let api = ApiClient::new(streams_base.clone());
1493
1494 let id = 123456u32;
1495
1496 let params = AllBookTickerStreamParams::builder().id(Some(id)).build().unwrap();
1497
1498 let AllBookTickerStreamParams {
1499 id,
1500 } = params.clone();
1501
1502 let pairs: &[(&str, Option<String>)] = &[
1503 ("id",
1504 id.map(|v| v.to_string())
1505 ),
1506 ];
1507
1508 let vars: HashMap<_, _> = pairs
1509 .iter()
1510 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1511 .collect();
1512 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1513
1514 let ws_stream = api.all_book_ticker_stream(params).await.unwrap();
1515
1516 let called = Arc::new(AtomicBool::new(false));
1517 let called_with_message = called.clone();
1518 ws_stream.on_message(move |_payload: models::AllBookTickerStreamResponse| {
1519 called_with_message.store(true, Ordering::SeqCst);
1520 });
1521
1522 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!({}));
1523 let msg = json!({
1524 "stream": stream,
1525 "data": payload,
1526 });
1527
1528 streams_base.on_message(msg.to_string(), conn.clone()).await;
1529 yield_now().await;
1530
1531 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1532 });
1533 }
1534
1535 #[test]
1536 fn all_book_ticker_stream_should_not_fire_after_unsubscribe() {
1537 TOKIO_SHARED_RT.block_on(async {
1538 let (streams_base, conn) = make_streams_base().await;
1539 let api = ApiClient::new(streams_base.clone());
1540
1541 let id = 123456u32;
1542
1543 let params = AllBookTickerStreamParams::builder().id(Some(id)).build().unwrap();
1544
1545 let AllBookTickerStreamParams {
1546 id,
1547 } = params.clone();
1548
1549 let pairs: &[(&str, Option<String>)] = &[
1550 ("id",
1551 id.map(|v| v.to_string())
1552 ),
1553 ];
1554
1555 let vars: HashMap<_, _> = pairs
1556 .iter()
1557 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1558 .collect();
1559 let stream = replace_websocket_streams_placeholders("/!bookTicker", &vars);
1560
1561 let ws_stream = api.all_book_ticker_stream(params).await.unwrap();
1562
1563 let called = Arc::new(AtomicBool::new(false));
1564 let called_clone = called.clone();
1565 ws_stream.on_message(move |_payload: models::AllBookTickerStreamResponse| {
1566 called_clone.store(true, Ordering::SeqCst);
1567 });
1568
1569 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1570
1571 ws_stream.unsubscribe().await;
1572
1573 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!({}));
1574 let msg = json!({
1575 "stream": stream,
1576 "data": payload,
1577 });
1578
1579 streams_base.on_message(msg.to_string(), conn.clone()).await;
1580
1581 yield_now().await;
1582
1583 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1584 });
1585 }
1586
1587 #[test]
1588 fn all_mini_ticker_stream_should_execute_successfully() {
1589 TOKIO_SHARED_RT.block_on(async {
1590 let (streams_base, _) = make_streams_base().await;
1591 let api = ApiClient::new(streams_base.clone());
1592
1593 let id = 123456u32;
1594
1595 let params = AllMiniTickerStreamParams::builder()
1596 .id(Some(id))
1597 .build()
1598 .unwrap();
1599
1600 let AllMiniTickerStreamParams { id } = params.clone();
1601
1602 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
1603
1604 let vars: HashMap<_, _> = pairs
1605 .iter()
1606 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1607 .collect();
1608 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1609 let ws_stream = api
1610 .all_mini_ticker_stream(params)
1611 .await
1612 .expect("all_mini_ticker_stream should return a WebsocketStream");
1613
1614 assert!(
1615 streams_base.is_subscribed(&stream).await,
1616 "expected stream '{stream}' to be subscribed"
1617 );
1618 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
1619 });
1620 }
1621
1622 #[test]
1623 fn all_mini_ticker_stream_should_handle_incoming_message() {
1624 TOKIO_SHARED_RT.block_on(async {
1625 let (streams_base, conn) = make_streams_base().await;
1626 let api = ApiClient::new(streams_base.clone());
1627
1628 let id = 123456u32;
1629
1630 let params = AllMiniTickerStreamParams::builder().id(Some(id)).build().unwrap();
1631
1632 let AllMiniTickerStreamParams {
1633 id,
1634 } = params.clone();
1635
1636 let pairs: &[(&str, Option<String>)] = &[
1637 ("id",
1638 id.map(|v| v.to_string())
1639 ),
1640 ];
1641
1642 let vars: HashMap<_, _> = pairs
1643 .iter()
1644 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1645 .collect();
1646 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1647
1648 let ws_stream = api.all_mini_ticker_stream(params).await.unwrap();
1649
1650 let called = Arc::new(AtomicBool::new(false));
1651 let called_with_message = called.clone();
1652 ws_stream.on_message(move |_payload: models::AllMiniTickerStreamResponse| {
1653 called_with_message.store(true, Ordering::SeqCst);
1654 });
1655
1656 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!({}));
1657 let msg = json!({
1658 "stream": stream,
1659 "data": payload,
1660 });
1661
1662 streams_base.on_message(msg.to_string(), conn.clone()).await;
1663 yield_now().await;
1664
1665 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1666 });
1667 }
1668
1669 #[test]
1670 fn all_mini_ticker_stream_should_not_fire_after_unsubscribe() {
1671 TOKIO_SHARED_RT.block_on(async {
1672 let (streams_base, conn) = make_streams_base().await;
1673 let api = ApiClient::new(streams_base.clone());
1674
1675 let id = 123456u32;
1676
1677 let params = AllMiniTickerStreamParams::builder().id(Some(id)).build().unwrap();
1678
1679 let AllMiniTickerStreamParams {
1680 id,
1681 } = params.clone();
1682
1683 let pairs: &[(&str, Option<String>)] = &[
1684 ("id",
1685 id.map(|v| v.to_string())
1686 ),
1687 ];
1688
1689 let vars: HashMap<_, _> = pairs
1690 .iter()
1691 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1692 .collect();
1693 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1694
1695 let ws_stream = api.all_mini_ticker_stream(params).await.unwrap();
1696
1697 let called = Arc::new(AtomicBool::new(false));
1698 let called_clone = called.clone();
1699 ws_stream.on_message(move |_payload: models::AllMiniTickerStreamResponse| {
1700 called_clone.store(true, Ordering::SeqCst);
1701 });
1702
1703 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1704
1705 ws_stream.unsubscribe().await;
1706
1707 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!({}));
1708 let msg = json!({
1709 "stream": stream,
1710 "data": payload,
1711 });
1712
1713 streams_base.on_message(msg.to_string(), conn.clone()).await;
1714
1715 yield_now().await;
1716
1717 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1718 });
1719 }
1720
1721 #[test]
1722 fn all_ticker_stream_should_execute_successfully() {
1723 TOKIO_SHARED_RT.block_on(async {
1724 let (streams_base, _) = make_streams_base().await;
1725 let api = ApiClient::new(streams_base.clone());
1726
1727 let id = 123456u32;
1728
1729 let params = AllTickerStreamParams::builder()
1730 .id(Some(id))
1731 .build()
1732 .unwrap();
1733
1734 let AllTickerStreamParams { id } = params.clone();
1735
1736 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
1737
1738 let vars: HashMap<_, _> = pairs
1739 .iter()
1740 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1741 .collect();
1742 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1743 let ws_stream = api
1744 .all_ticker_stream(params)
1745 .await
1746 .expect("all_ticker_stream should return a WebsocketStream");
1747
1748 assert!(
1749 streams_base.is_subscribed(&stream).await,
1750 "expected stream '{stream}' to be subscribed"
1751 );
1752 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
1753 });
1754 }
1755
1756 #[test]
1757 fn all_ticker_stream_should_handle_incoming_message() {
1758 TOKIO_SHARED_RT.block_on(async {
1759 let (streams_base, conn) = make_streams_base().await;
1760 let api = ApiClient::new(streams_base.clone());
1761
1762 let id = 123456u32;
1763
1764 let params = AllTickerStreamParams::builder().id(Some(id)).build().unwrap();
1765
1766 let AllTickerStreamParams {
1767 id,
1768 } = params.clone();
1769
1770 let pairs: &[(&str, Option<String>)] = &[
1771 ("id",
1772 id.map(|v| v.to_string())
1773 ),
1774 ];
1775
1776 let vars: HashMap<_, _> = pairs
1777 .iter()
1778 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1779 .collect();
1780 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1781
1782 let ws_stream = api.all_ticker_stream(params).await.unwrap();
1783
1784 let called = Arc::new(AtomicBool::new(false));
1785 let called_with_message = called.clone();
1786 ws_stream.on_message(move |_payload: models::AllTickerStreamResponse| {
1787 called_with_message.store(true, Ordering::SeqCst);
1788 });
1789
1790 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!({}));
1791 let msg = json!({
1792 "stream": stream,
1793 "data": payload,
1794 });
1795
1796 streams_base.on_message(msg.to_string(), conn.clone()).await;
1797 yield_now().await;
1798
1799 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1800 });
1801 }
1802
1803 #[test]
1804 fn all_ticker_stream_should_not_fire_after_unsubscribe() {
1805 TOKIO_SHARED_RT.block_on(async {
1806 let (streams_base, conn) = make_streams_base().await;
1807 let api = ApiClient::new(streams_base.clone());
1808
1809 let id = 123456u32;
1810
1811 let params = AllTickerStreamParams::builder().id(Some(id)).build().unwrap();
1812
1813 let AllTickerStreamParams {
1814 id,
1815 } = params.clone();
1816
1817 let pairs: &[(&str, Option<String>)] = &[
1818 ("id",
1819 id.map(|v| v.to_string())
1820 ),
1821 ];
1822
1823 let vars: HashMap<_, _> = pairs
1824 .iter()
1825 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1826 .collect();
1827 let stream = replace_websocket_streams_placeholders("/!ticker@arr", &vars);
1828
1829 let ws_stream = api.all_ticker_stream(params).await.unwrap();
1830
1831 let called = Arc::new(AtomicBool::new(false));
1832 let called_clone = called.clone();
1833 ws_stream.on_message(move |_payload: models::AllTickerStreamResponse| {
1834 called_clone.store(true, Ordering::SeqCst);
1835 });
1836
1837 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1838
1839 ws_stream.unsubscribe().await;
1840
1841 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!({}));
1842 let msg = json!({
1843 "stream": stream,
1844 "data": payload,
1845 });
1846
1847 streams_base.on_message(msg.to_string(), conn.clone()).await;
1848
1849 yield_now().await;
1850
1851 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1852 });
1853 }
1854
1855 #[test]
1856 fn all_tokens24h_ticker_stream_should_execute_successfully() {
1857 TOKIO_SHARED_RT.block_on(async {
1858 let (streams_base, _) = make_streams_base().await;
1859 let api = ApiClient::new(streams_base.clone());
1860
1861 let id = 123456u32;
1862
1863 let params = AllTokens24hTickerStreamParams::builder()
1864 .id(Some(id))
1865 .build()
1866 .unwrap();
1867
1868 let AllTokens24hTickerStreamParams { id } = params.clone();
1869
1870 let pairs: &[(&str, Option<String>)] = &[("id", id.map(|v| v.to_string()))];
1871
1872 let vars: HashMap<_, _> = pairs
1873 .iter()
1874 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1875 .collect();
1876 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1877 let ws_stream = api
1878 .all_tokens24h_ticker_stream(params)
1879 .await
1880 .expect("all_tokens24h_ticker_stream should return a WebsocketStream");
1881
1882 assert!(
1883 streams_base.is_subscribed(&stream).await,
1884 "expected stream '{stream}' to be subscribed"
1885 );
1886 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
1887 });
1888 }
1889
1890 #[test]
1891 fn all_tokens24h_ticker_stream_should_handle_incoming_message() {
1892 TOKIO_SHARED_RT.block_on(async {
1893 let (streams_base, conn) = make_streams_base().await;
1894 let api = ApiClient::new(streams_base.clone());
1895
1896 let id = 123456u32;
1897
1898 let params = AllTokens24hTickerStreamParams::builder().id(Some(id)).build().unwrap();
1899
1900 let AllTokens24hTickerStreamParams {
1901 id,
1902 } = params.clone();
1903
1904 let pairs: &[(&str, Option<String>)] = &[
1905 ("id",
1906 id.map(|v| v.to_string())
1907 ),
1908 ];
1909
1910 let vars: HashMap<_, _> = pairs
1911 .iter()
1912 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1913 .collect();
1914 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1915
1916 let ws_stream = api.all_tokens24h_ticker_stream(params).await.unwrap();
1917
1918 let called = Arc::new(AtomicBool::new(false));
1919 let called_with_message = called.clone();
1920 ws_stream.on_message(move |_payload: models::AllTokens24hTickerStreamResponse| {
1921 called_with_message.store(true, Ordering::SeqCst);
1922 });
1923
1924 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!({}));
1925 let msg = json!({
1926 "stream": stream,
1927 "data": payload,
1928 });
1929
1930 streams_base.on_message(msg.to_string(), conn.clone()).await;
1931 yield_now().await;
1932
1933 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1934 });
1935 }
1936
1937 #[test]
1938 fn all_tokens24h_ticker_stream_should_not_fire_after_unsubscribe() {
1939 TOKIO_SHARED_RT.block_on(async {
1940 let (streams_base, conn) = make_streams_base().await;
1941 let api = ApiClient::new(streams_base.clone());
1942
1943 let id = 123456u32;
1944
1945 let params = AllTokens24hTickerStreamParams::builder().id(Some(id)).build().unwrap();
1946
1947 let AllTokens24hTickerStreamParams {
1948 id,
1949 } = params.clone();
1950
1951 let pairs: &[(&str, Option<String>)] = &[
1952 ("id",
1953 id.map(|v| v.to_string())
1954 ),
1955 ];
1956
1957 let vars: HashMap<_, _> = pairs
1958 .iter()
1959 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1960 .collect();
1961 let stream = replace_websocket_streams_placeholders("/came@allTokens@ticker24", &vars);
1962
1963 let ws_stream = api.all_tokens24h_ticker_stream(params).await.unwrap();
1964
1965 let called = Arc::new(AtomicBool::new(false));
1966 let called_clone = called.clone();
1967 ws_stream.on_message(move |_payload: models::AllTokens24hTickerStreamResponse| {
1968 called_clone.store(true, Ordering::SeqCst);
1969 });
1970
1971 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1972
1973 ws_stream.unsubscribe().await;
1974
1975 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!({}));
1976 let msg = json!({
1977 "stream": stream,
1978 "data": payload,
1979 });
1980
1981 streams_base.on_message(msg.to_string(), conn.clone()).await;
1982
1983 yield_now().await;
1984
1985 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1986 });
1987 }
1988
1989 #[test]
1990 fn book_ticker_stream_should_execute_successfully() {
1991 TOKIO_SHARED_RT.block_on(async {
1992 let (streams_base, _) = make_streams_base().await;
1993 let api = ApiClient::new(streams_base.clone());
1994
1995 let id = 123456u32;
1996
1997 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string())
1998 .id(Some(id))
1999 .build()
2000 .unwrap();
2001
2002 let BookTickerStreamParams { symbol, id } = params.clone();
2003
2004 let pairs: &[(&str, Option<String>)] = &[
2005 ("symbol", Some(symbol.clone())),
2006 ("id", id.map(|v| v.to_string())),
2007 ];
2008
2009 let vars: HashMap<_, _> = pairs
2010 .iter()
2011 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2012 .collect();
2013 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2014 let ws_stream = api
2015 .book_ticker_stream(params)
2016 .await
2017 .expect("book_ticker_stream should return a WebsocketStream");
2018
2019 assert!(
2020 streams_base.is_subscribed(&stream).await,
2021 "expected stream '{stream}' to be subscribed"
2022 );
2023 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2024 });
2025 }
2026
2027 #[test]
2028 fn book_ticker_stream_should_handle_incoming_message() {
2029 TOKIO_SHARED_RT.block_on(async {
2030 let (streams_base, conn) = make_streams_base().await;
2031 let api = ApiClient::new(streams_base.clone());
2032
2033 let id = 123456u32;
2034
2035 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
2036
2037 let BookTickerStreamParams {
2038 symbol,id,
2039 } = params.clone();
2040
2041 let pairs: &[(&str, Option<String>)] = &[
2042 ("symbol",
2043 Some(symbol.clone())
2044 ),
2045 ("id",
2046 id.map(|v| v.to_string())
2047 ),
2048 ];
2049
2050 let vars: HashMap<_, _> = pairs
2051 .iter()
2052 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2053 .collect();
2054 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2055
2056 let ws_stream = api.book_ticker_stream(params).await.unwrap();
2057
2058 let called = Arc::new(AtomicBool::new(false));
2059 let called_with_message = called.clone();
2060 ws_stream.on_message(move |_payload: models::BookTickerStreamResponse| {
2061 called_with_message.store(true, Ordering::SeqCst);
2062 });
2063
2064 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!({}));
2065 let msg = json!({
2066 "stream": stream,
2067 "data": payload,
2068 });
2069
2070 streams_base.on_message(msg.to_string(), conn.clone()).await;
2071 yield_now().await;
2072
2073 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2074 });
2075 }
2076
2077 #[test]
2078 fn book_ticker_stream_should_not_fire_after_unsubscribe() {
2079 TOKIO_SHARED_RT.block_on(async {
2080 let (streams_base, conn) = make_streams_base().await;
2081 let api = ApiClient::new(streams_base.clone());
2082
2083 let id = 123456u32;
2084
2085 let params = BookTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
2086
2087 let BookTickerStreamParams {
2088 symbol,id,
2089 } = params.clone();
2090
2091 let pairs: &[(&str, Option<String>)] = &[
2092 ("symbol",
2093 Some(symbol.clone())
2094 ),
2095 ("id",
2096 id.map(|v| v.to_string())
2097 ),
2098 ];
2099
2100 let vars: HashMap<_, _> = pairs
2101 .iter()
2102 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2103 .collect();
2104 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2105
2106 let ws_stream = api.book_ticker_stream(params).await.unwrap();
2107
2108 let called = Arc::new(AtomicBool::new(false));
2109 let called_clone = called.clone();
2110 ws_stream.on_message(move |_payload: models::BookTickerStreamResponse| {
2111 called_clone.store(true, Ordering::SeqCst);
2112 });
2113
2114 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2115
2116 ws_stream.unsubscribe().await;
2117
2118 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!({}));
2119 let msg = json!({
2120 "stream": stream,
2121 "data": payload,
2122 });
2123
2124 streams_base.on_message(msg.to_string(), conn.clone()).await;
2125
2126 yield_now().await;
2127
2128 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2129 });
2130 }
2131
2132 #[test]
2133 fn contract_kline_stream_should_execute_successfully() {
2134 TOKIO_SHARED_RT.block_on(async {
2135 let (streams_base, _) = make_streams_base().await;
2136 let api = ApiClient::new(streams_base.clone());
2137
2138 let id = 123456u32;
2139
2140 let params = ContractKlineStreamParams::builder(
2141 "G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),
2142 "CT_501".to_string(),
2143 ContractKlineStreamIntervalEnum::Interval1s,
2144 )
2145 .id(Some(id))
2146 .build()
2147 .unwrap();
2148
2149 let ContractKlineStreamParams {
2150 contract_address,
2151 chain_id,
2152 interval,
2153 id,
2154 } = params.clone();
2155
2156 let pairs: &[(&str, Option<String>)] = &[
2157 ("contractAddress", Some(contract_address.clone())),
2158 ("chainId", Some(chain_id.clone())),
2159 ("interval", Some(interval.as_str().to_string())),
2160 ("id", id.map(|v| v.to_string())),
2161 ];
2162
2163 let vars: HashMap<_, _> = pairs
2164 .iter()
2165 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2166 .collect();
2167 let stream = replace_websocket_streams_placeholders(
2168 "/came@<contractAddress>@<chainId>@kline_<interval>",
2169 &vars,
2170 );
2171 let ws_stream = api
2172 .contract_kline_stream(params)
2173 .await
2174 .expect("contract_kline_stream should return a WebsocketStream");
2175
2176 assert!(
2177 streams_base.is_subscribed(&stream).await,
2178 "expected stream '{stream}' to be subscribed"
2179 );
2180 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2181 });
2182 }
2183
2184 #[test]
2185 fn contract_kline_stream_should_handle_incoming_message() {
2186 TOKIO_SHARED_RT.block_on(async {
2187 let (streams_base, conn) = make_streams_base().await;
2188 let api = ApiClient::new(streams_base.clone());
2189
2190 let id = 123456u32;
2191
2192 let params = ContractKlineStreamParams::builder("G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),"CT_501".to_string(),ContractKlineStreamIntervalEnum::Interval1s,).id(Some(id)).build().unwrap();
2193
2194 let ContractKlineStreamParams {
2195 contract_address,chain_id,interval,id,
2196 } = params.clone();
2197
2198 let pairs: &[(&str, Option<String>)] = &[
2199 ("contractAddress",
2200 Some(contract_address.clone())
2201 ),
2202 ("chainId",
2203 Some(chain_id.clone())
2204 ),
2205 ("interval",
2206 Some(interval.as_str().to_string())
2207 ),
2208 ("id",
2209 id.map(|v| v.to_string())
2210 ),
2211 ];
2212
2213 let vars: HashMap<_, _> = pairs
2214 .iter()
2215 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2216 .collect();
2217 let stream = replace_websocket_streams_placeholders("/came@<contractAddress>@<chainId>@kline_<interval>", &vars);
2218
2219 let ws_stream = api.contract_kline_stream(params).await.unwrap();
2220
2221 let called = Arc::new(AtomicBool::new(false));
2222 let called_with_message = called.clone();
2223 ws_stream.on_message(move |_payload: models::ContractKlineStreamResponse| {
2224 called_with_message.store(true, Ordering::SeqCst);
2225 });
2226
2227 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!({}));
2228 let msg = json!({
2229 "stream": stream,
2230 "data": payload,
2231 });
2232
2233 streams_base.on_message(msg.to_string(), conn.clone()).await;
2234 yield_now().await;
2235
2236 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2237 });
2238 }
2239
2240 #[test]
2241 fn contract_kline_stream_should_not_fire_after_unsubscribe() {
2242 TOKIO_SHARED_RT.block_on(async {
2243 let (streams_base, conn) = make_streams_base().await;
2244 let api = ApiClient::new(streams_base.clone());
2245
2246 let id = 123456u32;
2247
2248 let params = ContractKlineStreamParams::builder("G7vQWurMkMMm2dU3iZpXYFTHT9Biio4F4gZCrwFpKNwG".to_string(),"CT_501".to_string(),ContractKlineStreamIntervalEnum::Interval1s,).id(Some(id)).build().unwrap();
2249
2250 let ContractKlineStreamParams {
2251 contract_address,chain_id,interval,id,
2252 } = params.clone();
2253
2254 let pairs: &[(&str, Option<String>)] = &[
2255 ("contractAddress",
2256 Some(contract_address.clone())
2257 ),
2258 ("chainId",
2259 Some(chain_id.clone())
2260 ),
2261 ("interval",
2262 Some(interval.as_str().to_string())
2263 ),
2264 ("id",
2265 id.map(|v| v.to_string())
2266 ),
2267 ];
2268
2269 let vars: HashMap<_, _> = pairs
2270 .iter()
2271 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2272 .collect();
2273 let stream = replace_websocket_streams_placeholders("/came@<contractAddress>@<chainId>@kline_<interval>", &vars);
2274
2275 let ws_stream = api.contract_kline_stream(params).await.unwrap();
2276
2277 let called = Arc::new(AtomicBool::new(false));
2278 let called_clone = called.clone();
2279 ws_stream.on_message(move |_payload: models::ContractKlineStreamResponse| {
2280 called_clone.store(true, Ordering::SeqCst);
2281 });
2282
2283 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2284
2285 ws_stream.unsubscribe().await;
2286
2287 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!({}));
2288 let msg = json!({
2289 "stream": stream,
2290 "data": payload,
2291 });
2292
2293 streams_base.on_message(msg.to_string(), conn.clone()).await;
2294
2295 yield_now().await;
2296
2297 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2298 });
2299 }
2300
2301 #[test]
2302 fn full_depth_stream_should_execute_successfully() {
2303 TOKIO_SHARED_RT.block_on(async {
2304 let (streams_base, _) = make_streams_base().await;
2305 let api = ApiClient::new(streams_base.clone());
2306
2307 let id = 123456u32;
2308
2309 let params = FullDepthStreamParams::builder(
2310 "alpha_116usdt".to_string(),
2311 FullDepthStreamIntervalEnum::Interval0ms,
2312 )
2313 .id(Some(id))
2314 .build()
2315 .unwrap();
2316
2317 let FullDepthStreamParams {
2318 symbol,
2319 interval,
2320 id,
2321 } = params.clone();
2322
2323 let pairs: &[(&str, Option<String>)] = &[
2324 ("symbol", Some(symbol.clone())),
2325 ("interval", Some(interval.as_str().to_string())),
2326 ("id", id.map(|v| v.to_string())),
2327 ];
2328
2329 let vars: HashMap<_, _> = pairs
2330 .iter()
2331 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2332 .collect();
2333 let stream =
2334 replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2335 let ws_stream = api
2336 .full_depth_stream(params)
2337 .await
2338 .expect("full_depth_stream should return a WebsocketStream");
2339
2340 assert!(
2341 streams_base.is_subscribed(&stream).await,
2342 "expected stream '{stream}' to be subscribed"
2343 );
2344 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2345 });
2346 }
2347
2348 #[test]
2349 fn full_depth_stream_should_handle_incoming_message() {
2350 TOKIO_SHARED_RT.block_on(async {
2351 let (streams_base, conn) = make_streams_base().await;
2352 let api = ApiClient::new(streams_base.clone());
2353
2354 let id = 123456u32;
2355
2356 let params = FullDepthStreamParams::builder("alpha_116usdt".to_string(),FullDepthStreamIntervalEnum::Interval0ms,).id(Some(id)).build().unwrap();
2357
2358 let FullDepthStreamParams {
2359 symbol,interval,id,
2360 } = params.clone();
2361
2362 let pairs: &[(&str, Option<String>)] = &[
2363 ("symbol",
2364 Some(symbol.clone())
2365 ),
2366 ("interval",
2367 Some(interval.as_str().to_string())
2368 ),
2369 ("id",
2370 id.map(|v| v.to_string())
2371 ),
2372 ];
2373
2374 let vars: HashMap<_, _> = pairs
2375 .iter()
2376 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2377 .collect();
2378 let stream = replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2379
2380 let ws_stream = api.full_depth_stream(params).await.unwrap();
2381
2382 let called = Arc::new(AtomicBool::new(false));
2383 let called_with_message = called.clone();
2384 ws_stream.on_message(move |_payload: models::FullDepthStreamResponse| {
2385 called_with_message.store(true, Ordering::SeqCst);
2386 });
2387
2388 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!({}));
2389 let msg = json!({
2390 "stream": stream,
2391 "data": payload,
2392 });
2393
2394 streams_base.on_message(msg.to_string(), conn.clone()).await;
2395 yield_now().await;
2396
2397 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2398 });
2399 }
2400
2401 #[test]
2402 fn full_depth_stream_should_not_fire_after_unsubscribe() {
2403 TOKIO_SHARED_RT.block_on(async {
2404 let (streams_base, conn) = make_streams_base().await;
2405 let api = ApiClient::new(streams_base.clone());
2406
2407 let id = 123456u32;
2408
2409 let params = FullDepthStreamParams::builder("alpha_116usdt".to_string(),FullDepthStreamIntervalEnum::Interval0ms,).id(Some(id)).build().unwrap();
2410
2411 let FullDepthStreamParams {
2412 symbol,interval,id,
2413 } = params.clone();
2414
2415 let pairs: &[(&str, Option<String>)] = &[
2416 ("symbol",
2417 Some(symbol.clone())
2418 ),
2419 ("interval",
2420 Some(interval.as_str().to_string())
2421 ),
2422 ("id",
2423 id.map(|v| v.to_string())
2424 ),
2425 ];
2426
2427 let vars: HashMap<_, _> = pairs
2428 .iter()
2429 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2430 .collect();
2431 let stream = replace_websocket_streams_placeholders("/<symbol>@fulldepth@<interval>", &vars);
2432
2433 let ws_stream = api.full_depth_stream(params).await.unwrap();
2434
2435 let called = Arc::new(AtomicBool::new(false));
2436 let called_clone = called.clone();
2437 ws_stream.on_message(move |_payload: models::FullDepthStreamResponse| {
2438 called_clone.store(true, Ordering::SeqCst);
2439 });
2440
2441 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2442
2443 ws_stream.unsubscribe().await;
2444
2445 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!({}));
2446 let msg = json!({
2447 "stream": stream,
2448 "data": payload,
2449 });
2450
2451 streams_base.on_message(msg.to_string(), conn.clone()).await;
2452
2453 yield_now().await;
2454
2455 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2456 });
2457 }
2458
2459 #[test]
2460 fn kline_stream_should_execute_successfully() {
2461 TOKIO_SHARED_RT.block_on(async {
2462 let (streams_base, _) = make_streams_base().await;
2463 let api = ApiClient::new(streams_base.clone());
2464
2465 let id = 123456u32;
2466
2467 let params = KlineStreamParams::builder(
2468 "alpha_116usdt".to_string(),
2469 KlineStreamIntervalEnum::Interval1m,
2470 )
2471 .id(Some(id))
2472 .build()
2473 .unwrap();
2474
2475 let KlineStreamParams {
2476 symbol,
2477 interval,
2478 id,
2479 } = params.clone();
2480
2481 let pairs: &[(&str, Option<String>)] = &[
2482 ("symbol", Some(symbol.clone())),
2483 ("interval", Some(interval.as_str().to_string())),
2484 ("id", id.map(|v| v.to_string())),
2485 ];
2486
2487 let vars: HashMap<_, _> = pairs
2488 .iter()
2489 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2490 .collect();
2491 let stream =
2492 replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2493 let ws_stream = api
2494 .kline_stream(params)
2495 .await
2496 .expect("kline_stream should return a WebsocketStream");
2497
2498 assert!(
2499 streams_base.is_subscribed(&stream).await,
2500 "expected stream '{stream}' to be subscribed"
2501 );
2502 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2503 });
2504 }
2505
2506 #[test]
2507 fn kline_stream_should_handle_incoming_message() {
2508 TOKIO_SHARED_RT.block_on(async {
2509 let (streams_base, conn) = make_streams_base().await;
2510 let api = ApiClient::new(streams_base.clone());
2511
2512 let id = 123456u32;
2513
2514 let params = KlineStreamParams::builder("alpha_116usdt".to_string(),KlineStreamIntervalEnum::Interval1m,).id(Some(id)).build().unwrap();
2515
2516 let KlineStreamParams {
2517 symbol,interval,id,
2518 } = params.clone();
2519
2520 let pairs: &[(&str, Option<String>)] = &[
2521 ("symbol",
2522 Some(symbol.clone())
2523 ),
2524 ("interval",
2525 Some(interval.as_str().to_string())
2526 ),
2527 ("id",
2528 id.map(|v| v.to_string())
2529 ),
2530 ];
2531
2532 let vars: HashMap<_, _> = pairs
2533 .iter()
2534 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2535 .collect();
2536 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2537
2538 let ws_stream = api.kline_stream(params).await.unwrap();
2539
2540 let called = Arc::new(AtomicBool::new(false));
2541 let called_with_message = called.clone();
2542 ws_stream.on_message(move |_payload: models::KlineStreamResponse| {
2543 called_with_message.store(true, Ordering::SeqCst);
2544 });
2545
2546 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!({}));
2547 let msg = json!({
2548 "stream": stream,
2549 "data": payload,
2550 });
2551
2552 streams_base.on_message(msg.to_string(), conn.clone()).await;
2553 yield_now().await;
2554
2555 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2556 });
2557 }
2558
2559 #[test]
2560 fn kline_stream_should_not_fire_after_unsubscribe() {
2561 TOKIO_SHARED_RT.block_on(async {
2562 let (streams_base, conn) = make_streams_base().await;
2563 let api = ApiClient::new(streams_base.clone());
2564
2565 let id = 123456u32;
2566
2567 let params = KlineStreamParams::builder("alpha_116usdt".to_string(),KlineStreamIntervalEnum::Interval1m,).id(Some(id)).build().unwrap();
2568
2569 let KlineStreamParams {
2570 symbol,interval,id,
2571 } = params.clone();
2572
2573 let pairs: &[(&str, Option<String>)] = &[
2574 ("symbol",
2575 Some(symbol.clone())
2576 ),
2577 ("interval",
2578 Some(interval.as_str().to_string())
2579 ),
2580 ("id",
2581 id.map(|v| v.to_string())
2582 ),
2583 ];
2584
2585 let vars: HashMap<_, _> = pairs
2586 .iter()
2587 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2588 .collect();
2589 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2590
2591 let ws_stream = api.kline_stream(params).await.unwrap();
2592
2593 let called = Arc::new(AtomicBool::new(false));
2594 let called_clone = called.clone();
2595 ws_stream.on_message(move |_payload: models::KlineStreamResponse| {
2596 called_clone.store(true, Ordering::SeqCst);
2597 });
2598
2599 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2600
2601 ws_stream.unsubscribe().await;
2602
2603 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!({}));
2604 let msg = json!({
2605 "stream": stream,
2606 "data": payload,
2607 });
2608
2609 streams_base.on_message(msg.to_string(), conn.clone()).await;
2610
2611 yield_now().await;
2612
2613 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2614 });
2615 }
2616
2617 #[test]
2618 fn mini_ticker_stream_should_execute_successfully() {
2619 TOKIO_SHARED_RT.block_on(async {
2620 let (streams_base, _) = make_streams_base().await;
2621 let api = ApiClient::new(streams_base.clone());
2622
2623 let id = 123456u32;
2624
2625 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string())
2626 .id(Some(id))
2627 .build()
2628 .unwrap();
2629
2630 let MiniTickerStreamParams { symbol, id } = params.clone();
2631
2632 let pairs: &[(&str, Option<String>)] = &[
2633 ("symbol", Some(symbol.clone())),
2634 ("id", id.map(|v| v.to_string())),
2635 ];
2636
2637 let vars: HashMap<_, _> = pairs
2638 .iter()
2639 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2640 .collect();
2641 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2642 let ws_stream = api
2643 .mini_ticker_stream(params)
2644 .await
2645 .expect("mini_ticker_stream should return a WebsocketStream");
2646
2647 assert!(
2648 streams_base.is_subscribed(&stream).await,
2649 "expected stream '{stream}' to be subscribed"
2650 );
2651 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2652 });
2653 }
2654
2655 #[test]
2656 fn mini_ticker_stream_should_handle_incoming_message() {
2657 TOKIO_SHARED_RT.block_on(async {
2658 let (streams_base, conn) = make_streams_base().await;
2659 let api = ApiClient::new(streams_base.clone());
2660
2661 let id = 123456u32;
2662
2663 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
2664
2665 let MiniTickerStreamParams {
2666 symbol,id,
2667 } = params.clone();
2668
2669 let pairs: &[(&str, Option<String>)] = &[
2670 ("symbol",
2671 Some(symbol.clone())
2672 ),
2673 ("id",
2674 id.map(|v| v.to_string())
2675 ),
2676 ];
2677
2678 let vars: HashMap<_, _> = pairs
2679 .iter()
2680 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2681 .collect();
2682 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2683
2684 let ws_stream = api.mini_ticker_stream(params).await.unwrap();
2685
2686 let called = Arc::new(AtomicBool::new(false));
2687 let called_with_message = called.clone();
2688 ws_stream.on_message(move |_payload: models::MiniTickerStreamResponse| {
2689 called_with_message.store(true, Ordering::SeqCst);
2690 });
2691
2692 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!({}));
2693 let msg = json!({
2694 "stream": stream,
2695 "data": payload,
2696 });
2697
2698 streams_base.on_message(msg.to_string(), conn.clone()).await;
2699 yield_now().await;
2700
2701 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2702 });
2703 }
2704
2705 #[test]
2706 fn mini_ticker_stream_should_not_fire_after_unsubscribe() {
2707 TOKIO_SHARED_RT.block_on(async {
2708 let (streams_base, conn) = make_streams_base().await;
2709 let api = ApiClient::new(streams_base.clone());
2710
2711 let id = 123456u32;
2712
2713 let params = MiniTickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
2714
2715 let MiniTickerStreamParams {
2716 symbol,id,
2717 } = params.clone();
2718
2719 let pairs: &[(&str, Option<String>)] = &[
2720 ("symbol",
2721 Some(symbol.clone())
2722 ),
2723 ("id",
2724 id.map(|v| v.to_string())
2725 ),
2726 ];
2727
2728 let vars: HashMap<_, _> = pairs
2729 .iter()
2730 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2731 .collect();
2732 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2733
2734 let ws_stream = api.mini_ticker_stream(params).await.unwrap();
2735
2736 let called = Arc::new(AtomicBool::new(false));
2737 let called_clone = called.clone();
2738 ws_stream.on_message(move |_payload: models::MiniTickerStreamResponse| {
2739 called_clone.store(true, Ordering::SeqCst);
2740 });
2741
2742 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2743
2744 ws_stream.unsubscribe().await;
2745
2746 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!({}));
2747 let msg = json!({
2748 "stream": stream,
2749 "data": payload,
2750 });
2751
2752 streams_base.on_message(msg.to_string(), conn.clone()).await;
2753
2754 yield_now().await;
2755
2756 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2757 });
2758 }
2759
2760 #[test]
2761 fn partial_depth_stream_should_execute_successfully() {
2762 TOKIO_SHARED_RT.block_on(async {
2763 let (streams_base, _) = make_streams_base().await;
2764 let api = ApiClient::new(streams_base.clone());
2765
2766 let id = 123456u32;
2767
2768 let params = PartialDepthStreamParams::builder(
2769 "alpha_116usdt".to_string(),
2770 PartialDepthStreamLevelsEnum::Levels5,
2771 PartialDepthStreamIntervalEnum::Interval0ms,
2772 )
2773 .id(Some(id))
2774 .build()
2775 .unwrap();
2776
2777 let PartialDepthStreamParams {
2778 symbol,
2779 levels,
2780 interval,
2781 id,
2782 } = params.clone();
2783
2784 let pairs: &[(&str, Option<String>)] = &[
2785 ("symbol", Some(symbol.clone())),
2786 ("levels", Some(levels.as_str().to_string())),
2787 ("interval", Some(interval.as_str().to_string())),
2788 ("id", id.map(|v| v.to_string())),
2789 ];
2790
2791 let vars: HashMap<_, _> = pairs
2792 .iter()
2793 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2794 .collect();
2795 let stream =
2796 replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2797 let ws_stream = api
2798 .partial_depth_stream(params)
2799 .await
2800 .expect("partial_depth_stream should return a WebsocketStream");
2801
2802 assert!(
2803 streams_base.is_subscribed(&stream).await,
2804 "expected stream '{stream}' to be subscribed"
2805 );
2806 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2807 });
2808 }
2809
2810 #[test]
2811 fn partial_depth_stream_should_handle_incoming_message() {
2812 TOKIO_SHARED_RT.block_on(async {
2813 let (streams_base, conn) = make_streams_base().await;
2814 let api = ApiClient::new(streams_base.clone());
2815
2816 let id = 123456u32;
2817
2818 let params = PartialDepthStreamParams::builder("alpha_116usdt".to_string(),PartialDepthStreamLevelsEnum::Levels5,PartialDepthStreamIntervalEnum::Interval0ms,).id(Some(id)).build().unwrap();
2819
2820 let PartialDepthStreamParams {
2821 symbol,levels,interval,id,
2822 } = params.clone();
2823
2824 let pairs: &[(&str, Option<String>)] = &[
2825 ("symbol",
2826 Some(symbol.clone())
2827 ),
2828 ("levels",
2829 Some(levels.as_str().to_string())
2830 ),
2831 ("interval",
2832 Some(interval.as_str().to_string())
2833 ),
2834 ("id",
2835 id.map(|v| v.to_string())
2836 ),
2837 ];
2838
2839 let vars: HashMap<_, _> = pairs
2840 .iter()
2841 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2842 .collect();
2843 let stream = replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2844
2845 let ws_stream = api.partial_depth_stream(params).await.unwrap();
2846
2847 let called = Arc::new(AtomicBool::new(false));
2848 let called_with_message = called.clone();
2849 ws_stream.on_message(move |_payload: models::PartialDepthStreamResponse| {
2850 called_with_message.store(true, Ordering::SeqCst);
2851 });
2852
2853 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!({}));
2854 let msg = json!({
2855 "stream": stream,
2856 "data": payload,
2857 });
2858
2859 streams_base.on_message(msg.to_string(), conn.clone()).await;
2860 yield_now().await;
2861
2862 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2863 });
2864 }
2865
2866 #[test]
2867 fn partial_depth_stream_should_not_fire_after_unsubscribe() {
2868 TOKIO_SHARED_RT.block_on(async {
2869 let (streams_base, conn) = make_streams_base().await;
2870 let api = ApiClient::new(streams_base.clone());
2871
2872 let id = 123456u32;
2873
2874 let params = PartialDepthStreamParams::builder("alpha_116usdt".to_string(),PartialDepthStreamLevelsEnum::Levels5,PartialDepthStreamIntervalEnum::Interval0ms,).id(Some(id)).build().unwrap();
2875
2876 let PartialDepthStreamParams {
2877 symbol,levels,interval,id,
2878 } = params.clone();
2879
2880 let pairs: &[(&str, Option<String>)] = &[
2881 ("symbol",
2882 Some(symbol.clone())
2883 ),
2884 ("levels",
2885 Some(levels.as_str().to_string())
2886 ),
2887 ("interval",
2888 Some(interval.as_str().to_string())
2889 ),
2890 ("id",
2891 id.map(|v| v.to_string())
2892 ),
2893 ];
2894
2895 let vars: HashMap<_, _> = pairs
2896 .iter()
2897 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2898 .collect();
2899 let stream = replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<interval>", &vars);
2900
2901 let ws_stream = api.partial_depth_stream(params).await.unwrap();
2902
2903 let called = Arc::new(AtomicBool::new(false));
2904 let called_clone = called.clone();
2905 ws_stream.on_message(move |_payload: models::PartialDepthStreamResponse| {
2906 called_clone.store(true, Ordering::SeqCst);
2907 });
2908
2909 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2910
2911 ws_stream.unsubscribe().await;
2912
2913 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!({}));
2914 let msg = json!({
2915 "stream": stream,
2916 "data": payload,
2917 });
2918
2919 streams_base.on_message(msg.to_string(), conn.clone()).await;
2920
2921 yield_now().await;
2922
2923 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2924 });
2925 }
2926
2927 #[test]
2928 fn ticker_stream_should_execute_successfully() {
2929 TOKIO_SHARED_RT.block_on(async {
2930 let (streams_base, _) = make_streams_base().await;
2931 let api = ApiClient::new(streams_base.clone());
2932
2933 let id = 123456u32;
2934
2935 let params = TickerStreamParams::builder("alpha_116usdt".to_string())
2936 .id(Some(id))
2937 .build()
2938 .unwrap();
2939
2940 let TickerStreamParams { symbol, id } = params.clone();
2941
2942 let pairs: &[(&str, Option<String>)] = &[
2943 ("symbol", Some(symbol.clone())),
2944 ("id", id.map(|v| v.to_string())),
2945 ];
2946
2947 let vars: HashMap<_, _> = pairs
2948 .iter()
2949 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2950 .collect();
2951 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
2952 let ws_stream = api
2953 .ticker_stream(params)
2954 .await
2955 .expect("ticker_stream should return a WebsocketStream");
2956
2957 assert!(
2958 streams_base.is_subscribed(&stream).await,
2959 "expected stream '{stream}' to be subscribed"
2960 );
2961 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
2962 });
2963 }
2964
2965 #[test]
2966 fn ticker_stream_should_handle_incoming_message() {
2967 TOKIO_SHARED_RT.block_on(async {
2968 let (streams_base, conn) = make_streams_base().await;
2969 let api = ApiClient::new(streams_base.clone());
2970
2971 let id = 123456u32;
2972
2973 let params = TickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
2974
2975 let TickerStreamParams {
2976 symbol,id,
2977 } = params.clone();
2978
2979 let pairs: &[(&str, Option<String>)] = &[
2980 ("symbol",
2981 Some(symbol.clone())
2982 ),
2983 ("id",
2984 id.map(|v| v.to_string())
2985 ),
2986 ];
2987
2988 let vars: HashMap<_, _> = pairs
2989 .iter()
2990 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2991 .collect();
2992 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
2993
2994 let ws_stream = api.ticker_stream(params).await.unwrap();
2995
2996 let called = Arc::new(AtomicBool::new(false));
2997 let called_with_message = called.clone();
2998 ws_stream.on_message(move |_payload: models::TickerStreamResponse| {
2999 called_with_message.store(true, Ordering::SeqCst);
3000 });
3001
3002 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!({}));
3003 let msg = json!({
3004 "stream": stream,
3005 "data": payload,
3006 });
3007
3008 streams_base.on_message(msg.to_string(), conn.clone()).await;
3009 yield_now().await;
3010
3011 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3012 });
3013 }
3014
3015 #[test]
3016 fn ticker_stream_should_not_fire_after_unsubscribe() {
3017 TOKIO_SHARED_RT.block_on(async {
3018 let (streams_base, conn) = make_streams_base().await;
3019 let api = ApiClient::new(streams_base.clone());
3020
3021 let id = 123456u32;
3022
3023 let params = TickerStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
3024
3025 let TickerStreamParams {
3026 symbol,id,
3027 } = params.clone();
3028
3029 let pairs: &[(&str, Option<String>)] = &[
3030 ("symbol",
3031 Some(symbol.clone())
3032 ),
3033 ("id",
3034 id.map(|v| v.to_string())
3035 ),
3036 ];
3037
3038 let vars: HashMap<_, _> = pairs
3039 .iter()
3040 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3041 .collect();
3042 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
3043
3044 let ws_stream = api.ticker_stream(params).await.unwrap();
3045
3046 let called = Arc::new(AtomicBool::new(false));
3047 let called_clone = called.clone();
3048 ws_stream.on_message(move |_payload: models::TickerStreamResponse| {
3049 called_clone.store(true, Ordering::SeqCst);
3050 });
3051
3052 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3053
3054 ws_stream.unsubscribe().await;
3055
3056 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!({}));
3057 let msg = json!({
3058 "stream": stream,
3059 "data": payload,
3060 });
3061
3062 streams_base.on_message(msg.to_string(), conn.clone()).await;
3063
3064 yield_now().await;
3065
3066 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3067 });
3068 }
3069
3070 #[test]
3071 fn trade_stream_should_execute_successfully() {
3072 TOKIO_SHARED_RT.block_on(async {
3073 let (streams_base, _) = make_streams_base().await;
3074 let api = ApiClient::new(streams_base.clone());
3075
3076 let id = 123456u32;
3077
3078 let params = TradeStreamParams::builder("alpha_116usdt".to_string())
3079 .id(Some(id))
3080 .build()
3081 .unwrap();
3082
3083 let TradeStreamParams { symbol, id } = params.clone();
3084
3085 let pairs: &[(&str, Option<String>)] = &[
3086 ("symbol", Some(symbol.clone())),
3087 ("id", id.map(|v| v.to_string())),
3088 ];
3089
3090 let vars: HashMap<_, _> = pairs
3091 .iter()
3092 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3093 .collect();
3094 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3095 let ws_stream = api
3096 .trade_stream(params)
3097 .await
3098 .expect("trade_stream should return a WebsocketStream");
3099
3100 assert!(
3101 streams_base.is_subscribed(&stream).await,
3102 "expected stream '{stream}' to be subscribed"
3103 );
3104 assert_eq!(ws_stream.id, Some(StreamId::Number(123456u32)));
3105 });
3106 }
3107
3108 #[test]
3109 fn trade_stream_should_handle_incoming_message() {
3110 TOKIO_SHARED_RT.block_on(async {
3111 let (streams_base, conn) = make_streams_base().await;
3112 let api = ApiClient::new(streams_base.clone());
3113
3114 let id = 123456u32;
3115
3116 let params = TradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
3117
3118 let TradeStreamParams {
3119 symbol,id,
3120 } = params.clone();
3121
3122 let pairs: &[(&str, Option<String>)] = &[
3123 ("symbol",
3124 Some(symbol.clone())
3125 ),
3126 ("id",
3127 id.map(|v| v.to_string())
3128 ),
3129 ];
3130
3131 let vars: HashMap<_, _> = pairs
3132 .iter()
3133 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3134 .collect();
3135 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3136
3137 let ws_stream = api.trade_stream(params).await.unwrap();
3138
3139 let called = Arc::new(AtomicBool::new(false));
3140 let called_with_message = called.clone();
3141 ws_stream.on_message(move |_payload: models::TradeStreamResponse| {
3142 called_with_message.store(true, Ordering::SeqCst);
3143 });
3144
3145 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!({}));
3146 let msg = json!({
3147 "stream": stream,
3148 "data": payload,
3149 });
3150
3151 streams_base.on_message(msg.to_string(), conn.clone()).await;
3152 yield_now().await;
3153
3154 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3155 });
3156 }
3157
3158 #[test]
3159 fn trade_stream_should_not_fire_after_unsubscribe() {
3160 TOKIO_SHARED_RT.block_on(async {
3161 let (streams_base, conn) = make_streams_base().await;
3162 let api = ApiClient::new(streams_base.clone());
3163
3164 let id = 123456u32;
3165
3166 let params = TradeStreamParams::builder("alpha_116usdt".to_string(),).id(Some(id)).build().unwrap();
3167
3168 let TradeStreamParams {
3169 symbol,id,
3170 } = params.clone();
3171
3172 let pairs: &[(&str, Option<String>)] = &[
3173 ("symbol",
3174 Some(symbol.clone())
3175 ),
3176 ("id",
3177 id.map(|v| v.to_string())
3178 ),
3179 ];
3180
3181 let vars: HashMap<_, _> = pairs
3182 .iter()
3183 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3184 .collect();
3185 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3186
3187 let ws_stream = api.trade_stream(params).await.unwrap();
3188
3189 let called = Arc::new(AtomicBool::new(false));
3190 let called_clone = called.clone();
3191 ws_stream.on_message(move |_payload: models::TradeStreamResponse| {
3192 called_clone.store(true, Ordering::SeqCst);
3193 });
3194
3195 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3196
3197 ws_stream.unsubscribe().await;
3198
3199 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!({}));
3200 let msg = json!({
3201 "stream": stream,
3202 "data": payload,
3203 });
3204
3205 streams_base.on_message(msg.to_string(), conn.clone()).await;
3206
3207 yield_now().await;
3208
3209 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3210 });
3211 }
3212}