Skip to main content

binance_sdk/spot/websocket_streams/apis/
api.rs

1/*
2 * Spot WebSocket Market Streams
3 *
4 * Access market data, manage accounts, and trade on Binance Spot.
5 *
6 * The version of the OpenAPI document: 1.0.0
7 *
8 *
9 * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
10 * https://openapi-generator.tech
11 * Do not edit the class manually.
12 */
13
14#![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/// Request parameters for the [`agg_trade`] operation.
443///
444/// This struct holds all of the inputs you can pass when calling
445/// [`agg_trade`](#method.agg_trade).
446#[derive(Clone, Debug, Builder, Deserialize)]
447#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
448pub struct AggTradeParams {
449    /// Symbol to query
450    ///
451    /// This field is **required.
452    #[builder(setter(into))]
453    #[serde(rename = "symbol")]
454    pub symbol: String,
455    /// Unique WebSocket request ID.
456    ///
457    /// This field is **optional.
458    #[builder(setter(into), default)]
459    #[serde(rename = "id", default)]
460    pub id: Option<String>,
461}
462
463impl AggTradeParams {
464    /// Create a builder for [`agg_trade`].
465    ///
466    /// Required parameters:
467    ///
468    /// * `symbol` — Symbol to query
469    ///
470    #[must_use]
471    pub fn builder(symbol: String) -> AggTradeParamsBuilder {
472        AggTradeParamsBuilder::default().symbol(symbol)
473    }
474}
475/// Request parameters for the [`all_market_rolling_window_ticker`] operation.
476///
477/// This struct holds all of the inputs you can pass when calling
478/// [`all_market_rolling_window_ticker`](#method.all_market_rolling_window_ticker).
479#[derive(Clone, Debug, Builder, Deserialize)]
480#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
481pub struct AllMarketRollingWindowTickerParams {
482    ///
483    /// The `window_size` parameter.
484    ///
485    /// This field is **required.
486    #[builder(setter(into))]
487    #[serde(rename = "windowSize")]
488    pub window_size: AllMarketRollingWindowTickerWindowSizeEnum,
489    /// Unique WebSocket request ID.
490    ///
491    /// This field is **optional.
492    #[builder(setter(into), default)]
493    #[serde(rename = "id", default)]
494    pub id: Option<String>,
495}
496
497impl AllMarketRollingWindowTickerParams {
498    /// Create a builder for [`all_market_rolling_window_ticker`].
499    ///
500    /// Required parameters:
501    ///
502    /// * `window_size` — String
503    ///
504    #[must_use]
505    pub fn builder(
506        window_size: AllMarketRollingWindowTickerWindowSizeEnum,
507    ) -> AllMarketRollingWindowTickerParamsBuilder {
508        AllMarketRollingWindowTickerParamsBuilder::default().window_size(window_size)
509    }
510}
511/// Request parameters for the [`all_mini_ticker`] operation.
512///
513/// This struct holds all of the inputs you can pass when calling
514/// [`all_mini_ticker`](#method.all_mini_ticker).
515#[derive(Clone, Debug, Builder, Deserialize, Default)]
516#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
517pub struct AllMiniTickerParams {
518    /// Unique WebSocket request ID.
519    ///
520    /// This field is **optional.
521    #[builder(setter(into), default)]
522    #[serde(rename = "id", default)]
523    pub id: Option<String>,
524}
525
526impl AllMiniTickerParams {
527    /// Create a builder for [`all_mini_ticker`].
528    ///
529    #[must_use]
530    pub fn builder() -> AllMiniTickerParamsBuilder {
531        AllMiniTickerParamsBuilder::default()
532    }
533}
534/// Request parameters for the [`avg_price`] operation.
535///
536/// This struct holds all of the inputs you can pass when calling
537/// [`avg_price`](#method.avg_price).
538#[derive(Clone, Debug, Builder, Deserialize)]
539#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
540pub struct AvgPriceParams {
541    /// Symbol to query
542    ///
543    /// This field is **required.
544    #[builder(setter(into))]
545    #[serde(rename = "symbol")]
546    pub symbol: String,
547    /// Unique WebSocket request ID.
548    ///
549    /// This field is **optional.
550    #[builder(setter(into), default)]
551    #[serde(rename = "id", default)]
552    pub id: Option<String>,
553}
554
555impl AvgPriceParams {
556    /// Create a builder for [`avg_price`].
557    ///
558    /// Required parameters:
559    ///
560    /// * `symbol` — Symbol to query
561    ///
562    #[must_use]
563    pub fn builder(symbol: String) -> AvgPriceParamsBuilder {
564        AvgPriceParamsBuilder::default().symbol(symbol)
565    }
566}
567/// Request parameters for the [`block_trade`] operation.
568///
569/// This struct holds all of the inputs you can pass when calling
570/// [`block_trade`](#method.block_trade).
571#[derive(Clone, Debug, Builder, Deserialize)]
572#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
573pub struct BlockTradeParams {
574    /// Symbol to query
575    ///
576    /// This field is **required.
577    #[builder(setter(into))]
578    #[serde(rename = "symbol")]
579    pub symbol: String,
580    /// Unique WebSocket request ID.
581    ///
582    /// This field is **optional.
583    #[builder(setter(into), default)]
584    #[serde(rename = "id", default)]
585    pub id: Option<String>,
586}
587
588impl BlockTradeParams {
589    /// Create a builder for [`block_trade`].
590    ///
591    /// Required parameters:
592    ///
593    /// * `symbol` — Symbol to query
594    ///
595    #[must_use]
596    pub fn builder(symbol: String) -> BlockTradeParamsBuilder {
597        BlockTradeParamsBuilder::default().symbol(symbol)
598    }
599}
600/// Request parameters for the [`book_ticker`] operation.
601///
602/// This struct holds all of the inputs you can pass when calling
603/// [`book_ticker`](#method.book_ticker).
604#[derive(Clone, Debug, Builder, Deserialize)]
605#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
606pub struct BookTickerParams {
607    /// Symbol to query
608    ///
609    /// This field is **required.
610    #[builder(setter(into))]
611    #[serde(rename = "symbol")]
612    pub symbol: String,
613    /// Unique WebSocket request ID.
614    ///
615    /// This field is **optional.
616    #[builder(setter(into), default)]
617    #[serde(rename = "id", default)]
618    pub id: Option<String>,
619}
620
621impl BookTickerParams {
622    /// Create a builder for [`book_ticker`].
623    ///
624    /// Required parameters:
625    ///
626    /// * `symbol` — Symbol to query
627    ///
628    #[must_use]
629    pub fn builder(symbol: String) -> BookTickerParamsBuilder {
630        BookTickerParamsBuilder::default().symbol(symbol)
631    }
632}
633/// Request parameters for the [`diff_book_depth`] operation.
634///
635/// This struct holds all of the inputs you can pass when calling
636/// [`diff_book_depth`](#method.diff_book_depth).
637#[derive(Clone, Debug, Builder, Deserialize)]
638#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
639pub struct DiffBookDepthParams {
640    /// Symbol to query
641    ///
642    /// This field is **required.
643    #[builder(setter(into))]
644    #[serde(rename = "symbol")]
645    pub symbol: String,
646    /// Unique WebSocket request ID.
647    ///
648    /// This field is **optional.
649    #[builder(setter(into), default)]
650    #[serde(rename = "id", default)]
651    pub id: Option<String>,
652    /// Optional stream update speed suffix
653    ///
654    /// This field is **optional.
655    #[builder(setter(into), default)]
656    #[serde(rename = "updateSpeed", default)]
657    pub update_speed: Option<DiffBookDepthUpdateSpeedEnum>,
658}
659
660impl DiffBookDepthParams {
661    /// Create a builder for [`diff_book_depth`].
662    ///
663    /// Required parameters:
664    ///
665    /// * `symbol` — Symbol to query
666    ///
667    #[must_use]
668    pub fn builder(symbol: String) -> DiffBookDepthParamsBuilder {
669        DiffBookDepthParamsBuilder::default().symbol(symbol)
670    }
671}
672/// Request parameters for the [`kline`] operation.
673///
674/// This struct holds all of the inputs you can pass when calling
675/// [`kline`](#method.kline).
676#[derive(Clone, Debug, Builder, Deserialize)]
677#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
678pub struct KlineParams {
679    /// Symbol to query
680    ///
681    /// This field is **required.
682    #[builder(setter(into))]
683    #[serde(rename = "symbol")]
684    pub symbol: String,
685    ///
686    /// The `interval` parameter.
687    ///
688    /// This field is **required.
689    #[builder(setter(into))]
690    #[serde(rename = "interval")]
691    pub interval: KlineIntervalEnum,
692    /// Unique WebSocket request ID.
693    ///
694    /// This field is **optional.
695    #[builder(setter(into), default)]
696    #[serde(rename = "id", default)]
697    pub id: Option<String>,
698}
699
700impl KlineParams {
701    /// Create a builder for [`kline`].
702    ///
703    /// Required parameters:
704    ///
705    /// * `symbol` — Symbol to query
706    /// * `interval` — String
707    ///
708    #[must_use]
709    pub fn builder(symbol: String, interval: KlineIntervalEnum) -> KlineParamsBuilder {
710        KlineParamsBuilder::default()
711            .symbol(symbol)
712            .interval(interval)
713    }
714}
715/// Request parameters for the [`kline_offset`] operation.
716///
717/// This struct holds all of the inputs you can pass when calling
718/// [`kline_offset`](#method.kline_offset).
719#[derive(Clone, Debug, Builder, Deserialize)]
720#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
721pub struct KlineOffsetParams {
722    /// Symbol to query
723    ///
724    /// This field is **required.
725    #[builder(setter(into))]
726    #[serde(rename = "symbol")]
727    pub symbol: String,
728    ///
729    /// The `interval` parameter.
730    ///
731    /// This field is **required.
732    #[builder(setter(into))]
733    #[serde(rename = "interval")]
734    pub interval: KlineOffsetIntervalEnum,
735    /// Unique WebSocket request ID.
736    ///
737    /// This field is **optional.
738    #[builder(setter(into), default)]
739    #[serde(rename = "id", default)]
740    pub id: Option<String>,
741}
742
743impl KlineOffsetParams {
744    /// Create a builder for [`kline_offset`].
745    ///
746    /// Required parameters:
747    ///
748    /// * `symbol` — Symbol to query
749    /// * `interval` — String
750    ///
751    #[must_use]
752    pub fn builder(symbol: String, interval: KlineOffsetIntervalEnum) -> KlineOffsetParamsBuilder {
753        KlineOffsetParamsBuilder::default()
754            .symbol(symbol)
755            .interval(interval)
756    }
757}
758/// Request parameters for the [`mini_ticker`] operation.
759///
760/// This struct holds all of the inputs you can pass when calling
761/// [`mini_ticker`](#method.mini_ticker).
762#[derive(Clone, Debug, Builder, Deserialize)]
763#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
764pub struct MiniTickerParams {
765    /// Symbol to query
766    ///
767    /// This field is **required.
768    #[builder(setter(into))]
769    #[serde(rename = "symbol")]
770    pub symbol: String,
771    /// Unique WebSocket request ID.
772    ///
773    /// This field is **optional.
774    #[builder(setter(into), default)]
775    #[serde(rename = "id", default)]
776    pub id: Option<String>,
777}
778
779impl MiniTickerParams {
780    /// Create a builder for [`mini_ticker`].
781    ///
782    /// Required parameters:
783    ///
784    /// * `symbol` — Symbol to query
785    ///
786    #[must_use]
787    pub fn builder(symbol: String) -> MiniTickerParamsBuilder {
788        MiniTickerParamsBuilder::default().symbol(symbol)
789    }
790}
791/// Request parameters for the [`partial_book_depth`] operation.
792///
793/// This struct holds all of the inputs you can pass when calling
794/// [`partial_book_depth`](#method.partial_book_depth).
795#[derive(Clone, Debug, Builder, Deserialize)]
796#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
797pub struct PartialBookDepthParams {
798    /// Symbol to query
799    ///
800    /// This field is **required.
801    #[builder(setter(into))]
802    #[serde(rename = "symbol")]
803    pub symbol: String,
804    ///
805    /// The `levels` parameter.
806    ///
807    /// This field is **required.
808    #[builder(setter(into))]
809    #[serde(rename = "levels")]
810    pub levels: PartialBookDepthLevelsEnum,
811    /// Unique WebSocket request ID.
812    ///
813    /// This field is **optional.
814    #[builder(setter(into), default)]
815    #[serde(rename = "id", default)]
816    pub id: Option<String>,
817    /// Optional stream update speed suffix
818    ///
819    /// This field is **optional.
820    #[builder(setter(into), default)]
821    #[serde(rename = "updateSpeed", default)]
822    pub update_speed: Option<PartialBookDepthUpdateSpeedEnum>,
823}
824
825impl PartialBookDepthParams {
826    /// Create a builder for [`partial_book_depth`].
827    ///
828    /// Required parameters:
829    ///
830    /// * `symbol` — Symbol to query
831    /// * `levels` — String
832    ///
833    #[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/// Request parameters for the [`reference_price`] operation.
844///
845/// This struct holds all of the inputs you can pass when calling
846/// [`reference_price`](#method.reference_price).
847#[derive(Clone, Debug, Builder, Deserialize)]
848#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
849pub struct ReferencePriceParams {
850    /// Symbol to query
851    ///
852    /// This field is **required.
853    #[builder(setter(into))]
854    #[serde(rename = "symbol")]
855    pub symbol: String,
856    /// Unique WebSocket request ID.
857    ///
858    /// This field is **optional.
859    #[builder(setter(into), default)]
860    #[serde(rename = "id", default)]
861    pub id: Option<String>,
862}
863
864impl ReferencePriceParams {
865    /// Create a builder for [`reference_price`].
866    ///
867    /// Required parameters:
868    ///
869    /// * `symbol` — Symbol to query
870    ///
871    #[must_use]
872    pub fn builder(symbol: String) -> ReferencePriceParamsBuilder {
873        ReferencePriceParamsBuilder::default().symbol(symbol)
874    }
875}
876/// Request parameters for the [`rolling_window_ticker`] operation.
877///
878/// This struct holds all of the inputs you can pass when calling
879/// [`rolling_window_ticker`](#method.rolling_window_ticker).
880#[derive(Clone, Debug, Builder, Deserialize)]
881#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
882pub struct RollingWindowTickerParams {
883    /// Symbol to query
884    ///
885    /// This field is **required.
886    #[builder(setter(into))]
887    #[serde(rename = "symbol")]
888    pub symbol: String,
889    ///
890    /// The `window_size` parameter.
891    ///
892    /// This field is **required.
893    #[builder(setter(into))]
894    #[serde(rename = "windowSize")]
895    pub window_size: RollingWindowTickerWindowSizeEnum,
896    /// Unique WebSocket request ID.
897    ///
898    /// This field is **optional.
899    #[builder(setter(into), default)]
900    #[serde(rename = "id", default)]
901    pub id: Option<String>,
902}
903
904impl RollingWindowTickerParams {
905    /// Create a builder for [`rolling_window_ticker`].
906    ///
907    /// Required parameters:
908    ///
909    /// * `symbol` — Symbol to query
910    /// * `window_size` — String
911    ///
912    #[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/// Request parameters for the [`ticker`] operation.
923///
924/// This struct holds all of the inputs you can pass when calling
925/// [`ticker`](#method.ticker).
926#[derive(Clone, Debug, Builder, Deserialize)]
927#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
928pub struct TickerParams {
929    /// Symbol to query
930    ///
931    /// This field is **required.
932    #[builder(setter(into))]
933    #[serde(rename = "symbol")]
934    pub symbol: String,
935    /// Unique WebSocket request ID.
936    ///
937    /// This field is **optional.
938    #[builder(setter(into), default)]
939    #[serde(rename = "id", default)]
940    pub id: Option<String>,
941}
942
943impl TickerParams {
944    /// Create a builder for [`ticker`].
945    ///
946    /// Required parameters:
947    ///
948    /// * `symbol` — Symbol to query
949    ///
950    #[must_use]
951    pub fn builder(symbol: String) -> TickerParamsBuilder {
952        TickerParamsBuilder::default().symbol(symbol)
953    }
954}
955/// Request parameters for the [`trade`] operation.
956///
957/// This struct holds all of the inputs you can pass when calling
958/// [`trade`](#method.trade).
959#[derive(Clone, Debug, Builder, Deserialize)]
960#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
961pub struct TradeParams {
962    /// Symbol to query
963    ///
964    /// This field is **required.
965    #[builder(setter(into))]
966    #[serde(rename = "symbol")]
967    pub symbol: String,
968    /// Unique WebSocket request ID.
969    ///
970    /// This field is **optional.
971    #[builder(setter(into), default)]
972    #[serde(rename = "id", default)]
973    pub id: Option<String>,
974}
975
976impl TradeParams {
977    /// Create a builder for [`trade`].
978    ///
979    /// Required parameters:
980    ///
981    /// * `symbol` — Symbol to query
982    ///
983    #[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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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").map(std::string::ToString::to_string);
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}