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::common::{
22 models::ParamBuildError,
23 utils::replace_websocket_streams_placeholders,
24 websocket::{WebsocketBase, WebsocketStream, WebsocketStreams, create_stream_handler},
25};
26use crate::models::StreamId;
27use crate::spot::websocket_streams::models;
28
29#[async_trait]
30pub trait Api: Send + Sync {
31 async fn agg_trade(
32 &self,
33 params: AggTradeParams,
34 ) -> anyhow::Result<Arc<WebsocketStream<models::AggTradeResponse>>>;
35 async fn all_market_rolling_window_ticker(
36 &self,
37 params: AllMarketRollingWindowTickerParams,
38 ) -> anyhow::Result<Arc<WebsocketStream<Vec<models::AllMarketRollingWindowTickerResponseInner>>>>;
39 async fn all_mini_ticker(
40 &self,
41 params: AllMiniTickerParams,
42 ) -> anyhow::Result<Arc<WebsocketStream<Vec<models::AllMiniTickerResponseInner>>>>;
43 async fn avg_price(
44 &self,
45 params: AvgPriceParams,
46 ) -> anyhow::Result<Arc<WebsocketStream<models::AvgPriceResponse>>>;
47 async fn block_trade(
48 &self,
49 params: BlockTradeParams,
50 ) -> anyhow::Result<Arc<WebsocketStream<models::BlockTradeResponse>>>;
51 async fn book_ticker(
52 &self,
53 params: BookTickerParams,
54 ) -> anyhow::Result<Arc<WebsocketStream<models::BookTickerResponse>>>;
55 async fn diff_book_depth(
56 &self,
57 params: DiffBookDepthParams,
58 ) -> anyhow::Result<Arc<WebsocketStream<models::DiffBookDepthResponse>>>;
59 async fn kline(
60 &self,
61 params: KlineParams,
62 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineResponse>>>;
63 async fn kline_offset(
64 &self,
65 params: KlineOffsetParams,
66 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineOffsetResponse>>>;
67 async fn mini_ticker(
68 &self,
69 params: MiniTickerParams,
70 ) -> anyhow::Result<Arc<WebsocketStream<models::MiniTickerResponse>>>;
71 async fn partial_book_depth(
72 &self,
73 params: PartialBookDepthParams,
74 ) -> anyhow::Result<Arc<WebsocketStream<models::PartialBookDepthResponse>>>;
75 async fn reference_price(
76 &self,
77 params: ReferencePriceParams,
78 ) -> anyhow::Result<Arc<WebsocketStream<models::ReferencePriceResponse>>>;
79 async fn rolling_window_ticker(
80 &self,
81 params: RollingWindowTickerParams,
82 ) -> anyhow::Result<Arc<WebsocketStream<models::RollingWindowTickerResponse>>>;
83 async fn ticker(
84 &self,
85 params: TickerParams,
86 ) -> anyhow::Result<Arc<WebsocketStream<models::TickerResponse>>>;
87 async fn trade(
88 &self,
89 params: TradeParams,
90 ) -> anyhow::Result<Arc<WebsocketStream<models::TradeResponse>>>;
91}
92
93pub struct ApiClient {
94 websocket_streams_base: Arc<WebsocketStreams>,
95}
96
97impl ApiClient {
98 pub fn new(websocket_streams_base: Arc<WebsocketStreams>) -> Self {
99 Self {
100 websocket_streams_base,
101 }
102 }
103}
104
105#[allow(non_camel_case_types)]
106#[derive(Debug, Clone, Serialize, Deserialize)]
107pub enum AllMarketRollingWindowTickerWindowSizeEnum {
108 #[serde(rename = "1h")]
109 WindowSize1h,
110 #[serde(rename = "4h")]
111 WindowSize4h,
112 #[serde(rename = "1d")]
113 WindowSize1d,
114}
115
116impl AllMarketRollingWindowTickerWindowSizeEnum {
117 #[must_use]
118 pub fn as_str(&self) -> &'static str {
119 match self {
120 Self::WindowSize1h => "1h",
121 Self::WindowSize4h => "4h",
122 Self::WindowSize1d => "1d",
123 }
124 }
125}
126
127impl std::str::FromStr for AllMarketRollingWindowTickerWindowSizeEnum {
128 type Err = Box<dyn std::error::Error + Send + Sync>;
129
130 fn from_str(s: &str) -> Result<Self, Self::Err> {
131 match s {
132 "1h" => Ok(Self::WindowSize1h),
133 "4h" => Ok(Self::WindowSize4h),
134 "1d" => Ok(Self::WindowSize1d),
135 other => Err(format!(
136 "invalid AllMarketRollingWindowTickerWindowSizeEnum: {}",
137 other
138 )
139 .into()),
140 }
141 }
142}
143
144#[allow(non_camel_case_types)]
145#[derive(Debug, Clone, Serialize, Deserialize)]
146pub enum DiffBookDepthUpdateSpeedEnum {
147 #[serde(rename = "100ms")]
148 UpdateSpeed100ms,
149}
150
151impl DiffBookDepthUpdateSpeedEnum {
152 #[must_use]
153 pub fn as_str(&self) -> &'static str {
154 match self {
155 Self::UpdateSpeed100ms => "100ms",
156 }
157 }
158}
159
160impl std::str::FromStr for DiffBookDepthUpdateSpeedEnum {
161 type Err = Box<dyn std::error::Error + Send + Sync>;
162
163 fn from_str(s: &str) -> Result<Self, Self::Err> {
164 match s {
165 "100ms" => Ok(Self::UpdateSpeed100ms),
166 other => Err(format!("invalid DiffBookDepthUpdateSpeedEnum: {}", other).into()),
167 }
168 }
169}
170
171#[allow(non_camel_case_types)]
172#[derive(Debug, Clone, Serialize, Deserialize)]
173pub enum KlineIntervalEnum {
174 #[serde(rename = "1s")]
175 Interval1s,
176 #[serde(rename = "1m")]
177 Interval1m,
178 #[serde(rename = "3m")]
179 Interval3m,
180 #[serde(rename = "5m")]
181 Interval5m,
182 #[serde(rename = "15m")]
183 Interval15m,
184 #[serde(rename = "30m")]
185 Interval30m,
186 #[serde(rename = "1h")]
187 Interval1h,
188 #[serde(rename = "2h")]
189 Interval2h,
190 #[serde(rename = "4h")]
191 Interval4h,
192 #[serde(rename = "6h")]
193 Interval6h,
194 #[serde(rename = "8h")]
195 Interval8h,
196 #[serde(rename = "12h")]
197 Interval12h,
198 #[serde(rename = "1d")]
199 Interval1d,
200 #[serde(rename = "3d")]
201 Interval3d,
202 #[serde(rename = "1w")]
203 Interval1w,
204 #[serde(rename = "1M")]
205 Interval1M,
206}
207
208impl KlineIntervalEnum {
209 #[must_use]
210 pub fn as_str(&self) -> &'static str {
211 match self {
212 Self::Interval1s => "1s",
213 Self::Interval1m => "1m",
214 Self::Interval3m => "3m",
215 Self::Interval5m => "5m",
216 Self::Interval15m => "15m",
217 Self::Interval30m => "30m",
218 Self::Interval1h => "1h",
219 Self::Interval2h => "2h",
220 Self::Interval4h => "4h",
221 Self::Interval6h => "6h",
222 Self::Interval8h => "8h",
223 Self::Interval12h => "12h",
224 Self::Interval1d => "1d",
225 Self::Interval3d => "3d",
226 Self::Interval1w => "1w",
227 Self::Interval1M => "1M",
228 }
229 }
230}
231
232impl std::str::FromStr for KlineIntervalEnum {
233 type Err = Box<dyn std::error::Error + Send + Sync>;
234
235 fn from_str(s: &str) -> Result<Self, Self::Err> {
236 match s {
237 "1s" => Ok(Self::Interval1s),
238 "1m" => Ok(Self::Interval1m),
239 "3m" => Ok(Self::Interval3m),
240 "5m" => Ok(Self::Interval5m),
241 "15m" => Ok(Self::Interval15m),
242 "30m" => Ok(Self::Interval30m),
243 "1h" => Ok(Self::Interval1h),
244 "2h" => Ok(Self::Interval2h),
245 "4h" => Ok(Self::Interval4h),
246 "6h" => Ok(Self::Interval6h),
247 "8h" => Ok(Self::Interval8h),
248 "12h" => Ok(Self::Interval12h),
249 "1d" => Ok(Self::Interval1d),
250 "3d" => Ok(Self::Interval3d),
251 "1w" => Ok(Self::Interval1w),
252 "1M" => Ok(Self::Interval1M),
253 other => Err(format!("invalid KlineIntervalEnum: {}", other).into()),
254 }
255 }
256}
257
258#[allow(non_camel_case_types)]
259#[derive(Debug, Clone, Serialize, Deserialize)]
260pub enum KlineOffsetIntervalEnum {
261 #[serde(rename = "1s")]
262 Interval1s,
263 #[serde(rename = "1m")]
264 Interval1m,
265 #[serde(rename = "3m")]
266 Interval3m,
267 #[serde(rename = "5m")]
268 Interval5m,
269 #[serde(rename = "15m")]
270 Interval15m,
271 #[serde(rename = "30m")]
272 Interval30m,
273 #[serde(rename = "1h")]
274 Interval1h,
275 #[serde(rename = "2h")]
276 Interval2h,
277 #[serde(rename = "4h")]
278 Interval4h,
279 #[serde(rename = "6h")]
280 Interval6h,
281 #[serde(rename = "8h")]
282 Interval8h,
283 #[serde(rename = "12h")]
284 Interval12h,
285 #[serde(rename = "1d")]
286 Interval1d,
287 #[serde(rename = "3d")]
288 Interval3d,
289 #[serde(rename = "1w")]
290 Interval1w,
291 #[serde(rename = "1M")]
292 Interval1M,
293}
294
295impl KlineOffsetIntervalEnum {
296 #[must_use]
297 pub fn as_str(&self) -> &'static str {
298 match self {
299 Self::Interval1s => "1s",
300 Self::Interval1m => "1m",
301 Self::Interval3m => "3m",
302 Self::Interval5m => "5m",
303 Self::Interval15m => "15m",
304 Self::Interval30m => "30m",
305 Self::Interval1h => "1h",
306 Self::Interval2h => "2h",
307 Self::Interval4h => "4h",
308 Self::Interval6h => "6h",
309 Self::Interval8h => "8h",
310 Self::Interval12h => "12h",
311 Self::Interval1d => "1d",
312 Self::Interval3d => "3d",
313 Self::Interval1w => "1w",
314 Self::Interval1M => "1M",
315 }
316 }
317}
318
319impl std::str::FromStr for KlineOffsetIntervalEnum {
320 type Err = Box<dyn std::error::Error + Send + Sync>;
321
322 fn from_str(s: &str) -> Result<Self, Self::Err> {
323 match s {
324 "1s" => Ok(Self::Interval1s),
325 "1m" => Ok(Self::Interval1m),
326 "3m" => Ok(Self::Interval3m),
327 "5m" => Ok(Self::Interval5m),
328 "15m" => Ok(Self::Interval15m),
329 "30m" => Ok(Self::Interval30m),
330 "1h" => Ok(Self::Interval1h),
331 "2h" => Ok(Self::Interval2h),
332 "4h" => Ok(Self::Interval4h),
333 "6h" => Ok(Self::Interval6h),
334 "8h" => Ok(Self::Interval8h),
335 "12h" => Ok(Self::Interval12h),
336 "1d" => Ok(Self::Interval1d),
337 "3d" => Ok(Self::Interval3d),
338 "1w" => Ok(Self::Interval1w),
339 "1M" => Ok(Self::Interval1M),
340 other => Err(format!("invalid KlineOffsetIntervalEnum: {}", other).into()),
341 }
342 }
343}
344
345#[allow(non_camel_case_types)]
346#[derive(Debug, Clone, Serialize, Deserialize)]
347pub enum PartialBookDepthLevelsEnum {
348 #[serde(rename = "5")]
349 Levels5,
350 #[serde(rename = "10")]
351 Levels10,
352 #[serde(rename = "20")]
353 Levels20,
354}
355
356impl PartialBookDepthLevelsEnum {
357 #[must_use]
358 pub fn as_str(&self) -> &'static str {
359 match self {
360 Self::Levels5 => "5",
361 Self::Levels10 => "10",
362 Self::Levels20 => "20",
363 }
364 }
365}
366
367impl std::str::FromStr for PartialBookDepthLevelsEnum {
368 type Err = Box<dyn std::error::Error + Send + Sync>;
369
370 fn from_str(s: &str) -> Result<Self, Self::Err> {
371 match s {
372 "5" => Ok(Self::Levels5),
373 "10" => Ok(Self::Levels10),
374 "20" => Ok(Self::Levels20),
375 other => Err(format!("invalid PartialBookDepthLevelsEnum: {}", other).into()),
376 }
377 }
378}
379
380#[allow(non_camel_case_types)]
381#[derive(Debug, Clone, Serialize, Deserialize)]
382pub enum PartialBookDepthUpdateSpeedEnum {
383 #[serde(rename = "100ms")]
384 UpdateSpeed100ms,
385}
386
387impl PartialBookDepthUpdateSpeedEnum {
388 #[must_use]
389 pub fn as_str(&self) -> &'static str {
390 match self {
391 Self::UpdateSpeed100ms => "100ms",
392 }
393 }
394}
395
396impl std::str::FromStr for PartialBookDepthUpdateSpeedEnum {
397 type Err = Box<dyn std::error::Error + Send + Sync>;
398
399 fn from_str(s: &str) -> Result<Self, Self::Err> {
400 match s {
401 "100ms" => Ok(Self::UpdateSpeed100ms),
402 other => Err(format!("invalid PartialBookDepthUpdateSpeedEnum: {}", other).into()),
403 }
404 }
405}
406
407#[allow(non_camel_case_types)]
408#[derive(Debug, Clone, Serialize, Deserialize)]
409pub enum RollingWindowTickerWindowSizeEnum {
410 #[serde(rename = "1h")]
411 WindowSize1h,
412 #[serde(rename = "4h")]
413 WindowSize4h,
414 #[serde(rename = "1d")]
415 WindowSize1d,
416}
417
418impl RollingWindowTickerWindowSizeEnum {
419 #[must_use]
420 pub fn as_str(&self) -> &'static str {
421 match self {
422 Self::WindowSize1h => "1h",
423 Self::WindowSize4h => "4h",
424 Self::WindowSize1d => "1d",
425 }
426 }
427}
428
429impl std::str::FromStr for RollingWindowTickerWindowSizeEnum {
430 type Err = Box<dyn std::error::Error + Send + Sync>;
431
432 fn from_str(s: &str) -> Result<Self, Self::Err> {
433 match s {
434 "1h" => Ok(Self::WindowSize1h),
435 "4h" => Ok(Self::WindowSize4h),
436 "1d" => Ok(Self::WindowSize1d),
437 other => Err(format!("invalid RollingWindowTickerWindowSizeEnum: {}", other).into()),
438 }
439 }
440}
441
442#[derive(Clone, Debug, Builder, Deserialize)]
447#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
448pub struct AggTradeParams {
449 #[builder(setter(into))]
453 #[serde(rename = "symbol")]
454 pub symbol: String,
455 #[builder(setter(into), default)]
459 #[serde(rename = "id", default)]
460 pub id: Option<String>,
461}
462
463impl AggTradeParams {
464 #[must_use]
471 pub fn builder(symbol: String) -> AggTradeParamsBuilder {
472 AggTradeParamsBuilder::default().symbol(symbol)
473 }
474}
475#[derive(Clone, Debug, Builder, Deserialize)]
480#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
481pub struct AllMarketRollingWindowTickerParams {
482 #[builder(setter(into))]
487 #[serde(rename = "windowSize")]
488 pub window_size: AllMarketRollingWindowTickerWindowSizeEnum,
489 #[builder(setter(into), default)]
493 #[serde(rename = "id", default)]
494 pub id: Option<String>,
495}
496
497impl AllMarketRollingWindowTickerParams {
498 #[must_use]
505 pub fn builder(
506 window_size: AllMarketRollingWindowTickerWindowSizeEnum,
507 ) -> AllMarketRollingWindowTickerParamsBuilder {
508 AllMarketRollingWindowTickerParamsBuilder::default().window_size(window_size)
509 }
510}
511#[derive(Clone, Debug, Builder, Deserialize, Default)]
516#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
517pub struct AllMiniTickerParams {
518 #[builder(setter(into), default)]
522 #[serde(rename = "id", default)]
523 pub id: Option<String>,
524}
525
526impl AllMiniTickerParams {
527 #[must_use]
530 pub fn builder() -> AllMiniTickerParamsBuilder {
531 AllMiniTickerParamsBuilder::default()
532 }
533}
534#[derive(Clone, Debug, Builder, Deserialize)]
539#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
540pub struct AvgPriceParams {
541 #[builder(setter(into))]
545 #[serde(rename = "symbol")]
546 pub symbol: String,
547 #[builder(setter(into), default)]
551 #[serde(rename = "id", default)]
552 pub id: Option<String>,
553}
554
555impl AvgPriceParams {
556 #[must_use]
563 pub fn builder(symbol: String) -> AvgPriceParamsBuilder {
564 AvgPriceParamsBuilder::default().symbol(symbol)
565 }
566}
567#[derive(Clone, Debug, Builder, Deserialize)]
572#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
573pub struct BlockTradeParams {
574 #[builder(setter(into))]
578 #[serde(rename = "symbol")]
579 pub symbol: String,
580 #[builder(setter(into), default)]
584 #[serde(rename = "id", default)]
585 pub id: Option<String>,
586}
587
588impl BlockTradeParams {
589 #[must_use]
596 pub fn builder(symbol: String) -> BlockTradeParamsBuilder {
597 BlockTradeParamsBuilder::default().symbol(symbol)
598 }
599}
600#[derive(Clone, Debug, Builder, Deserialize)]
605#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
606pub struct BookTickerParams {
607 #[builder(setter(into))]
611 #[serde(rename = "symbol")]
612 pub symbol: String,
613 #[builder(setter(into), default)]
617 #[serde(rename = "id", default)]
618 pub id: Option<String>,
619}
620
621impl BookTickerParams {
622 #[must_use]
629 pub fn builder(symbol: String) -> BookTickerParamsBuilder {
630 BookTickerParamsBuilder::default().symbol(symbol)
631 }
632}
633#[derive(Clone, Debug, Builder, Deserialize)]
638#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
639pub struct DiffBookDepthParams {
640 #[builder(setter(into))]
644 #[serde(rename = "symbol")]
645 pub symbol: String,
646 #[builder(setter(into), default)]
650 #[serde(rename = "id", default)]
651 pub id: Option<String>,
652 #[builder(setter(into), default)]
656 #[serde(rename = "updateSpeed", default)]
657 pub update_speed: Option<DiffBookDepthUpdateSpeedEnum>,
658}
659
660impl DiffBookDepthParams {
661 #[must_use]
668 pub fn builder(symbol: String) -> DiffBookDepthParamsBuilder {
669 DiffBookDepthParamsBuilder::default().symbol(symbol)
670 }
671}
672#[derive(Clone, Debug, Builder, Deserialize)]
677#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
678pub struct KlineParams {
679 #[builder(setter(into))]
683 #[serde(rename = "symbol")]
684 pub symbol: String,
685 #[builder(setter(into))]
690 #[serde(rename = "interval")]
691 pub interval: KlineIntervalEnum,
692 #[builder(setter(into), default)]
696 #[serde(rename = "id", default)]
697 pub id: Option<String>,
698}
699
700impl KlineParams {
701 #[must_use]
709 pub fn builder(symbol: String, interval: KlineIntervalEnum) -> KlineParamsBuilder {
710 KlineParamsBuilder::default()
711 .symbol(symbol)
712 .interval(interval)
713 }
714}
715#[derive(Clone, Debug, Builder, Deserialize)]
720#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
721pub struct KlineOffsetParams {
722 #[builder(setter(into))]
726 #[serde(rename = "symbol")]
727 pub symbol: String,
728 #[builder(setter(into))]
733 #[serde(rename = "interval")]
734 pub interval: KlineOffsetIntervalEnum,
735 #[builder(setter(into), default)]
739 #[serde(rename = "id", default)]
740 pub id: Option<String>,
741}
742
743impl KlineOffsetParams {
744 #[must_use]
752 pub fn builder(symbol: String, interval: KlineOffsetIntervalEnum) -> KlineOffsetParamsBuilder {
753 KlineOffsetParamsBuilder::default()
754 .symbol(symbol)
755 .interval(interval)
756 }
757}
758#[derive(Clone, Debug, Builder, Deserialize)]
763#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
764pub struct MiniTickerParams {
765 #[builder(setter(into))]
769 #[serde(rename = "symbol")]
770 pub symbol: String,
771 #[builder(setter(into), default)]
775 #[serde(rename = "id", default)]
776 pub id: Option<String>,
777}
778
779impl MiniTickerParams {
780 #[must_use]
787 pub fn builder(symbol: String) -> MiniTickerParamsBuilder {
788 MiniTickerParamsBuilder::default().symbol(symbol)
789 }
790}
791#[derive(Clone, Debug, Builder, Deserialize)]
796#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
797pub struct PartialBookDepthParams {
798 #[builder(setter(into))]
802 #[serde(rename = "symbol")]
803 pub symbol: String,
804 #[builder(setter(into))]
809 #[serde(rename = "levels")]
810 pub levels: PartialBookDepthLevelsEnum,
811 #[builder(setter(into), default)]
815 #[serde(rename = "id", default)]
816 pub id: Option<String>,
817 #[builder(setter(into), default)]
821 #[serde(rename = "updateSpeed", default)]
822 pub update_speed: Option<PartialBookDepthUpdateSpeedEnum>,
823}
824
825impl PartialBookDepthParams {
826 #[must_use]
834 pub fn builder(
835 symbol: String,
836 levels: PartialBookDepthLevelsEnum,
837 ) -> PartialBookDepthParamsBuilder {
838 PartialBookDepthParamsBuilder::default()
839 .symbol(symbol)
840 .levels(levels)
841 }
842}
843#[derive(Clone, Debug, Builder, Deserialize)]
848#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
849pub struct ReferencePriceParams {
850 #[builder(setter(into))]
854 #[serde(rename = "symbol")]
855 pub symbol: String,
856 #[builder(setter(into), default)]
860 #[serde(rename = "id", default)]
861 pub id: Option<String>,
862}
863
864impl ReferencePriceParams {
865 #[must_use]
872 pub fn builder(symbol: String) -> ReferencePriceParamsBuilder {
873 ReferencePriceParamsBuilder::default().symbol(symbol)
874 }
875}
876#[derive(Clone, Debug, Builder, Deserialize)]
881#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
882pub struct RollingWindowTickerParams {
883 #[builder(setter(into))]
887 #[serde(rename = "symbol")]
888 pub symbol: String,
889 #[builder(setter(into))]
894 #[serde(rename = "windowSize")]
895 pub window_size: RollingWindowTickerWindowSizeEnum,
896 #[builder(setter(into), default)]
900 #[serde(rename = "id", default)]
901 pub id: Option<String>,
902}
903
904impl RollingWindowTickerParams {
905 #[must_use]
913 pub fn builder(
914 symbol: String,
915 window_size: RollingWindowTickerWindowSizeEnum,
916 ) -> RollingWindowTickerParamsBuilder {
917 RollingWindowTickerParamsBuilder::default()
918 .symbol(symbol)
919 .window_size(window_size)
920 }
921}
922#[derive(Clone, Debug, Builder, Deserialize)]
927#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
928pub struct TickerParams {
929 #[builder(setter(into))]
933 #[serde(rename = "symbol")]
934 pub symbol: String,
935 #[builder(setter(into), default)]
939 #[serde(rename = "id", default)]
940 pub id: Option<String>,
941}
942
943impl TickerParams {
944 #[must_use]
951 pub fn builder(symbol: String) -> TickerParamsBuilder {
952 TickerParamsBuilder::default().symbol(symbol)
953 }
954}
955#[derive(Clone, Debug, Builder, Deserialize)]
960#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
961pub struct TradeParams {
962 #[builder(setter(into))]
966 #[serde(rename = "symbol")]
967 pub symbol: String,
968 #[builder(setter(into), default)]
972 #[serde(rename = "id", default)]
973 pub id: Option<String>,
974}
975
976impl TradeParams {
977 #[must_use]
984 pub fn builder(symbol: String) -> TradeParamsBuilder {
985 TradeParamsBuilder::default().symbol(symbol)
986 }
987}
988
989#[async_trait]
990impl Api for ApiClient {
991 async fn agg_trade(
992 &self,
993 params: AggTradeParams,
994 ) -> anyhow::Result<Arc<WebsocketStream<models::AggTradeResponse>>> {
995 let AggTradeParams { symbol, id } = params;
996
997 let pairs: &[(&str, Option<String>)] =
998 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
999
1000 let vars: HashMap<_, _> = pairs
1001 .iter()
1002 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1003 .collect();
1004
1005 let id_opt: Option<String> = vars.get("id").cloned();
1006
1007 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1008
1009 Ok(create_stream_handler::<models::AggTradeResponse>(
1010 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1011 stream,
1012 id_opt.map(|s| {
1013 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1014 if let Ok(n) = s.parse::<u32>() {
1015 return StreamId::Number(n);
1016 }
1017 }
1018 StreamId::Str(s)
1019 }),
1020 None,
1021 )
1022 .await)
1023 }
1024
1025 async fn all_market_rolling_window_ticker(
1026 &self,
1027 params: AllMarketRollingWindowTickerParams,
1028 ) -> anyhow::Result<Arc<WebsocketStream<Vec<models::AllMarketRollingWindowTickerResponseInner>>>>
1029 {
1030 let AllMarketRollingWindowTickerParams { window_size, id } = params;
1031
1032 let pairs: &[(&str, Option<String>)] = &[
1033 ("windowSize", Some(window_size.as_str().to_string())),
1034 ("id", id.clone()),
1035 ];
1036
1037 let vars: HashMap<_, _> = pairs
1038 .iter()
1039 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1040 .collect();
1041
1042 let id_opt: Option<String> = vars.get("id").cloned();
1043
1044 let stream = replace_websocket_streams_placeholders("/!ticker_<windowSize>@arr", &vars);
1045
1046 Ok(
1047 create_stream_handler::<Vec<models::AllMarketRollingWindowTickerResponseInner>>(
1048 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1049 stream,
1050 id_opt.map(|s| {
1051 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1052 if let Ok(n) = s.parse::<u32>() {
1053 return StreamId::Number(n);
1054 }
1055 }
1056 StreamId::Str(s)
1057 }),
1058 None,
1059 )
1060 .await,
1061 )
1062 }
1063
1064 async fn all_mini_ticker(
1065 &self,
1066 params: AllMiniTickerParams,
1067 ) -> anyhow::Result<Arc<WebsocketStream<Vec<models::AllMiniTickerResponseInner>>>> {
1068 let AllMiniTickerParams { id } = params;
1069
1070 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1071
1072 let vars: HashMap<_, _> = pairs
1073 .iter()
1074 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1075 .collect();
1076
1077 let id_opt: Option<String> = vars.get("id").cloned();
1078
1079 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1080
1081 Ok(
1082 create_stream_handler::<Vec<models::AllMiniTickerResponseInner>>(
1083 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1084 stream,
1085 id_opt.map(|s| {
1086 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1087 if let Ok(n) = s.parse::<u32>() {
1088 return StreamId::Number(n);
1089 }
1090 }
1091 StreamId::Str(s)
1092 }),
1093 None,
1094 )
1095 .await,
1096 )
1097 }
1098
1099 async fn avg_price(
1100 &self,
1101 params: AvgPriceParams,
1102 ) -> anyhow::Result<Arc<WebsocketStream<models::AvgPriceResponse>>> {
1103 let AvgPriceParams { symbol, id } = params;
1104
1105 let pairs: &[(&str, Option<String>)] =
1106 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1107
1108 let vars: HashMap<_, _> = pairs
1109 .iter()
1110 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1111 .collect();
1112
1113 let id_opt: Option<String> = vars.get("id").cloned();
1114
1115 let stream = replace_websocket_streams_placeholders("/<symbol>@avgPrice", &vars);
1116
1117 Ok(create_stream_handler::<models::AvgPriceResponse>(
1118 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1119 stream,
1120 id_opt.map(|s| {
1121 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1122 if let Ok(n) = s.parse::<u32>() {
1123 return StreamId::Number(n);
1124 }
1125 }
1126 StreamId::Str(s)
1127 }),
1128 None,
1129 )
1130 .await)
1131 }
1132
1133 async fn block_trade(
1134 &self,
1135 params: BlockTradeParams,
1136 ) -> anyhow::Result<Arc<WebsocketStream<models::BlockTradeResponse>>> {
1137 let BlockTradeParams { symbol, id } = params;
1138
1139 let pairs: &[(&str, Option<String>)] =
1140 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1141
1142 let vars: HashMap<_, _> = pairs
1143 .iter()
1144 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1145 .collect();
1146
1147 let id_opt: Option<String> = vars.get("id").cloned();
1148
1149 let stream = replace_websocket_streams_placeholders("/<symbol>@blockTrade", &vars);
1150
1151 Ok(create_stream_handler::<models::BlockTradeResponse>(
1152 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1153 stream,
1154 id_opt.map(|s| {
1155 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1156 if let Ok(n) = s.parse::<u32>() {
1157 return StreamId::Number(n);
1158 }
1159 }
1160 StreamId::Str(s)
1161 }),
1162 None,
1163 )
1164 .await)
1165 }
1166
1167 async fn book_ticker(
1168 &self,
1169 params: BookTickerParams,
1170 ) -> anyhow::Result<Arc<WebsocketStream<models::BookTickerResponse>>> {
1171 let BookTickerParams { symbol, id } = params;
1172
1173 let pairs: &[(&str, Option<String>)] =
1174 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1175
1176 let vars: HashMap<_, _> = pairs
1177 .iter()
1178 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1179 .collect();
1180
1181 let id_opt: Option<String> = vars.get("id").cloned();
1182
1183 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
1184
1185 Ok(create_stream_handler::<models::BookTickerResponse>(
1186 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1187 stream,
1188 id_opt.map(|s| {
1189 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1190 if let Ok(n) = s.parse::<u32>() {
1191 return StreamId::Number(n);
1192 }
1193 }
1194 StreamId::Str(s)
1195 }),
1196 None,
1197 )
1198 .await)
1199 }
1200
1201 async fn diff_book_depth(
1202 &self,
1203 params: DiffBookDepthParams,
1204 ) -> anyhow::Result<Arc<WebsocketStream<models::DiffBookDepthResponse>>> {
1205 let DiffBookDepthParams {
1206 symbol,
1207 id,
1208 update_speed,
1209 } = params;
1210
1211 let pairs: &[(&str, Option<String>)] = &[
1212 ("symbol", Some(symbol.clone())),
1213 ("id", id.clone()),
1214 (
1215 "updateSpeed",
1216 update_speed.clone().map(|v| v.as_str().to_string()),
1217 ),
1218 ];
1219
1220 let vars: HashMap<_, _> = pairs
1221 .iter()
1222 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1223 .collect();
1224
1225 let id_opt: Option<String> = vars.get("id").cloned();
1226
1227 let stream = replace_websocket_streams_placeholders("/<symbol>@depth@<updateSpeed>", &vars);
1228
1229 Ok(create_stream_handler::<models::DiffBookDepthResponse>(
1230 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1231 stream,
1232 id_opt.map(|s| {
1233 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1234 if let Ok(n) = s.parse::<u32>() {
1235 return StreamId::Number(n);
1236 }
1237 }
1238 StreamId::Str(s)
1239 }),
1240 None,
1241 )
1242 .await)
1243 }
1244
1245 async fn kline(
1246 &self,
1247 params: KlineParams,
1248 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineResponse>>> {
1249 let KlineParams {
1250 symbol,
1251 interval,
1252 id,
1253 } = params;
1254
1255 let pairs: &[(&str, Option<String>)] = &[
1256 ("symbol", Some(symbol.clone())),
1257 ("interval", Some(interval.as_str().to_string())),
1258 ("id", id.clone()),
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>@kline_<interval>", &vars);
1269
1270 Ok(create_stream_handler::<models::KlineResponse>(
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 async fn kline_offset(
1287 &self,
1288 params: KlineOffsetParams,
1289 ) -> anyhow::Result<Arc<WebsocketStream<models::KlineOffsetResponse>>> {
1290 let KlineOffsetParams {
1291 symbol,
1292 interval,
1293 id,
1294 } = params;
1295
1296 let pairs: &[(&str, Option<String>)] = &[
1297 ("symbol", Some(symbol.clone())),
1298 ("interval", Some(interval.as_str().to_string())),
1299 ("id", id.clone()),
1300 ];
1301
1302 let vars: HashMap<_, _> = pairs
1303 .iter()
1304 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1305 .collect();
1306
1307 let id_opt: Option<String> = vars.get("id").cloned();
1308
1309 let stream =
1310 replace_websocket_streams_placeholders("/<symbol>@kline_<interval>@+08:00", &vars);
1311
1312 Ok(create_stream_handler::<models::KlineOffsetResponse>(
1313 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1314 stream,
1315 id_opt.map(|s| {
1316 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1317 if let Ok(n) = s.parse::<u32>() {
1318 return StreamId::Number(n);
1319 }
1320 }
1321 StreamId::Str(s)
1322 }),
1323 None,
1324 )
1325 .await)
1326 }
1327
1328 async fn mini_ticker(
1329 &self,
1330 params: MiniTickerParams,
1331 ) -> anyhow::Result<Arc<WebsocketStream<models::MiniTickerResponse>>> {
1332 let MiniTickerParams { symbol, id } = params;
1333
1334 let pairs: &[(&str, Option<String>)] =
1335 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1336
1337 let vars: HashMap<_, _> = pairs
1338 .iter()
1339 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1340 .collect();
1341
1342 let id_opt: Option<String> = vars.get("id").cloned();
1343
1344 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
1345
1346 Ok(create_stream_handler::<models::MiniTickerResponse>(
1347 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1348 stream,
1349 id_opt.map(|s| {
1350 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1351 if let Ok(n) = s.parse::<u32>() {
1352 return StreamId::Number(n);
1353 }
1354 }
1355 StreamId::Str(s)
1356 }),
1357 None,
1358 )
1359 .await)
1360 }
1361
1362 async fn partial_book_depth(
1363 &self,
1364 params: PartialBookDepthParams,
1365 ) -> anyhow::Result<Arc<WebsocketStream<models::PartialBookDepthResponse>>> {
1366 let PartialBookDepthParams {
1367 symbol,
1368 levels,
1369 id,
1370 update_speed,
1371 } = params;
1372
1373 let pairs: &[(&str, Option<String>)] = &[
1374 ("symbol", Some(symbol.clone())),
1375 ("levels", Some(levels.as_str().to_string())),
1376 ("id", id.clone()),
1377 (
1378 "updateSpeed",
1379 update_speed.clone().map(|v| v.as_str().to_string()),
1380 ),
1381 ];
1382
1383 let vars: HashMap<_, _> = pairs
1384 .iter()
1385 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1386 .collect();
1387
1388 let id_opt: Option<String> = vars.get("id").cloned();
1389
1390 let stream =
1391 replace_websocket_streams_placeholders("/<symbol>@depth<levels>@<updateSpeed>", &vars);
1392
1393 Ok(create_stream_handler::<models::PartialBookDepthResponse>(
1394 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1395 stream,
1396 id_opt.map(|s| {
1397 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1398 if let Ok(n) = s.parse::<u32>() {
1399 return StreamId::Number(n);
1400 }
1401 }
1402 StreamId::Str(s)
1403 }),
1404 None,
1405 )
1406 .await)
1407 }
1408
1409 async fn reference_price(
1410 &self,
1411 params: ReferencePriceParams,
1412 ) -> anyhow::Result<Arc<WebsocketStream<models::ReferencePriceResponse>>> {
1413 let ReferencePriceParams { symbol, id } = params;
1414
1415 let pairs: &[(&str, Option<String>)] =
1416 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1417
1418 let vars: HashMap<_, _> = pairs
1419 .iter()
1420 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1421 .collect();
1422
1423 let id_opt: Option<String> = vars.get("id").cloned();
1424
1425 let stream = replace_websocket_streams_placeholders("/<symbol>@referencePrice", &vars);
1426
1427 Ok(create_stream_handler::<models::ReferencePriceResponse>(
1428 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1429 stream,
1430 id_opt.map(|s| {
1431 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1432 if let Ok(n) = s.parse::<u32>() {
1433 return StreamId::Number(n);
1434 }
1435 }
1436 StreamId::Str(s)
1437 }),
1438 None,
1439 )
1440 .await)
1441 }
1442
1443 async fn rolling_window_ticker(
1444 &self,
1445 params: RollingWindowTickerParams,
1446 ) -> anyhow::Result<Arc<WebsocketStream<models::RollingWindowTickerResponse>>> {
1447 let RollingWindowTickerParams {
1448 symbol,
1449 window_size,
1450 id,
1451 } = params;
1452
1453 let pairs: &[(&str, Option<String>)] = &[
1454 ("symbol", Some(symbol.clone())),
1455 ("windowSize", Some(window_size.as_str().to_string())),
1456 ("id", id.clone()),
1457 ];
1458
1459 let vars: HashMap<_, _> = pairs
1460 .iter()
1461 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1462 .collect();
1463
1464 let id_opt: Option<String> = vars.get("id").cloned();
1465
1466 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker_<windowSize>", &vars);
1467
1468 Ok(
1469 create_stream_handler::<models::RollingWindowTickerResponse>(
1470 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1471 stream,
1472 id_opt.map(|s| {
1473 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1474 if let Ok(n) = s.parse::<u32>() {
1475 return StreamId::Number(n);
1476 }
1477 }
1478 StreamId::Str(s)
1479 }),
1480 None,
1481 )
1482 .await,
1483 )
1484 }
1485
1486 async fn ticker(
1487 &self,
1488 params: TickerParams,
1489 ) -> anyhow::Result<Arc<WebsocketStream<models::TickerResponse>>> {
1490 let TickerParams { symbol, id } = params;
1491
1492 let pairs: &[(&str, Option<String>)] =
1493 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1494
1495 let vars: HashMap<_, _> = pairs
1496 .iter()
1497 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1498 .collect();
1499
1500 let id_opt: Option<String> = vars.get("id").cloned();
1501
1502 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
1503
1504 Ok(create_stream_handler::<models::TickerResponse>(
1505 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1506 stream,
1507 id_opt.map(|s| {
1508 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1509 if let Ok(n) = s.parse::<u32>() {
1510 return StreamId::Number(n);
1511 }
1512 }
1513 StreamId::Str(s)
1514 }),
1515 None,
1516 )
1517 .await)
1518 }
1519
1520 async fn trade(
1521 &self,
1522 params: TradeParams,
1523 ) -> anyhow::Result<Arc<WebsocketStream<models::TradeResponse>>> {
1524 let TradeParams { symbol, id } = params;
1525
1526 let pairs: &[(&str, Option<String>)] =
1527 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1528
1529 let vars: HashMap<_, _> = pairs
1530 .iter()
1531 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1532 .collect();
1533
1534 let id_opt: Option<String> = vars.get("id").cloned();
1535
1536 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
1537
1538 Ok(create_stream_handler::<models::TradeResponse>(
1539 WebsocketBase::WebsocketStreams(Arc::clone(&self.websocket_streams_base)),
1540 stream,
1541 id_opt.map(|s| {
1542 if !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()) {
1543 if let Ok(n) = s.parse::<u32>() {
1544 return StreamId::Number(n);
1545 }
1546 }
1547 StreamId::Str(s)
1548 }),
1549 None,
1550 )
1551 .await)
1552 }
1553}
1554
1555#[cfg(all(test, feature = "spot"))]
1556mod tests {
1557 use super::*;
1558 use crate::TOKIO_SHARED_RT;
1559 use crate::{
1560 common::websocket::{WebsocketConnection, WebsocketHandler},
1561 config::ConfigurationWebsocketStreams,
1562 };
1563 use serde_json::json;
1564 use std::sync::atomic::{AtomicBool, Ordering};
1565 use tokio::task::yield_now;
1566
1567 async fn make_streams_base() -> (Arc<WebsocketStreams>, Arc<WebsocketConnection>) {
1568 let conn = WebsocketConnection::new("test");
1569 let config = ConfigurationWebsocketStreams::builder()
1570 .build()
1571 .expect("Failed to build configuration");
1572 let streams_base = WebsocketStreams::new(config, vec![conn.clone()], vec![]);
1573 conn.set_handler(streams_base.clone() as Arc<dyn WebsocketHandler>)
1574 .await;
1575 (streams_base, conn)
1576 }
1577
1578 #[test]
1579 fn agg_trade_should_execute_successfully() {
1580 TOKIO_SHARED_RT.block_on(async {
1581 let (streams_base, _) = make_streams_base().await;
1582 let api = ApiClient::new(streams_base.clone());
1583
1584 let id = "test-id-123".to_string();
1585
1586 let params = AggTradeParams::builder("bnbusdt".to_string())
1587 .id(Some(id.clone()))
1588 .build()
1589 .unwrap();
1590
1591 let AggTradeParams { symbol, id } = params.clone();
1592
1593 let pairs: &[(&str, Option<String>)] =
1594 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
1595
1596 let vars: HashMap<_, _> = pairs
1597 .iter()
1598 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1599 .collect();
1600 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1601 let ws_stream = api
1602 .agg_trade(params)
1603 .await
1604 .expect("agg_trade should return a WebsocketStream");
1605
1606 assert!(
1607 streams_base.is_subscribed(&stream).await,
1608 "expected stream '{stream}' to be subscribed"
1609 );
1610 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1611 });
1612 }
1613
1614 #[test]
1615 fn agg_trade_should_handle_incoming_message() {
1616 TOKIO_SHARED_RT.block_on(async {
1617 let (streams_base, conn) = make_streams_base().await;
1618 let api = ApiClient::new(streams_base.clone());
1619
1620 let id = "test-id-123".to_string();
1621
1622 let params = AggTradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
1623
1624 let AggTradeParams {
1625 symbol,id,
1626 } = params.clone();
1627
1628 let pairs: &[(&str, Option<String>)] = &[
1629 ("symbol",
1630 Some(symbol.clone())
1631 ),
1632 ("id",
1633 id.clone()
1634 ),
1635 ];
1636
1637 let vars: HashMap<_, _> = pairs
1638 .iter()
1639 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1640 .collect();
1641 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1642
1643 let ws_stream = api.agg_trade(params).await.unwrap();
1644
1645 let called = Arc::new(AtomicBool::new(false));
1646 let called_with_message = called.clone();
1647 ws_stream.on_message(move |_payload: models::AggTradeResponse| {
1648 called_with_message.store(true, Ordering::SeqCst);
1649 });
1650
1651 let payload: Value = serde_json::from_str(r#"{"e":"aggTrade","E":1672515782136,"s":"BNBBTC","a":12345,"f":100,"l":105,"T":1672515782136,"m":true,"M":true}"#).unwrap_or_else(|_| serde_json::json!({}));
1652 let msg = json!({
1653 "stream": stream,
1654 "data": payload,
1655 });
1656
1657 streams_base.on_message(msg.to_string(), conn.clone()).await;
1658 yield_now().await;
1659
1660 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1661 });
1662 }
1663
1664 #[test]
1665 fn agg_trade_should_not_fire_after_unsubscribe() {
1666 TOKIO_SHARED_RT.block_on(async {
1667 let (streams_base, conn) = make_streams_base().await;
1668 let api = ApiClient::new(streams_base.clone());
1669
1670 let id = "test-id-123".to_string();
1671
1672 let params = AggTradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
1673
1674 let AggTradeParams {
1675 symbol,id,
1676 } = params.clone();
1677
1678 let pairs: &[(&str, Option<String>)] = &[
1679 ("symbol",
1680 Some(symbol.clone())
1681 ),
1682 ("id",
1683 id.clone()
1684 ),
1685 ];
1686
1687 let vars: HashMap<_, _> = pairs
1688 .iter()
1689 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1690 .collect();
1691 let stream = replace_websocket_streams_placeholders("/<symbol>@aggTrade", &vars);
1692
1693 let ws_stream = api.agg_trade(params).await.unwrap();
1694
1695 let called = Arc::new(AtomicBool::new(false));
1696 let called_clone = called.clone();
1697 ws_stream.on_message(move |_payload: models::AggTradeResponse| {
1698 called_clone.store(true, Ordering::SeqCst);
1699 });
1700
1701 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1702
1703 ws_stream.unsubscribe().await;
1704
1705 let payload: Value = serde_json::from_str(r#"{"e":"aggTrade","E":1672515782136,"s":"BNBBTC","a":12345,"f":100,"l":105,"T":1672515782136,"m":true,"M":true}"#).unwrap_or_else(|_| serde_json::json!({}));
1706 let msg = json!({
1707 "stream": stream,
1708 "data": payload,
1709 });
1710
1711 streams_base.on_message(msg.to_string(), conn.clone()).await;
1712
1713 yield_now().await;
1714
1715 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1716 });
1717 }
1718
1719 #[test]
1720 fn all_market_rolling_window_ticker_should_execute_successfully() {
1721 TOKIO_SHARED_RT.block_on(async {
1722 let (streams_base, _) = make_streams_base().await;
1723 let api = ApiClient::new(streams_base.clone());
1724
1725 let id = "test-id-123".to_string();
1726
1727 let params = AllMarketRollingWindowTickerParams::builder(
1728 AllMarketRollingWindowTickerWindowSizeEnum::WindowSize1h,
1729 )
1730 .id(Some(id.clone()))
1731 .build()
1732 .unwrap();
1733
1734 let AllMarketRollingWindowTickerParams { window_size, id } = params.clone();
1735
1736 let pairs: &[(&str, Option<String>)] = &[
1737 ("windowSize", Some(window_size.as_str().to_string())),
1738 ("id", id.clone()),
1739 ];
1740
1741 let vars: HashMap<_, _> = pairs
1742 .iter()
1743 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1744 .collect();
1745 let stream = replace_websocket_streams_placeholders("/!ticker_<windowSize>@arr", &vars);
1746 let ws_stream = api
1747 .all_market_rolling_window_ticker(params)
1748 .await
1749 .expect("all_market_rolling_window_ticker should return a WebsocketStream");
1750
1751 assert!(
1752 streams_base.is_subscribed(&stream).await,
1753 "expected stream '{stream}' to be subscribed"
1754 );
1755 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1756 });
1757 }
1758
1759 #[test]
1760 fn all_market_rolling_window_ticker_should_handle_incoming_message() {
1761 TOKIO_SHARED_RT.block_on(async {
1762 let (streams_base, conn) = make_streams_base().await;
1763 let api = ApiClient::new(streams_base.clone());
1764
1765 let id = "test-id-123".to_string();
1766
1767 let params = AllMarketRollingWindowTickerParams::builder(AllMarketRollingWindowTickerWindowSizeEnum::WindowSize1h,).id(Some(id.clone())).build().unwrap();
1768
1769 let AllMarketRollingWindowTickerParams {
1770 window_size,id,
1771 } = params.clone();
1772
1773 let pairs: &[(&str, Option<String>)] = &[
1774 ("windowSize",
1775 Some(window_size.as_str().to_string())
1776 ),
1777 ("id",
1778 id.clone()
1779 ),
1780 ];
1781
1782 let vars: HashMap<_, _> = pairs
1783 .iter()
1784 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1785 .collect();
1786 let stream = replace_websocket_streams_placeholders("/!ticker_<windowSize>@arr", &vars);
1787
1788 let ws_stream = api.all_market_rolling_window_ticker(params).await.unwrap();
1789
1790 let called = Arc::new(AtomicBool::new(false));
1791 let called_with_message = called.clone();
1792 ws_stream.on_message(move |_payload: Vec<models::AllMarketRollingWindowTickerResponseInner>| {
1793 called_with_message.store(true, Ordering::SeqCst);
1794 });
1795
1796 let payload: Value = serde_json::from_str(r#"[{"e":"1hTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}]"#).unwrap_or_else(|_| serde_json::json!({}));
1797 let msg = json!({
1798 "stream": stream,
1799 "data": payload,
1800 });
1801
1802 streams_base.on_message(msg.to_string(), conn.clone()).await;
1803 yield_now().await;
1804
1805 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
1806 });
1807 }
1808
1809 #[test]
1810 fn all_market_rolling_window_ticker_should_not_fire_after_unsubscribe() {
1811 TOKIO_SHARED_RT.block_on(async {
1812 let (streams_base, conn) = make_streams_base().await;
1813 let api = ApiClient::new(streams_base.clone());
1814
1815 let id = "test-id-123".to_string();
1816
1817 let params = AllMarketRollingWindowTickerParams::builder(AllMarketRollingWindowTickerWindowSizeEnum::WindowSize1h,).id(Some(id.clone())).build().unwrap();
1818
1819 let AllMarketRollingWindowTickerParams {
1820 window_size,id,
1821 } = params.clone();
1822
1823 let pairs: &[(&str, Option<String>)] = &[
1824 ("windowSize",
1825 Some(window_size.as_str().to_string())
1826 ),
1827 ("id",
1828 id.clone()
1829 ),
1830 ];
1831
1832 let vars: HashMap<_, _> = pairs
1833 .iter()
1834 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1835 .collect();
1836 let stream = replace_websocket_streams_placeholders("/!ticker_<windowSize>@arr", &vars);
1837
1838 let ws_stream = api.all_market_rolling_window_ticker(params).await.unwrap();
1839
1840 let called = Arc::new(AtomicBool::new(false));
1841 let called_clone = called.clone();
1842 ws_stream.on_message(move |_payload: Vec<models::AllMarketRollingWindowTickerResponseInner>| {
1843 called_clone.store(true, Ordering::SeqCst);
1844 });
1845
1846 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
1847
1848 ws_stream.unsubscribe().await;
1849
1850 let payload: Value = serde_json::from_str(r#"[{"e":"1hTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}]"#).unwrap_or_else(|_| serde_json::json!({}));
1851 let msg = json!({
1852 "stream": stream,
1853 "data": payload,
1854 });
1855
1856 streams_base.on_message(msg.to_string(), conn.clone()).await;
1857
1858 yield_now().await;
1859
1860 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
1861 });
1862 }
1863
1864 #[test]
1865 fn all_mini_ticker_should_execute_successfully() {
1866 TOKIO_SHARED_RT.block_on(async {
1867 let (streams_base, _) = make_streams_base().await;
1868 let api = ApiClient::new(streams_base.clone());
1869
1870 let id = "test-id-123".to_string();
1871
1872 let params = AllMiniTickerParams::builder()
1873 .id(Some(id.clone()))
1874 .build()
1875 .unwrap();
1876
1877 let AllMiniTickerParams { id } = params.clone();
1878
1879 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1880
1881 let vars: HashMap<_, _> = pairs
1882 .iter()
1883 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1884 .collect();
1885 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1886 let ws_stream = api
1887 .all_mini_ticker(params)
1888 .await
1889 .expect("all_mini_ticker should return a WebsocketStream");
1890
1891 assert!(
1892 streams_base.is_subscribed(&stream).await,
1893 "expected stream '{stream}' to be subscribed"
1894 );
1895 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
1896 });
1897 }
1898
1899 #[test]
1900 fn all_mini_ticker_should_handle_incoming_message() {
1901 TOKIO_SHARED_RT.block_on(async {
1902 let (streams_base, conn) = make_streams_base().await;
1903 let api = ApiClient::new(streams_base.clone());
1904
1905 let id = "test-id-123".to_string();
1906
1907 let params = AllMiniTickerParams::builder()
1908 .id(Some(id.clone()))
1909 .build()
1910 .unwrap();
1911
1912 let AllMiniTickerParams { id } = params.clone();
1913
1914 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1915
1916 let vars: HashMap<_, _> = pairs
1917 .iter()
1918 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1919 .collect();
1920 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1921
1922 let ws_stream = api.all_mini_ticker(params).await.unwrap();
1923
1924 let called = Arc::new(AtomicBool::new(false));
1925 let called_with_message = called.clone();
1926 ws_stream.on_message(move |_payload: Vec<models::AllMiniTickerResponseInner>| {
1927 called_with_message.store(true, Ordering::SeqCst);
1928 });
1929
1930 let payload: Value =
1931 serde_json::from_str(r#"[{"e":"24hrMiniTicker","E":1672515782136,"s":"BNBBTC"}]"#)
1932 .unwrap_or_else(|_| serde_json::json!({}));
1933 let msg = json!({
1934 "stream": stream,
1935 "data": payload,
1936 });
1937
1938 streams_base.on_message(msg.to_string(), conn.clone()).await;
1939 yield_now().await;
1940
1941 assert!(
1942 called.load(Ordering::SeqCst),
1943 "expected our callback to have been invoked"
1944 );
1945 });
1946 }
1947
1948 #[test]
1949 fn all_mini_ticker_should_not_fire_after_unsubscribe() {
1950 TOKIO_SHARED_RT.block_on(async {
1951 let (streams_base, conn) = make_streams_base().await;
1952 let api = ApiClient::new(streams_base.clone());
1953
1954 let id = "test-id-123".to_string();
1955
1956 let params = AllMiniTickerParams::builder()
1957 .id(Some(id.clone()))
1958 .build()
1959 .unwrap();
1960
1961 let AllMiniTickerParams { id } = params.clone();
1962
1963 let pairs: &[(&str, Option<String>)] = &[("id", id.clone())];
1964
1965 let vars: HashMap<_, _> = pairs
1966 .iter()
1967 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
1968 .collect();
1969 let stream = replace_websocket_streams_placeholders("/!miniTicker@arr", &vars);
1970
1971 let ws_stream = api.all_mini_ticker(params).await.unwrap();
1972
1973 let called = Arc::new(AtomicBool::new(false));
1974 let called_clone = called.clone();
1975 ws_stream.on_message(move |_payload: Vec<models::AllMiniTickerResponseInner>| {
1976 called_clone.store(true, Ordering::SeqCst);
1977 });
1978
1979 assert!(
1980 streams_base.is_subscribed(&stream).await,
1981 "should be subscribed before unsubscribe"
1982 );
1983
1984 ws_stream.unsubscribe().await;
1985
1986 let payload: Value =
1987 serde_json::from_str(r#"[{"e":"24hrMiniTicker","E":1672515782136,"s":"BNBBTC"}]"#)
1988 .unwrap_or_else(|_| serde_json::json!({}));
1989 let msg = json!({
1990 "stream": stream,
1991 "data": payload,
1992 });
1993
1994 streams_base.on_message(msg.to_string(), conn.clone()).await;
1995
1996 yield_now().await;
1997
1998 assert!(
1999 !called.load(Ordering::SeqCst),
2000 "callback should not be invoked after unsubscribe"
2001 );
2002 });
2003 }
2004
2005 #[test]
2006 fn avg_price_should_execute_successfully() {
2007 TOKIO_SHARED_RT.block_on(async {
2008 let (streams_base, _) = make_streams_base().await;
2009 let api = ApiClient::new(streams_base.clone());
2010
2011 let id = "test-id-123".to_string();
2012
2013 let params = AvgPriceParams::builder("bnbusdt".to_string())
2014 .id(Some(id.clone()))
2015 .build()
2016 .unwrap();
2017
2018 let AvgPriceParams { symbol, id } = params.clone();
2019
2020 let pairs: &[(&str, Option<String>)] =
2021 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2022
2023 let vars: HashMap<_, _> = pairs
2024 .iter()
2025 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2026 .collect();
2027 let stream = replace_websocket_streams_placeholders("/<symbol>@avgPrice", &vars);
2028 let ws_stream = api
2029 .avg_price(params)
2030 .await
2031 .expect("avg_price should return a WebsocketStream");
2032
2033 assert!(
2034 streams_base.is_subscribed(&stream).await,
2035 "expected stream '{stream}' to be subscribed"
2036 );
2037 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2038 });
2039 }
2040
2041 #[test]
2042 fn avg_price_should_handle_incoming_message() {
2043 TOKIO_SHARED_RT.block_on(async {
2044 let (streams_base, conn) = make_streams_base().await;
2045 let api = ApiClient::new(streams_base.clone());
2046
2047 let id = "test-id-123".to_string();
2048
2049 let params = AvgPriceParams::builder("bnbusdt".to_string())
2050 .id(Some(id.clone()))
2051 .build()
2052 .unwrap();
2053
2054 let AvgPriceParams { symbol, id } = params.clone();
2055
2056 let pairs: &[(&str, Option<String>)] =
2057 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2058
2059 let vars: HashMap<_, _> = pairs
2060 .iter()
2061 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2062 .collect();
2063 let stream = replace_websocket_streams_placeholders("/<symbol>@avgPrice", &vars);
2064
2065 let ws_stream = api.avg_price(params).await.unwrap();
2066
2067 let called = Arc::new(AtomicBool::new(false));
2068 let called_with_message = called.clone();
2069 ws_stream.on_message(move |_payload: models::AvgPriceResponse| {
2070 called_with_message.store(true, Ordering::SeqCst);
2071 });
2072
2073 let payload: Value = serde_json::from_str(
2074 r#"{"e":"avgPrice","E":1693907033000,"s":"BTCUSDT","i":"5m","T":1693907032213}"#,
2075 )
2076 .unwrap_or_else(|_| serde_json::json!({}));
2077 let msg = json!({
2078 "stream": stream,
2079 "data": payload,
2080 });
2081
2082 streams_base.on_message(msg.to_string(), conn.clone()).await;
2083 yield_now().await;
2084
2085 assert!(
2086 called.load(Ordering::SeqCst),
2087 "expected our callback to have been invoked"
2088 );
2089 });
2090 }
2091
2092 #[test]
2093 fn avg_price_should_not_fire_after_unsubscribe() {
2094 TOKIO_SHARED_RT.block_on(async {
2095 let (streams_base, conn) = make_streams_base().await;
2096 let api = ApiClient::new(streams_base.clone());
2097
2098 let id = "test-id-123".to_string();
2099
2100 let params = AvgPriceParams::builder("bnbusdt".to_string())
2101 .id(Some(id.clone()))
2102 .build()
2103 .unwrap();
2104
2105 let AvgPriceParams { symbol, id } = params.clone();
2106
2107 let pairs: &[(&str, Option<String>)] =
2108 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2109
2110 let vars: HashMap<_, _> = pairs
2111 .iter()
2112 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2113 .collect();
2114 let stream = replace_websocket_streams_placeholders("/<symbol>@avgPrice", &vars);
2115
2116 let ws_stream = api.avg_price(params).await.unwrap();
2117
2118 let called = Arc::new(AtomicBool::new(false));
2119 let called_clone = called.clone();
2120 ws_stream.on_message(move |_payload: models::AvgPriceResponse| {
2121 called_clone.store(true, Ordering::SeqCst);
2122 });
2123
2124 assert!(
2125 streams_base.is_subscribed(&stream).await,
2126 "should be subscribed before unsubscribe"
2127 );
2128
2129 ws_stream.unsubscribe().await;
2130
2131 let payload: Value = serde_json::from_str(
2132 r#"{"e":"avgPrice","E":1693907033000,"s":"BTCUSDT","i":"5m","T":1693907032213}"#,
2133 )
2134 .unwrap_or_else(|_| serde_json::json!({}));
2135 let msg = json!({
2136 "stream": stream,
2137 "data": payload,
2138 });
2139
2140 streams_base.on_message(msg.to_string(), conn.clone()).await;
2141
2142 yield_now().await;
2143
2144 assert!(
2145 !called.load(Ordering::SeqCst),
2146 "callback should not be invoked after unsubscribe"
2147 );
2148 });
2149 }
2150
2151 #[test]
2152 fn block_trade_should_execute_successfully() {
2153 TOKIO_SHARED_RT.block_on(async {
2154 let (streams_base, _) = make_streams_base().await;
2155 let api = ApiClient::new(streams_base.clone());
2156
2157 let id = "test-id-123".to_string();
2158
2159 let params = BlockTradeParams::builder("bnbusdt".to_string())
2160 .id(Some(id.clone()))
2161 .build()
2162 .unwrap();
2163
2164 let BlockTradeParams { symbol, id } = params.clone();
2165
2166 let pairs: &[(&str, Option<String>)] =
2167 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2168
2169 let vars: HashMap<_, _> = pairs
2170 .iter()
2171 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2172 .collect();
2173 let stream = replace_websocket_streams_placeholders("/<symbol>@blockTrade", &vars);
2174 let ws_stream = api
2175 .block_trade(params)
2176 .await
2177 .expect("block_trade should return a WebsocketStream");
2178
2179 assert!(
2180 streams_base.is_subscribed(&stream).await,
2181 "expected stream '{stream}' to be subscribed"
2182 );
2183 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2184 });
2185 }
2186
2187 #[test]
2188 fn block_trade_should_handle_incoming_message() {
2189 TOKIO_SHARED_RT.block_on(async {
2190 let (streams_base, conn) = make_streams_base().await;
2191 let api = ApiClient::new(streams_base.clone());
2192
2193 let id = "test-id-123".to_string();
2194
2195 let params = BlockTradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
2196
2197 let BlockTradeParams {
2198 symbol,id,
2199 } = params.clone();
2200
2201 let pairs: &[(&str, Option<String>)] = &[
2202 ("symbol",
2203 Some(symbol.clone())
2204 ),
2205 ("id",
2206 id.clone()
2207 ),
2208 ];
2209
2210 let vars: HashMap<_, _> = pairs
2211 .iter()
2212 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2213 .collect();
2214 let stream = replace_websocket_streams_placeholders("/<symbol>@blockTrade", &vars);
2215
2216 let ws_stream = api.block_trade(params).await.unwrap();
2217
2218 let called = Arc::new(AtomicBool::new(false));
2219 let called_with_message = called.clone();
2220 ws_stream.on_message(move |_payload: models::BlockTradeResponse| {
2221 called_with_message.store(true, Ordering::SeqCst);
2222 });
2223
2224 let payload: Value = serde_json::from_str(r#"{"e":"blockTrade","E":1772506983582,"s":"BNBBTC","t":582,"p":"0.052","q":"5838","T":1772506983321,"m":true}"#).unwrap_or_else(|_| serde_json::json!({}));
2225 let msg = json!({
2226 "stream": stream,
2227 "data": payload,
2228 });
2229
2230 streams_base.on_message(msg.to_string(), conn.clone()).await;
2231 yield_now().await;
2232
2233 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2234 });
2235 }
2236
2237 #[test]
2238 fn block_trade_should_not_fire_after_unsubscribe() {
2239 TOKIO_SHARED_RT.block_on(async {
2240 let (streams_base, conn) = make_streams_base().await;
2241 let api = ApiClient::new(streams_base.clone());
2242
2243 let id = "test-id-123".to_string();
2244
2245 let params = BlockTradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
2246
2247 let BlockTradeParams {
2248 symbol,id,
2249 } = params.clone();
2250
2251 let pairs: &[(&str, Option<String>)] = &[
2252 ("symbol",
2253 Some(symbol.clone())
2254 ),
2255 ("id",
2256 id.clone()
2257 ),
2258 ];
2259
2260 let vars: HashMap<_, _> = pairs
2261 .iter()
2262 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2263 .collect();
2264 let stream = replace_websocket_streams_placeholders("/<symbol>@blockTrade", &vars);
2265
2266 let ws_stream = api.block_trade(params).await.unwrap();
2267
2268 let called = Arc::new(AtomicBool::new(false));
2269 let called_clone = called.clone();
2270 ws_stream.on_message(move |_payload: models::BlockTradeResponse| {
2271 called_clone.store(true, Ordering::SeqCst);
2272 });
2273
2274 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2275
2276 ws_stream.unsubscribe().await;
2277
2278 let payload: Value = serde_json::from_str(r#"{"e":"blockTrade","E":1772506983582,"s":"BNBBTC","t":582,"p":"0.052","q":"5838","T":1772506983321,"m":true}"#).unwrap_or_else(|_| serde_json::json!({}));
2279 let msg = json!({
2280 "stream": stream,
2281 "data": payload,
2282 });
2283
2284 streams_base.on_message(msg.to_string(), conn.clone()).await;
2285
2286 yield_now().await;
2287
2288 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2289 });
2290 }
2291
2292 #[test]
2293 fn book_ticker_should_execute_successfully() {
2294 TOKIO_SHARED_RT.block_on(async {
2295 let (streams_base, _) = make_streams_base().await;
2296 let api = ApiClient::new(streams_base.clone());
2297
2298 let id = "test-id-123".to_string();
2299
2300 let params = BookTickerParams::builder("bnbusdt".to_string())
2301 .id(Some(id.clone()))
2302 .build()
2303 .unwrap();
2304
2305 let BookTickerParams { symbol, id } = params.clone();
2306
2307 let pairs: &[(&str, Option<String>)] =
2308 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2309
2310 let vars: HashMap<_, _> = pairs
2311 .iter()
2312 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2313 .collect();
2314 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2315 let ws_stream = api
2316 .book_ticker(params)
2317 .await
2318 .expect("book_ticker should return a WebsocketStream");
2319
2320 assert!(
2321 streams_base.is_subscribed(&stream).await,
2322 "expected stream '{stream}' to be subscribed"
2323 );
2324 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2325 });
2326 }
2327
2328 #[test]
2329 fn book_ticker_should_handle_incoming_message() {
2330 TOKIO_SHARED_RT.block_on(async {
2331 let (streams_base, conn) = make_streams_base().await;
2332 let api = ApiClient::new(streams_base.clone());
2333
2334 let id = "test-id-123".to_string();
2335
2336 let params = BookTickerParams::builder("bnbusdt".to_string())
2337 .id(Some(id.clone()))
2338 .build()
2339 .unwrap();
2340
2341 let BookTickerParams { symbol, id } = params.clone();
2342
2343 let pairs: &[(&str, Option<String>)] =
2344 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2345
2346 let vars: HashMap<_, _> = pairs
2347 .iter()
2348 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2349 .collect();
2350 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2351
2352 let ws_stream = api.book_ticker(params).await.unwrap();
2353
2354 let called = Arc::new(AtomicBool::new(false));
2355 let called_with_message = called.clone();
2356 ws_stream.on_message(move |_payload: models::BookTickerResponse| {
2357 called_with_message.store(true, Ordering::SeqCst);
2358 });
2359
2360 let payload: Value = serde_json::from_str(r#"{"u":400900217,"s":"BNBUSDT"}"#)
2361 .unwrap_or_else(|_| serde_json::json!({}));
2362 let msg = json!({
2363 "stream": stream,
2364 "data": payload,
2365 });
2366
2367 streams_base.on_message(msg.to_string(), conn.clone()).await;
2368 yield_now().await;
2369
2370 assert!(
2371 called.load(Ordering::SeqCst),
2372 "expected our callback to have been invoked"
2373 );
2374 });
2375 }
2376
2377 #[test]
2378 fn book_ticker_should_not_fire_after_unsubscribe() {
2379 TOKIO_SHARED_RT.block_on(async {
2380 let (streams_base, conn) = make_streams_base().await;
2381 let api = ApiClient::new(streams_base.clone());
2382
2383 let id = "test-id-123".to_string();
2384
2385 let params = BookTickerParams::builder("bnbusdt".to_string())
2386 .id(Some(id.clone()))
2387 .build()
2388 .unwrap();
2389
2390 let BookTickerParams { symbol, id } = params.clone();
2391
2392 let pairs: &[(&str, Option<String>)] =
2393 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2394
2395 let vars: HashMap<_, _> = pairs
2396 .iter()
2397 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2398 .collect();
2399 let stream = replace_websocket_streams_placeholders("/<symbol>@bookTicker", &vars);
2400
2401 let ws_stream = api.book_ticker(params).await.unwrap();
2402
2403 let called = Arc::new(AtomicBool::new(false));
2404 let called_clone = called.clone();
2405 ws_stream.on_message(move |_payload: models::BookTickerResponse| {
2406 called_clone.store(true, Ordering::SeqCst);
2407 });
2408
2409 assert!(
2410 streams_base.is_subscribed(&stream).await,
2411 "should be subscribed before unsubscribe"
2412 );
2413
2414 ws_stream.unsubscribe().await;
2415
2416 let payload: Value = serde_json::from_str(r#"{"u":400900217,"s":"BNBUSDT"}"#)
2417 .unwrap_or_else(|_| serde_json::json!({}));
2418 let msg = json!({
2419 "stream": stream,
2420 "data": payload,
2421 });
2422
2423 streams_base.on_message(msg.to_string(), conn.clone()).await;
2424
2425 yield_now().await;
2426
2427 assert!(
2428 !called.load(Ordering::SeqCst),
2429 "callback should not be invoked after unsubscribe"
2430 );
2431 });
2432 }
2433
2434 #[test]
2435 fn diff_book_depth_should_execute_successfully() {
2436 TOKIO_SHARED_RT.block_on(async {
2437 let (streams_base, _) = make_streams_base().await;
2438 let api = ApiClient::new(streams_base.clone());
2439
2440 let id = "test-id-123".to_string();
2441
2442 let params = DiffBookDepthParams::builder("bnbusdt".to_string())
2443 .id(Some(id.clone()))
2444 .build()
2445 .unwrap();
2446
2447 let DiffBookDepthParams {
2448 symbol,
2449 id,
2450 update_speed,
2451 } = params.clone();
2452
2453 let pairs: &[(&str, Option<String>)] = &[
2454 ("symbol", Some(symbol.clone())),
2455 ("id", id.clone()),
2456 (
2457 "updateSpeed",
2458 update_speed.clone().map(|v| v.as_str().to_string()),
2459 ),
2460 ];
2461
2462 let vars: HashMap<_, _> = pairs
2463 .iter()
2464 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2465 .collect();
2466 let stream =
2467 replace_websocket_streams_placeholders("/<symbol>@depth@<updateSpeed>", &vars);
2468 let ws_stream = api
2469 .diff_book_depth(params)
2470 .await
2471 .expect("diff_book_depth should return a WebsocketStream");
2472
2473 assert!(
2474 streams_base.is_subscribed(&stream).await,
2475 "expected stream '{stream}' to be subscribed"
2476 );
2477 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2478 });
2479 }
2480
2481 #[test]
2482 fn diff_book_depth_should_handle_incoming_message() {
2483 TOKIO_SHARED_RT.block_on(async {
2484 let (streams_base, conn) = make_streams_base().await;
2485 let api = ApiClient::new(streams_base.clone());
2486
2487 let id = "test-id-123".to_string();
2488
2489 let params = DiffBookDepthParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
2490
2491 let DiffBookDepthParams {
2492 symbol,id,update_speed,
2493 } = params.clone();
2494
2495 let pairs: &[(&str, Option<String>)] = &[
2496 ("symbol",
2497 Some(symbol.clone())
2498 ),
2499 ("id",
2500 id.clone()
2501 ),
2502 ("updateSpeed",
2503 update_speed.clone().map(|v| v.as_str().to_string())
2504 ),
2505 ];
2506
2507 let vars: HashMap<_, _> = pairs
2508 .iter()
2509 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2510 .collect();
2511 let stream = replace_websocket_streams_placeholders("/<symbol>@depth@<updateSpeed>", &vars);
2512
2513 let ws_stream = api.diff_book_depth(params).await.unwrap();
2514
2515 let called = Arc::new(AtomicBool::new(false));
2516 let called_with_message = called.clone();
2517 ws_stream.on_message(move |_payload: models::DiffBookDepthResponse| {
2518 called_with_message.store(true, Ordering::SeqCst);
2519 });
2520
2521 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1672515782136,"s":"BNBBTC","U":157,"u":160,"b":[["0.0024","10"]],"a":[["0.0026","100"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2522 let msg = json!({
2523 "stream": stream,
2524 "data": payload,
2525 });
2526
2527 streams_base.on_message(msg.to_string(), conn.clone()).await;
2528 yield_now().await;
2529
2530 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2531 });
2532 }
2533
2534 #[test]
2535 fn diff_book_depth_should_not_fire_after_unsubscribe() {
2536 TOKIO_SHARED_RT.block_on(async {
2537 let (streams_base, conn) = make_streams_base().await;
2538 let api = ApiClient::new(streams_base.clone());
2539
2540 let id = "test-id-123".to_string();
2541
2542 let params = DiffBookDepthParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
2543
2544 let DiffBookDepthParams {
2545 symbol,id,update_speed,
2546 } = params.clone();
2547
2548 let pairs: &[(&str, Option<String>)] = &[
2549 ("symbol",
2550 Some(symbol.clone())
2551 ),
2552 ("id",
2553 id.clone()
2554 ),
2555 ("updateSpeed",
2556 update_speed.clone().map(|v| v.as_str().to_string())
2557 ),
2558 ];
2559
2560 let vars: HashMap<_, _> = pairs
2561 .iter()
2562 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2563 .collect();
2564 let stream = replace_websocket_streams_placeholders("/<symbol>@depth@<updateSpeed>", &vars);
2565
2566 let ws_stream = api.diff_book_depth(params).await.unwrap();
2567
2568 let called = Arc::new(AtomicBool::new(false));
2569 let called_clone = called.clone();
2570 ws_stream.on_message(move |_payload: models::DiffBookDepthResponse| {
2571 called_clone.store(true, Ordering::SeqCst);
2572 });
2573
2574 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2575
2576 ws_stream.unsubscribe().await;
2577
2578 let payload: Value = serde_json::from_str(r#"{"e":"depthUpdate","E":1672515782136,"s":"BNBBTC","U":157,"u":160,"b":[["0.0024","10"]],"a":[["0.0026","100"]]}"#).unwrap_or_else(|_| serde_json::json!({}));
2579 let msg = json!({
2580 "stream": stream,
2581 "data": payload,
2582 });
2583
2584 streams_base.on_message(msg.to_string(), conn.clone()).await;
2585
2586 yield_now().await;
2587
2588 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2589 });
2590 }
2591
2592 #[test]
2593 fn kline_should_execute_successfully() {
2594 TOKIO_SHARED_RT.block_on(async {
2595 let (streams_base, _) = make_streams_base().await;
2596 let api = ApiClient::new(streams_base.clone());
2597
2598 let id = "test-id-123".to_string();
2599
2600 let params = KlineParams::builder("bnbusdt".to_string(), KlineIntervalEnum::Interval1s)
2601 .id(Some(id.clone()))
2602 .build()
2603 .unwrap();
2604
2605 let KlineParams {
2606 symbol,
2607 interval,
2608 id,
2609 } = params.clone();
2610
2611 let pairs: &[(&str, Option<String>)] = &[
2612 ("symbol", Some(symbol.clone())),
2613 ("interval", Some(interval.as_str().to_string())),
2614 ("id", id.clone()),
2615 ];
2616
2617 let vars: HashMap<_, _> = pairs
2618 .iter()
2619 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2620 .collect();
2621 let stream =
2622 replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2623 let ws_stream = api
2624 .kline(params)
2625 .await
2626 .expect("kline should return a WebsocketStream");
2627
2628 assert!(
2629 streams_base.is_subscribed(&stream).await,
2630 "expected stream '{stream}' to be subscribed"
2631 );
2632 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2633 });
2634 }
2635
2636 #[test]
2637 fn kline_should_handle_incoming_message() {
2638 TOKIO_SHARED_RT.block_on(async {
2639 let (streams_base, conn) = make_streams_base().await;
2640 let api = ApiClient::new(streams_base.clone());
2641
2642 let id = "test-id-123".to_string();
2643
2644 let params = KlineParams::builder("bnbusdt".to_string(),KlineIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2645
2646 let KlineParams {
2647 symbol,interval,id,
2648 } = params.clone();
2649
2650 let pairs: &[(&str, Option<String>)] = &[
2651 ("symbol",
2652 Some(symbol.clone())
2653 ),
2654 ("interval",
2655 Some(interval.as_str().to_string())
2656 ),
2657 ("id",
2658 id.clone()
2659 ),
2660 ];
2661
2662 let vars: HashMap<_, _> = pairs
2663 .iter()
2664 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2665 .collect();
2666 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2667
2668 let ws_stream = api.kline(params).await.unwrap();
2669
2670 let called = Arc::new(AtomicBool::new(false));
2671 let called_with_message = called.clone();
2672 ws_stream.on_message(move |_payload: models::KlineResponse| {
2673 called_with_message.store(true, Ordering::SeqCst);
2674 });
2675
2676 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1672515782136,"s":"BNBBTC","k":{"t":1672515780000,"T":1672515839999,"s":"BNBBTC","i":"1m","f":100,"L":200,"n":100,"x":false}}"#).unwrap_or_else(|_| serde_json::json!({}));
2677 let msg = json!({
2678 "stream": stream,
2679 "data": payload,
2680 });
2681
2682 streams_base.on_message(msg.to_string(), conn.clone()).await;
2683 yield_now().await;
2684
2685 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2686 });
2687 }
2688
2689 #[test]
2690 fn kline_should_not_fire_after_unsubscribe() {
2691 TOKIO_SHARED_RT.block_on(async {
2692 let (streams_base, conn) = make_streams_base().await;
2693 let api = ApiClient::new(streams_base.clone());
2694
2695 let id = "test-id-123".to_string();
2696
2697 let params = KlineParams::builder("bnbusdt".to_string(),KlineIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2698
2699 let KlineParams {
2700 symbol,interval,id,
2701 } = params.clone();
2702
2703 let pairs: &[(&str, Option<String>)] = &[
2704 ("symbol",
2705 Some(symbol.clone())
2706 ),
2707 ("interval",
2708 Some(interval.as_str().to_string())
2709 ),
2710 ("id",
2711 id.clone()
2712 ),
2713 ];
2714
2715 let vars: HashMap<_, _> = pairs
2716 .iter()
2717 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2718 .collect();
2719 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>", &vars);
2720
2721 let ws_stream = api.kline(params).await.unwrap();
2722
2723 let called = Arc::new(AtomicBool::new(false));
2724 let called_clone = called.clone();
2725 ws_stream.on_message(move |_payload: models::KlineResponse| {
2726 called_clone.store(true, Ordering::SeqCst);
2727 });
2728
2729 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2730
2731 ws_stream.unsubscribe().await;
2732
2733 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1672515782136,"s":"BNBBTC","k":{"t":1672515780000,"T":1672515839999,"s":"BNBBTC","i":"1m","f":100,"L":200,"n":100,"x":false}}"#).unwrap_or_else(|_| serde_json::json!({}));
2734 let msg = json!({
2735 "stream": stream,
2736 "data": payload,
2737 });
2738
2739 streams_base.on_message(msg.to_string(), conn.clone()).await;
2740
2741 yield_now().await;
2742
2743 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2744 });
2745 }
2746
2747 #[test]
2748 fn kline_offset_should_execute_successfully() {
2749 TOKIO_SHARED_RT.block_on(async {
2750 let (streams_base, _) = make_streams_base().await;
2751 let api = ApiClient::new(streams_base.clone());
2752
2753 let id = "test-id-123".to_string();
2754
2755 let params = KlineOffsetParams::builder(
2756 "bnbusdt".to_string(),
2757 KlineOffsetIntervalEnum::Interval1s,
2758 )
2759 .id(Some(id.clone()))
2760 .build()
2761 .unwrap();
2762
2763 let KlineOffsetParams {
2764 symbol,
2765 interval,
2766 id,
2767 } = params.clone();
2768
2769 let pairs: &[(&str, Option<String>)] = &[
2770 ("symbol", Some(symbol.clone())),
2771 ("interval", Some(interval.as_str().to_string())),
2772 ("id", id.clone()),
2773 ];
2774
2775 let vars: HashMap<_, _> = pairs
2776 .iter()
2777 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2778 .collect();
2779 let stream =
2780 replace_websocket_streams_placeholders("/<symbol>@kline_<interval>@+08:00", &vars);
2781 let ws_stream = api
2782 .kline_offset(params)
2783 .await
2784 .expect("kline_offset should return a WebsocketStream");
2785
2786 assert!(
2787 streams_base.is_subscribed(&stream).await,
2788 "expected stream '{stream}' to be subscribed"
2789 );
2790 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2791 });
2792 }
2793
2794 #[test]
2795 fn kline_offset_should_handle_incoming_message() {
2796 TOKIO_SHARED_RT.block_on(async {
2797 let (streams_base, conn) = make_streams_base().await;
2798 let api = ApiClient::new(streams_base.clone());
2799
2800 let id = "test-id-123".to_string();
2801
2802 let params = KlineOffsetParams::builder("bnbusdt".to_string(),KlineOffsetIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2803
2804 let KlineOffsetParams {
2805 symbol,interval,id,
2806 } = params.clone();
2807
2808 let pairs: &[(&str, Option<String>)] = &[
2809 ("symbol",
2810 Some(symbol.clone())
2811 ),
2812 ("interval",
2813 Some(interval.as_str().to_string())
2814 ),
2815 ("id",
2816 id.clone()
2817 ),
2818 ];
2819
2820 let vars: HashMap<_, _> = pairs
2821 .iter()
2822 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2823 .collect();
2824 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>@+08:00", &vars);
2825
2826 let ws_stream = api.kline_offset(params).await.unwrap();
2827
2828 let called = Arc::new(AtomicBool::new(false));
2829 let called_with_message = called.clone();
2830 ws_stream.on_message(move |_payload: models::KlineOffsetResponse| {
2831 called_with_message.store(true, Ordering::SeqCst);
2832 });
2833
2834 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1672515782136,"s":"BNBBTC","k":{"t":1672515780000,"T":1672515839999,"s":"BNBBTC","i":"1m","f":100,"L":200,"n":100,"x":false}}"#).unwrap_or_else(|_| serde_json::json!({}));
2835 let msg = json!({
2836 "stream": stream,
2837 "data": payload,
2838 });
2839
2840 streams_base.on_message(msg.to_string(), conn.clone()).await;
2841 yield_now().await;
2842
2843 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
2844 });
2845 }
2846
2847 #[test]
2848 fn kline_offset_should_not_fire_after_unsubscribe() {
2849 TOKIO_SHARED_RT.block_on(async {
2850 let (streams_base, conn) = make_streams_base().await;
2851 let api = ApiClient::new(streams_base.clone());
2852
2853 let id = "test-id-123".to_string();
2854
2855 let params = KlineOffsetParams::builder("bnbusdt".to_string(),KlineOffsetIntervalEnum::Interval1s,).id(Some(id.clone())).build().unwrap();
2856
2857 let KlineOffsetParams {
2858 symbol,interval,id,
2859 } = params.clone();
2860
2861 let pairs: &[(&str, Option<String>)] = &[
2862 ("symbol",
2863 Some(symbol.clone())
2864 ),
2865 ("interval",
2866 Some(interval.as_str().to_string())
2867 ),
2868 ("id",
2869 id.clone()
2870 ),
2871 ];
2872
2873 let vars: HashMap<_, _> = pairs
2874 .iter()
2875 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2876 .collect();
2877 let stream = replace_websocket_streams_placeholders("/<symbol>@kline_<interval>@+08:00", &vars);
2878
2879 let ws_stream = api.kline_offset(params).await.unwrap();
2880
2881 let called = Arc::new(AtomicBool::new(false));
2882 let called_clone = called.clone();
2883 ws_stream.on_message(move |_payload: models::KlineOffsetResponse| {
2884 called_clone.store(true, Ordering::SeqCst);
2885 });
2886
2887 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
2888
2889 ws_stream.unsubscribe().await;
2890
2891 let payload: Value = serde_json::from_str(r#"{"e":"kline","E":1672515782136,"s":"BNBBTC","k":{"t":1672515780000,"T":1672515839999,"s":"BNBBTC","i":"1m","f":100,"L":200,"n":100,"x":false}}"#).unwrap_or_else(|_| serde_json::json!({}));
2892 let msg = json!({
2893 "stream": stream,
2894 "data": payload,
2895 });
2896
2897 streams_base.on_message(msg.to_string(), conn.clone()).await;
2898
2899 yield_now().await;
2900
2901 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
2902 });
2903 }
2904
2905 #[test]
2906 fn mini_ticker_should_execute_successfully() {
2907 TOKIO_SHARED_RT.block_on(async {
2908 let (streams_base, _) = make_streams_base().await;
2909 let api = ApiClient::new(streams_base.clone());
2910
2911 let id = "test-id-123".to_string();
2912
2913 let params = MiniTickerParams::builder("bnbusdt".to_string())
2914 .id(Some(id.clone()))
2915 .build()
2916 .unwrap();
2917
2918 let MiniTickerParams { symbol, id } = params.clone();
2919
2920 let pairs: &[(&str, Option<String>)] =
2921 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2922
2923 let vars: HashMap<_, _> = pairs
2924 .iter()
2925 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2926 .collect();
2927 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2928 let ws_stream = api
2929 .mini_ticker(params)
2930 .await
2931 .expect("mini_ticker should return a WebsocketStream");
2932
2933 assert!(
2934 streams_base.is_subscribed(&stream).await,
2935 "expected stream '{stream}' to be subscribed"
2936 );
2937 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
2938 });
2939 }
2940
2941 #[test]
2942 fn mini_ticker_should_handle_incoming_message() {
2943 TOKIO_SHARED_RT.block_on(async {
2944 let (streams_base, conn) = make_streams_base().await;
2945 let api = ApiClient::new(streams_base.clone());
2946
2947 let id = "test-id-123".to_string();
2948
2949 let params = MiniTickerParams::builder("bnbusdt".to_string())
2950 .id(Some(id.clone()))
2951 .build()
2952 .unwrap();
2953
2954 let MiniTickerParams { symbol, id } = params.clone();
2955
2956 let pairs: &[(&str, Option<String>)] =
2957 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
2958
2959 let vars: HashMap<_, _> = pairs
2960 .iter()
2961 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
2962 .collect();
2963 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
2964
2965 let ws_stream = api.mini_ticker(params).await.unwrap();
2966
2967 let called = Arc::new(AtomicBool::new(false));
2968 let called_with_message = called.clone();
2969 ws_stream.on_message(move |_payload: models::MiniTickerResponse| {
2970 called_with_message.store(true, Ordering::SeqCst);
2971 });
2972
2973 let payload: Value =
2974 serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1672515782136,"s":"BNBBTC"}"#)
2975 .unwrap_or_else(|_| serde_json::json!({}));
2976 let msg = json!({
2977 "stream": stream,
2978 "data": payload,
2979 });
2980
2981 streams_base.on_message(msg.to_string(), conn.clone()).await;
2982 yield_now().await;
2983
2984 assert!(
2985 called.load(Ordering::SeqCst),
2986 "expected our callback to have been invoked"
2987 );
2988 });
2989 }
2990
2991 #[test]
2992 fn mini_ticker_should_not_fire_after_unsubscribe() {
2993 TOKIO_SHARED_RT.block_on(async {
2994 let (streams_base, conn) = make_streams_base().await;
2995 let api = ApiClient::new(streams_base.clone());
2996
2997 let id = "test-id-123".to_string();
2998
2999 let params = MiniTickerParams::builder("bnbusdt".to_string())
3000 .id(Some(id.clone()))
3001 .build()
3002 .unwrap();
3003
3004 let MiniTickerParams { symbol, id } = params.clone();
3005
3006 let pairs: &[(&str, Option<String>)] =
3007 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3008
3009 let vars: HashMap<_, _> = pairs
3010 .iter()
3011 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3012 .collect();
3013 let stream = replace_websocket_streams_placeholders("/<symbol>@miniTicker", &vars);
3014
3015 let ws_stream = api.mini_ticker(params).await.unwrap();
3016
3017 let called = Arc::new(AtomicBool::new(false));
3018 let called_clone = called.clone();
3019 ws_stream.on_message(move |_payload: models::MiniTickerResponse| {
3020 called_clone.store(true, Ordering::SeqCst);
3021 });
3022
3023 assert!(
3024 streams_base.is_subscribed(&stream).await,
3025 "should be subscribed before unsubscribe"
3026 );
3027
3028 ws_stream.unsubscribe().await;
3029
3030 let payload: Value =
3031 serde_json::from_str(r#"{"e":"24hrMiniTicker","E":1672515782136,"s":"BNBBTC"}"#)
3032 .unwrap_or_else(|_| serde_json::json!({}));
3033 let msg = json!({
3034 "stream": stream,
3035 "data": payload,
3036 });
3037
3038 streams_base.on_message(msg.to_string(), conn.clone()).await;
3039
3040 yield_now().await;
3041
3042 assert!(
3043 !called.load(Ordering::SeqCst),
3044 "callback should not be invoked after unsubscribe"
3045 );
3046 });
3047 }
3048
3049 #[test]
3050 fn partial_book_depth_should_execute_successfully() {
3051 TOKIO_SHARED_RT.block_on(async {
3052 let (streams_base, _) = make_streams_base().await;
3053 let api = ApiClient::new(streams_base.clone());
3054
3055 let id = "test-id-123".to_string();
3056
3057 let params = PartialBookDepthParams::builder(
3058 "bnbusdt".to_string(),
3059 PartialBookDepthLevelsEnum::Levels5,
3060 )
3061 .id(Some(id.clone()))
3062 .build()
3063 .unwrap();
3064
3065 let PartialBookDepthParams {
3066 symbol,
3067 levels,
3068 id,
3069 update_speed,
3070 } = params.clone();
3071
3072 let pairs: &[(&str, Option<String>)] = &[
3073 ("symbol", Some(symbol.clone())),
3074 ("levels", Some(levels.as_str().to_string())),
3075 ("id", id.clone()),
3076 (
3077 "updateSpeed",
3078 update_speed.clone().map(|v| v.as_str().to_string()),
3079 ),
3080 ];
3081
3082 let vars: HashMap<_, _> = pairs
3083 .iter()
3084 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3085 .collect();
3086 let stream = replace_websocket_streams_placeholders(
3087 "/<symbol>@depth<levels>@<updateSpeed>",
3088 &vars,
3089 );
3090 let ws_stream = api
3091 .partial_book_depth(params)
3092 .await
3093 .expect("partial_book_depth should return a WebsocketStream");
3094
3095 assert!(
3096 streams_base.is_subscribed(&stream).await,
3097 "expected stream '{stream}' to be subscribed"
3098 );
3099 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3100 });
3101 }
3102
3103 #[test]
3104 fn partial_book_depth_should_handle_incoming_message() {
3105 TOKIO_SHARED_RT.block_on(async {
3106 let (streams_base, conn) = make_streams_base().await;
3107 let api = ApiClient::new(streams_base.clone());
3108
3109 let id = "test-id-123".to_string();
3110
3111 let params = PartialBookDepthParams::builder(
3112 "bnbusdt".to_string(),
3113 PartialBookDepthLevelsEnum::Levels5,
3114 )
3115 .id(Some(id.clone()))
3116 .build()
3117 .unwrap();
3118
3119 let PartialBookDepthParams {
3120 symbol,
3121 levels,
3122 id,
3123 update_speed,
3124 } = params.clone();
3125
3126 let pairs: &[(&str, Option<String>)] = &[
3127 ("symbol", Some(symbol.clone())),
3128 ("levels", Some(levels.as_str().to_string())),
3129 ("id", id.clone()),
3130 (
3131 "updateSpeed",
3132 update_speed.clone().map(|v| v.as_str().to_string()),
3133 ),
3134 ];
3135
3136 let vars: HashMap<_, _> = pairs
3137 .iter()
3138 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3139 .collect();
3140 let stream = replace_websocket_streams_placeholders(
3141 "/<symbol>@depth<levels>@<updateSpeed>",
3142 &vars,
3143 );
3144
3145 let ws_stream = api.partial_book_depth(params).await.unwrap();
3146
3147 let called = Arc::new(AtomicBool::new(false));
3148 let called_with_message = called.clone();
3149 ws_stream.on_message(move |_payload: models::PartialBookDepthResponse| {
3150 called_with_message.store(true, Ordering::SeqCst);
3151 });
3152
3153 let payload: Value = serde_json::from_str(
3154 r#"{"lastUpdateId":160,"bids":[["0.0024","10"]],"asks":[["0.0026","100"]]}"#,
3155 )
3156 .unwrap_or_else(|_| serde_json::json!({}));
3157 let msg = json!({
3158 "stream": stream,
3159 "data": payload,
3160 });
3161
3162 streams_base.on_message(msg.to_string(), conn.clone()).await;
3163 yield_now().await;
3164
3165 assert!(
3166 called.load(Ordering::SeqCst),
3167 "expected our callback to have been invoked"
3168 );
3169 });
3170 }
3171
3172 #[test]
3173 fn partial_book_depth_should_not_fire_after_unsubscribe() {
3174 TOKIO_SHARED_RT.block_on(async {
3175 let (streams_base, conn) = make_streams_base().await;
3176 let api = ApiClient::new(streams_base.clone());
3177
3178 let id = "test-id-123".to_string();
3179
3180 let params = PartialBookDepthParams::builder(
3181 "bnbusdt".to_string(),
3182 PartialBookDepthLevelsEnum::Levels5,
3183 )
3184 .id(Some(id.clone()))
3185 .build()
3186 .unwrap();
3187
3188 let PartialBookDepthParams {
3189 symbol,
3190 levels,
3191 id,
3192 update_speed,
3193 } = params.clone();
3194
3195 let pairs: &[(&str, Option<String>)] = &[
3196 ("symbol", Some(symbol.clone())),
3197 ("levels", Some(levels.as_str().to_string())),
3198 ("id", id.clone()),
3199 (
3200 "updateSpeed",
3201 update_speed.clone().map(|v| v.as_str().to_string()),
3202 ),
3203 ];
3204
3205 let vars: HashMap<_, _> = pairs
3206 .iter()
3207 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3208 .collect();
3209 let stream = replace_websocket_streams_placeholders(
3210 "/<symbol>@depth<levels>@<updateSpeed>",
3211 &vars,
3212 );
3213
3214 let ws_stream = api.partial_book_depth(params).await.unwrap();
3215
3216 let called = Arc::new(AtomicBool::new(false));
3217 let called_clone = called.clone();
3218 ws_stream.on_message(move |_payload: models::PartialBookDepthResponse| {
3219 called_clone.store(true, Ordering::SeqCst);
3220 });
3221
3222 assert!(
3223 streams_base.is_subscribed(&stream).await,
3224 "should be subscribed before unsubscribe"
3225 );
3226
3227 ws_stream.unsubscribe().await;
3228
3229 let payload: Value = serde_json::from_str(
3230 r#"{"lastUpdateId":160,"bids":[["0.0024","10"]],"asks":[["0.0026","100"]]}"#,
3231 )
3232 .unwrap_or_else(|_| serde_json::json!({}));
3233 let msg = json!({
3234 "stream": stream,
3235 "data": payload,
3236 });
3237
3238 streams_base.on_message(msg.to_string(), conn.clone()).await;
3239
3240 yield_now().await;
3241
3242 assert!(
3243 !called.load(Ordering::SeqCst),
3244 "callback should not be invoked after unsubscribe"
3245 );
3246 });
3247 }
3248
3249 #[test]
3250 fn reference_price_should_execute_successfully() {
3251 TOKIO_SHARED_RT.block_on(async {
3252 let (streams_base, _) = make_streams_base().await;
3253 let api = ApiClient::new(streams_base.clone());
3254
3255 let id = "test-id-123".to_string();
3256
3257 let params = ReferencePriceParams::builder("bazusd".to_string())
3258 .id(Some(id.clone()))
3259 .build()
3260 .unwrap();
3261
3262 let ReferencePriceParams { symbol, id } = params.clone();
3263
3264 let pairs: &[(&str, Option<String>)] =
3265 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3266
3267 let vars: HashMap<_, _> = pairs
3268 .iter()
3269 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3270 .collect();
3271 let stream = replace_websocket_streams_placeholders("/<symbol>@referencePrice", &vars);
3272 let ws_stream = api
3273 .reference_price(params)
3274 .await
3275 .expect("reference_price should return a WebsocketStream");
3276
3277 assert!(
3278 streams_base.is_subscribed(&stream).await,
3279 "expected stream '{stream}' to be subscribed"
3280 );
3281 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3282 });
3283 }
3284
3285 #[test]
3286 fn reference_price_should_handle_incoming_message() {
3287 TOKIO_SHARED_RT.block_on(async {
3288 let (streams_base, conn) = make_streams_base().await;
3289 let api = ApiClient::new(streams_base.clone());
3290
3291 let id = "test-id-123".to_string();
3292
3293 let params = ReferencePriceParams::builder("bazusd".to_string())
3294 .id(Some(id.clone()))
3295 .build()
3296 .unwrap();
3297
3298 let ReferencePriceParams { symbol, id } = params.clone();
3299
3300 let pairs: &[(&str, Option<String>)] =
3301 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3302
3303 let vars: HashMap<_, _> = pairs
3304 .iter()
3305 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3306 .collect();
3307 let stream = replace_websocket_streams_placeholders("/<symbol>@referencePrice", &vars);
3308
3309 let ws_stream = api.reference_price(params).await.unwrap();
3310
3311 let called = Arc::new(AtomicBool::new(false));
3312 let called_with_message = called.clone();
3313 ws_stream.on_message(move |_payload: models::ReferencePriceResponse| {
3314 called_with_message.store(true, Ordering::SeqCst);
3315 });
3316
3317 let payload: Value = serde_json::from_str(
3318 r#"{"e":"referencePrice","s":"BAZUSD","r":"1.00","t":1770313263917}"#,
3319 )
3320 .unwrap_or_else(|_| serde_json::json!({}));
3321 let msg = json!({
3322 "stream": stream,
3323 "data": payload,
3324 });
3325
3326 streams_base.on_message(msg.to_string(), conn.clone()).await;
3327 yield_now().await;
3328
3329 assert!(
3330 called.load(Ordering::SeqCst),
3331 "expected our callback to have been invoked"
3332 );
3333 });
3334 }
3335
3336 #[test]
3337 fn reference_price_should_not_fire_after_unsubscribe() {
3338 TOKIO_SHARED_RT.block_on(async {
3339 let (streams_base, conn) = make_streams_base().await;
3340 let api = ApiClient::new(streams_base.clone());
3341
3342 let id = "test-id-123".to_string();
3343
3344 let params = ReferencePriceParams::builder("bazusd".to_string())
3345 .id(Some(id.clone()))
3346 .build()
3347 .unwrap();
3348
3349 let ReferencePriceParams { symbol, id } = params.clone();
3350
3351 let pairs: &[(&str, Option<String>)] =
3352 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3353
3354 let vars: HashMap<_, _> = pairs
3355 .iter()
3356 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3357 .collect();
3358 let stream = replace_websocket_streams_placeholders("/<symbol>@referencePrice", &vars);
3359
3360 let ws_stream = api.reference_price(params).await.unwrap();
3361
3362 let called = Arc::new(AtomicBool::new(false));
3363 let called_clone = called.clone();
3364 ws_stream.on_message(move |_payload: models::ReferencePriceResponse| {
3365 called_clone.store(true, Ordering::SeqCst);
3366 });
3367
3368 assert!(
3369 streams_base.is_subscribed(&stream).await,
3370 "should be subscribed before unsubscribe"
3371 );
3372
3373 ws_stream.unsubscribe().await;
3374
3375 let payload: Value = serde_json::from_str(
3376 r#"{"e":"referencePrice","s":"BAZUSD","r":"1.00","t":1770313263917}"#,
3377 )
3378 .unwrap_or_else(|_| serde_json::json!({}));
3379 let msg = json!({
3380 "stream": stream,
3381 "data": payload,
3382 });
3383
3384 streams_base.on_message(msg.to_string(), conn.clone()).await;
3385
3386 yield_now().await;
3387
3388 assert!(
3389 !called.load(Ordering::SeqCst),
3390 "callback should not be invoked after unsubscribe"
3391 );
3392 });
3393 }
3394
3395 #[test]
3396 fn rolling_window_ticker_should_execute_successfully() {
3397 TOKIO_SHARED_RT.block_on(async {
3398 let (streams_base, _) = make_streams_base().await;
3399 let api = ApiClient::new(streams_base.clone());
3400
3401 let id = "test-id-123".to_string();
3402
3403 let params = RollingWindowTickerParams::builder(
3404 "bnbusdt".to_string(),
3405 RollingWindowTickerWindowSizeEnum::WindowSize1h,
3406 )
3407 .id(Some(id.clone()))
3408 .build()
3409 .unwrap();
3410
3411 let RollingWindowTickerParams {
3412 symbol,
3413 window_size,
3414 id,
3415 } = params.clone();
3416
3417 let pairs: &[(&str, Option<String>)] = &[
3418 ("symbol", Some(symbol.clone())),
3419 ("windowSize", Some(window_size.as_str().to_string())),
3420 ("id", id.clone()),
3421 ];
3422
3423 let vars: HashMap<_, _> = pairs
3424 .iter()
3425 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3426 .collect();
3427 let stream =
3428 replace_websocket_streams_placeholders("/<symbol>@ticker_<windowSize>", &vars);
3429 let ws_stream = api
3430 .rolling_window_ticker(params)
3431 .await
3432 .expect("rolling_window_ticker should return a WebsocketStream");
3433
3434 assert!(
3435 streams_base.is_subscribed(&stream).await,
3436 "expected stream '{stream}' to be subscribed"
3437 );
3438 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3439 });
3440 }
3441
3442 #[test]
3443 fn rolling_window_ticker_should_handle_incoming_message() {
3444 TOKIO_SHARED_RT.block_on(async {
3445 let (streams_base, conn) = make_streams_base().await;
3446 let api = ApiClient::new(streams_base.clone());
3447
3448 let id = "test-id-123".to_string();
3449
3450 let params = RollingWindowTickerParams::builder("bnbusdt".to_string(),RollingWindowTickerWindowSizeEnum::WindowSize1h,).id(Some(id.clone())).build().unwrap();
3451
3452 let RollingWindowTickerParams {
3453 symbol,window_size,id,
3454 } = params.clone();
3455
3456 let pairs: &[(&str, Option<String>)] = &[
3457 ("symbol",
3458 Some(symbol.clone())
3459 ),
3460 ("windowSize",
3461 Some(window_size.as_str().to_string())
3462 ),
3463 ("id",
3464 id.clone()
3465 ),
3466 ];
3467
3468 let vars: HashMap<_, _> = pairs
3469 .iter()
3470 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3471 .collect();
3472 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker_<windowSize>", &vars);
3473
3474 let ws_stream = api.rolling_window_ticker(params).await.unwrap();
3475
3476 let called = Arc::new(AtomicBool::new(false));
3477 let called_with_message = called.clone();
3478 ws_stream.on_message(move |_payload: models::RollingWindowTickerResponse| {
3479 called_with_message.store(true, Ordering::SeqCst);
3480 });
3481
3482 let payload: Value = serde_json::from_str(r#"{"e":"1hTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}"#).unwrap_or_else(|_| serde_json::json!({}));
3483 let msg = json!({
3484 "stream": stream,
3485 "data": payload,
3486 });
3487
3488 streams_base.on_message(msg.to_string(), conn.clone()).await;
3489 yield_now().await;
3490
3491 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3492 });
3493 }
3494
3495 #[test]
3496 fn rolling_window_ticker_should_not_fire_after_unsubscribe() {
3497 TOKIO_SHARED_RT.block_on(async {
3498 let (streams_base, conn) = make_streams_base().await;
3499 let api = ApiClient::new(streams_base.clone());
3500
3501 let id = "test-id-123".to_string();
3502
3503 let params = RollingWindowTickerParams::builder("bnbusdt".to_string(),RollingWindowTickerWindowSizeEnum::WindowSize1h,).id(Some(id.clone())).build().unwrap();
3504
3505 let RollingWindowTickerParams {
3506 symbol,window_size,id,
3507 } = params.clone();
3508
3509 let pairs: &[(&str, Option<String>)] = &[
3510 ("symbol",
3511 Some(symbol.clone())
3512 ),
3513 ("windowSize",
3514 Some(window_size.as_str().to_string())
3515 ),
3516 ("id",
3517 id.clone()
3518 ),
3519 ];
3520
3521 let vars: HashMap<_, _> = pairs
3522 .iter()
3523 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3524 .collect();
3525 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker_<windowSize>", &vars);
3526
3527 let ws_stream = api.rolling_window_ticker(params).await.unwrap();
3528
3529 let called = Arc::new(AtomicBool::new(false));
3530 let called_clone = called.clone();
3531 ws_stream.on_message(move |_payload: models::RollingWindowTickerResponse| {
3532 called_clone.store(true, Ordering::SeqCst);
3533 });
3534
3535 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3536
3537 ws_stream.unsubscribe().await;
3538
3539 let payload: Value = serde_json::from_str(r#"{"e":"1hTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}"#).unwrap_or_else(|_| serde_json::json!({}));
3540 let msg = json!({
3541 "stream": stream,
3542 "data": payload,
3543 });
3544
3545 streams_base.on_message(msg.to_string(), conn.clone()).await;
3546
3547 yield_now().await;
3548
3549 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3550 });
3551 }
3552
3553 #[test]
3554 fn ticker_should_execute_successfully() {
3555 TOKIO_SHARED_RT.block_on(async {
3556 let (streams_base, _) = make_streams_base().await;
3557 let api = ApiClient::new(streams_base.clone());
3558
3559 let id = "test-id-123".to_string();
3560
3561 let params = TickerParams::builder("bnbusdt".to_string())
3562 .id(Some(id.clone()))
3563 .build()
3564 .unwrap();
3565
3566 let TickerParams { symbol, id } = params.clone();
3567
3568 let pairs: &[(&str, Option<String>)] =
3569 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3570
3571 let vars: HashMap<_, _> = pairs
3572 .iter()
3573 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3574 .collect();
3575 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
3576 let ws_stream = api
3577 .ticker(params)
3578 .await
3579 .expect("ticker should return a WebsocketStream");
3580
3581 assert!(
3582 streams_base.is_subscribed(&stream).await,
3583 "expected stream '{stream}' to be subscribed"
3584 );
3585 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3586 });
3587 }
3588
3589 #[test]
3590 fn ticker_should_handle_incoming_message() {
3591 TOKIO_SHARED_RT.block_on(async {
3592 let (streams_base, conn) = make_streams_base().await;
3593 let api = ApiClient::new(streams_base.clone());
3594
3595 let id = "test-id-123".to_string();
3596
3597 let params = TickerParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
3598
3599 let TickerParams {
3600 symbol,id,
3601 } = params.clone();
3602
3603 let pairs: &[(&str, Option<String>)] = &[
3604 ("symbol",
3605 Some(symbol.clone())
3606 ),
3607 ("id",
3608 id.clone()
3609 ),
3610 ];
3611
3612 let vars: HashMap<_, _> = pairs
3613 .iter()
3614 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3615 .collect();
3616 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
3617
3618 let ws_stream = api.ticker(params).await.unwrap();
3619
3620 let called = Arc::new(AtomicBool::new(false));
3621 let called_with_message = called.clone();
3622 ws_stream.on_message(move |_payload: models::TickerResponse| {
3623 called_with_message.store(true, Ordering::SeqCst);
3624 });
3625
3626 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}"#).unwrap_or_else(|_| serde_json::json!({}));
3627 let msg = json!({
3628 "stream": stream,
3629 "data": payload,
3630 });
3631
3632 streams_base.on_message(msg.to_string(), conn.clone()).await;
3633 yield_now().await;
3634
3635 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3636 });
3637 }
3638
3639 #[test]
3640 fn ticker_should_not_fire_after_unsubscribe() {
3641 TOKIO_SHARED_RT.block_on(async {
3642 let (streams_base, conn) = make_streams_base().await;
3643 let api = ApiClient::new(streams_base.clone());
3644
3645 let id = "test-id-123".to_string();
3646
3647 let params = TickerParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
3648
3649 let TickerParams {
3650 symbol,id,
3651 } = params.clone();
3652
3653 let pairs: &[(&str, Option<String>)] = &[
3654 ("symbol",
3655 Some(symbol.clone())
3656 ),
3657 ("id",
3658 id.clone()
3659 ),
3660 ];
3661
3662 let vars: HashMap<_, _> = pairs
3663 .iter()
3664 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3665 .collect();
3666 let stream = replace_websocket_streams_placeholders("/<symbol>@ticker", &vars);
3667
3668 let ws_stream = api.ticker(params).await.unwrap();
3669
3670 let called = Arc::new(AtomicBool::new(false));
3671 let called_clone = called.clone();
3672 ws_stream.on_message(move |_payload: models::TickerResponse| {
3673 called_clone.store(true, Ordering::SeqCst);
3674 });
3675
3676 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3677
3678 ws_stream.unsubscribe().await;
3679
3680 let payload: Value = serde_json::from_str(r#"{"e":"24hrTicker","E":1672515782136,"s":"BNBBTC","O":0,"C":1675216573749,"F":0,"L":18150,"n":18151}"#).unwrap_or_else(|_| serde_json::json!({}));
3681 let msg = json!({
3682 "stream": stream,
3683 "data": payload,
3684 });
3685
3686 streams_base.on_message(msg.to_string(), conn.clone()).await;
3687
3688 yield_now().await;
3689
3690 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3691 });
3692 }
3693
3694 #[test]
3695 fn trade_should_execute_successfully() {
3696 TOKIO_SHARED_RT.block_on(async {
3697 let (streams_base, _) = make_streams_base().await;
3698 let api = ApiClient::new(streams_base.clone());
3699
3700 let id = "test-id-123".to_string();
3701
3702 let params = TradeParams::builder("bnbusdt".to_string())
3703 .id(Some(id.clone()))
3704 .build()
3705 .unwrap();
3706
3707 let TradeParams { symbol, id } = params.clone();
3708
3709 let pairs: &[(&str, Option<String>)] =
3710 &[("symbol", Some(symbol.clone())), ("id", id.clone())];
3711
3712 let vars: HashMap<_, _> = pairs
3713 .iter()
3714 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3715 .collect();
3716 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3717 let ws_stream = api
3718 .trade(params)
3719 .await
3720 .expect("trade should return a WebsocketStream");
3721
3722 assert!(
3723 streams_base.is_subscribed(&stream).await,
3724 "expected stream '{stream}' to be subscribed"
3725 );
3726 assert_eq!(ws_stream.id, Some(StreamId::Str("test-id-123".to_string())));
3727 });
3728 }
3729
3730 #[test]
3731 fn trade_should_handle_incoming_message() {
3732 TOKIO_SHARED_RT.block_on(async {
3733 let (streams_base, conn) = make_streams_base().await;
3734 let api = ApiClient::new(streams_base.clone());
3735
3736 let id = "test-id-123".to_string();
3737
3738 let params = TradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
3739
3740 let TradeParams {
3741 symbol,id,
3742 } = params.clone();
3743
3744 let pairs: &[(&str, Option<String>)] = &[
3745 ("symbol",
3746 Some(symbol.clone())
3747 ),
3748 ("id",
3749 id.clone()
3750 ),
3751 ];
3752
3753 let vars: HashMap<_, _> = pairs
3754 .iter()
3755 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3756 .collect();
3757 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3758
3759 let ws_stream = api.trade(params).await.unwrap();
3760
3761 let called = Arc::new(AtomicBool::new(false));
3762 let called_with_message = called.clone();
3763 ws_stream.on_message(move |_payload: models::TradeResponse| {
3764 called_with_message.store(true, Ordering::SeqCst);
3765 });
3766
3767 let payload: Value = serde_json::from_str(r#"{"e":"trade","E":1672515782136,"s":"BNBBTC","t":12345,"T":1672515782136,"m":true,"M":true}"#).unwrap_or_else(|_| serde_json::json!({}));
3768 let msg = json!({
3769 "stream": stream,
3770 "data": payload,
3771 });
3772
3773 streams_base.on_message(msg.to_string(), conn.clone()).await;
3774 yield_now().await;
3775
3776 assert!(called.load(Ordering::SeqCst), "expected our callback to have been invoked");
3777 });
3778 }
3779
3780 #[test]
3781 fn trade_should_not_fire_after_unsubscribe() {
3782 TOKIO_SHARED_RT.block_on(async {
3783 let (streams_base, conn) = make_streams_base().await;
3784 let api = ApiClient::new(streams_base.clone());
3785
3786 let id = "test-id-123".to_string();
3787
3788 let params = TradeParams::builder("bnbusdt".to_string(),).id(Some(id.clone())).build().unwrap();
3789
3790 let TradeParams {
3791 symbol,id,
3792 } = params.clone();
3793
3794 let pairs: &[(&str, Option<String>)] = &[
3795 ("symbol",
3796 Some(symbol.clone())
3797 ),
3798 ("id",
3799 id.clone()
3800 ),
3801 ];
3802
3803 let vars: HashMap<_, _> = pairs
3804 .iter()
3805 .filter_map(|&(k, ref v)| v.clone().map(|v| (k, v)))
3806 .collect();
3807 let stream = replace_websocket_streams_placeholders("/<symbol>@trade", &vars);
3808
3809 let ws_stream = api.trade(params).await.unwrap();
3810
3811 let called = Arc::new(AtomicBool::new(false));
3812 let called_clone = called.clone();
3813 ws_stream.on_message(move |_payload: models::TradeResponse| {
3814 called_clone.store(true, Ordering::SeqCst);
3815 });
3816
3817 assert!(streams_base.is_subscribed(&stream).await, "should be subscribed before unsubscribe");
3818
3819 ws_stream.unsubscribe().await;
3820
3821 let payload: Value = serde_json::from_str(r#"{"e":"trade","E":1672515782136,"s":"BNBBTC","t":12345,"T":1672515782136,"m":true,"M":true}"#).unwrap_or_else(|_| serde_json::json!({}));
3822 let msg = json!({
3823 "stream": stream,
3824 "data": payload,
3825 });
3826
3827 streams_base.on_message(msg.to_string(), conn.clone()).await;
3828
3829 yield_now().await;
3830
3831 assert!(!called.load(Ordering::SeqCst), "callback should not be invoked after unsubscribe");
3832 });
3833 }
3834}