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