Skip to main content

binance_sdk/alpha/websocket_streams/apis/
api.rs

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